Ingestor¶
The ingestor moves JSON events from Kafka into an Iceberg table. It writes each event to a local Parquet write-ahead log (WAL) first. It then uploads the finished files to your storage and commits them to the table in batches.
Open the interactive diagram for pan and zoom, search, and export to PNG or SVG.
The record path¶
- Consume. Each dataset has its own Kafka consumer on its own topic.
- Decode. The dataset's
customMapping(JMESPath) maps the JSON payload to columns. The ingestor drops a record that it cannot decode, and counts it. One bad record does not stop the ingestor. - Append. The ingestor adds the row to an open WAL segment, a file named
active_*.parquet. Each segment holds rows from one Kafka partition for one Iceberg partition. - Rotate. The ingestor closes the segment, syncs it to disk and renames it to
ready_{topic}_{partition}_{start}-{end}_{uuid}.parquet. The name records the Kafka offsets that the file holds. Then the ingestor commits the Kafka offsetend + 1. - Upload. The uploader scans the WAL directory every 5 seconds and writes each
ready_*file to Azure Blob. The object name comes from the file's UUID, so a repeated upload replaces the same object. If the dataset hasbloomFilterColumns, the uploader also writes a bloom-filter sidecar for the file. - Commit. The committer collects the uploaded files. At each interval it adds all the files of
a dataset to the table in one Iceberg append. After the append succeeds, the ingestor deletes
the local
ready_*files.
When a segment rotates¶
| Trigger | Configuration key | Default |
|---|---|---|
| The file reaches its maximum size | ingestion.wal.maxFileSizeBytes |
64 MiB |
| The file reaches its maximum age | ingestion.wal.maxFileAgeMinutes |
5 minutes |
| A Kafka partition has too many open segments | ingestion.wal.maxOpenSegmentsPerPartition |
256 |
| Kafka takes the partition away in a rebalance | ||
| The ingestor shuts down |
The ingestor checks the age of idle segments every 30 seconds. All the open segments of one Kafka partition rotate together, because one committed offset covers all of them.
Two rules that protect your data¶
The Kafka offset is committed at rotation, not after the table commit. When that offset
commit lands, Kafka does not deliver those rows again. Until the table commit succeeds, the local
ready_* file is their only copy.
The local file stays until the table commit succeeds. An upload alone does not delete it. If the upload fails, the file waits for the next scan. If the table commit fails, the committer tries the same batch again at the next interval.
The commit interval is ingestion.commitIntervalSeconds (default 20 seconds). Each pod changes the
interval by a small amount of its own, so that the pods do not all commit together. One append for each
dataset in each interval keeps the number of catalog commits low.
Backpressure¶
The ingestor does not drop records when your storage or your catalog is unavailable. The ready_*
files stay on the WAL volume. Every 5 seconds the ingestor measures how full the volume is:
- at
ingestion.wal.diskHighWatermarkPct(default 80 %) it pauses all its Kafka consumers; - 5 percentage points below that value it resumes them.
While the consumers are paused, the events stay in Kafka. Your Kafka retention is the limit on how long an outage can continue without data loss.
Recovery after a crash¶
The rename and the offset commit are two separate steps, so a crash can occur between them. At start, before the uploader runs, the ingestor compares each file in the WAL directory with the offset that Kafka has committed for its partition:
| File | Committed offset | Action | Reason |
|---|---|---|---|
active_* |
any | Delete | The offsets were not committed. Kafka delivers the records again. |
ready_* |
none, or not after the file's first offset | Delete | The offset commit did not occur. Kafka delivers the records again, and a second copy would duplicate rows. |
ready_* |
after the file's last offset, file not in the table | Keep | Kafka does not deliver the records again. The uploader and the committer finish the work. |
ready_* |
after the file's last offset, file already in the table | Delete | The data is in the table. Only the local copy goes. |
ready_* |
inside the file's offset range | Stop | This cannot occur in correct operation. The ingestor does not start. |
The check needs Kafka and the catalog. If the ingestor cannot read the committed offsets, it does not start. "In the table" means that a snapshot that the table still keeps refers to the file.
The diagram shows the same rules as states of one file. The top row is normal operation: a file goes through the numbered states in order, and the diagram draws no arrows between them. The ingestor does the check at each start, not only after a crash:
- A file that is
OpenorRotatedis discarded, because Kafka delivers its records again. - A file that is
ReadyorUploadedstays. The uploader writes it to storage, to the same object name as before, and the committer adds it to the table. - A file that is
In the tableloses only its local copy.
The last row of the table above is not a state of a file. It is an error, and the ingestor does not start.
An error while the ingestor writes the WAL is fatal. The process exits, the pod starts again, and this recovery runs.
Shutdown¶
On SIGTERM the ingestor does these steps in this order:
- Release the partitions. The ingestor starts to leave its Kafka consumer groups. Before Kafka takes the partitions, the ingestor rotates the open segments and commits their offsets. Then it leaves. Kafka gives the partitions to the other pods immediately; it does not wait for a session timeout.
- Upload. The ingestor writes the
ready_*files to Azure Blob. - Commit. The ingestor adds the uploaded files to the table, one append for each dataset, and deletes the local files.
- Exit.
The metrics endpoint stays available during these steps. The limit for all of them is
ingestion.drainTimeoutSeconds (default 96). Step 1 can use a quarter of that time at most, and
not more than 10 seconds. If Kafka has not taken the partitions by then, the ingestor rotates the
open segments itself and continues. If the time runs out, the ingestor exits, and recovery completes
the work at the next start.
The ingestor does not write records that arrive after it released a partition. The pod that gets the partition reads them from the committed offset.