Skip to content

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-cluster via 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).