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.