Skip to content

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.

Ingestor flow Ingestor flow

Open the interactive diagram for pan and zoom, search, and export to PNG or SVG.

The record path

  1. Consume. Each dataset has its own Kafka consumer on its own topic.
  2. 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.
  3. 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.
  4. 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 offset end + 1.
  5. 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 has bloomFilterColumns, the uploader also writes a bloom-filter sidecar for the file.
  6. 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 Open or Rotated is discarded, because Kafka delivers its records again.
  • A file that is Ready or Uploaded stays. 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 table loses 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.

WAL file states WAL file states

Open the interactive diagram.

An error while the ingestor writes the WAL is fatal. The process exits, the pod starts again, and this recovery runs.

Shutdown

Ingestor shutdown sequence Ingestor shutdown sequence

Open the interactive diagram.

On SIGTERM the ingestor does these steps in this order:

  1. 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.
  2. Upload. The ingestor writes the ready_* files to Azure Blob.
  3. Commit. The ingestor adds the uploaded files to the table, one append for each dataset, and deletes the local files.
  4. 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.