Skip to content

parse-node-outputs

Unified bounded/unbounded Flink job: HyperLiquid node outputs (JSONL) → typed rows → Parquet → Iceberg.

Modes

Mode Source Execution Trigger
Live Kafka per-table data-plane topics (hyperliquid.book-diffs / .order-statuses / .fills) — one message per line unbounded checkpoint interval (2–5 min) or writer size (128–512 MB)
Replay Kafka (identical — replay tap's progressive appends feed the same topics) unbounded same as live
Backfill file list (S3/SSD enumeration, no Kafka payloads) bounded one task per snapshot-aligned segment [S_i, S_i+1)

Design constraints

  • Data plane = streamed lines (keyed coin); seals (hyperliquid.node-files) drive reconciliation, not parsing. Backfill reads files directly — Kafka never carries bulk reprocessing.
  • Deterministic row keys (table, block_number, log_index); idempotent replays.
  • Checkpoint-aligned Iceberg commits (exactly-once); adaptive file sizing.
  • Backfill segments self-validate: replayed end-state must hash-match the trailing abci snapshot.

Full pipeline context: services/ingestion/README.md.

Status

Skeleton — not implemented.