hyperdata-ingestion-flink
Ingestion-stage Apache Flink jobs for the HyperData Platform. One repo per engine per layer (see hyperdata-platform ingestion design); transformation/modeling jobs live elsewhere.
Jobs
| Job | Mode | What it does |
|---|---|---|
parse-node-outputs |
bounded + unbounded (planned) | The core job. Streams node output lines and lands them as Parquet/Iceberg. One parse core, three drivers: live + replay via the per-table data-plane topics (hyperliquid.book-diffs / .order-statuses / .fills, one message per line), backfill via bounded file enumeration. |
samples/ (future) |
— | Learning walkthroughs from the Flink docs. Not production. |
Contract
- Source of truth for behavior: services/ingestion/README.md in the platform monorepo (commit duality, adaptive file sizing, idempotency keys).
- Data plane = streamed lines: per-table topics carry the node's JSON lines as they're appended (keyed
coin); the seal topic (hyperliquid.node-files, one message per finalized hour-file) carries durability watermarks + sha256 for reconciliation. Backfill bypasses Kafka entirely (bounded file enumeration). - Exactly-once into the warehouse: checkpoint-aligned Iceberg commits; deterministic row keys
(table, block_number, log_index)make replays no-ops. - Runtime: jobs are submitted to a session cluster (
platform/flink-clustervia compose); each job builds a self-contained fat jar (Java, connectors bundled) — the cluster image stays generic.
Status
Skeleton — job implementation not started. Language: Java fat-jar (decided: high message volume + Iceberg/Kafka connector maturity favor the native runtime; PyFlink reserved for samples/stateful showcase jobs).