Skip to content

Spark Structured Streaming from Kafka to Iceberg: how to write and keep tables healthy

Spark
Chad Harris·October 6, 2026·8 min read

A Spark Structured Streaming job writes a Kafka topic into an Apache Iceberg table by reading the topic with the Kafka source and writing with format("iceberg"), a checkpoint location and a trigger interval of at least one minute. The write is a few lines. What decides whether the table stays healthy is the trigger interval, how a partitioned table is written, and the maintenance a streaming table needs afterwards. Factor House has no Spark product and publishes this page as a guide to a common pipeline around the Kafka and Flink tools it does make.

The write

The Iceberg Spark Structured Streaming documentation shows the write as a DataStreamWriter with the iceberg format, an output mode, a trigger and a checkpoint location, finishing in toTable with the table name. Two output modes are supported. Append adds the rows of every micro-batch to the table, and complete replaces the table contents every micro-batch. The table must exist before the query starts. The same page says Iceberg does not support Spark’s experimental continuous processing, because it provides no interface to commit the output.

The read side is the Kafka source, which keeps its position in the checkpoint and commits no offsets to Kafka. So the checkpoint location is the record of which Kafka offsets are already in the table, and it is worth protecting like the table itself. Where Spark keeps its position covers what that means for lag.

Factor House publishes a working example. Lab 10 of Factor House Local builds a PySpark Structured Streaming job that ingests Kafka order data, deserializes Avro messages and writes the results to Iceberg, and it runs the Kafka environment with Kpow.

Set the trigger interval first

Every micro-batch is a commit, and every commit makes a snapshot. The Iceberg documentation says a high rate of commits produces data files, manifests and snapshots that lead to additional maintenance, and recommends a trigger interval of one minute at the minimum, increased if needed. A trigger of a few seconds is the most common cause of the table problems described in too many small files.

Write partitioned tables with the fanout writer

Iceberg requires data to be sorted by partition within each write task. The documentation says that sorting in a streaming query adds latency because repartition and sort are heavy operations, and that you can enable the fanout writer with the fanout-enabled option to remove the requirement. The fanout writer keeps a file open for every partition value until the write task finishes, and the documentation advises against it for batch writes, where an explicit sort is cheap.

Maintain a streaming table

The Iceberg documentation says streaming writes create new table versions quickly and creating lots of metadata, so maintenance is highly recommended. It names three jobs.

  • Expire old snapshots. Each batch produces a snapshot, and by default the procedure expires snapshots older than five days. See Iceberg storage keeps growing.
  • Compact data files. The data written by a streaming process is typically small, which leaves the table tracking many small files. Iceberg and Spark provide the rewrite_data_files procedure. See the compaction guide.
  • Rewrite manifests. To keep write latency low, Iceberg uses a fast append that does not compact manifests, which can produce many small manifest files.

If two writers touch the table, commits can conflict, covered in Iceberg commits failing or conflicting.

Spark or another writer

Spark is one of three common routes from Kafka to Iceberg. The Kafka Connect sink and Flink are the others, and Factor House’s ranking of tools for Kafka to Iceberg pipelines scores the Kafka and Flink tools that manage those two routes. Kpow, which Factor House makes, manages the Kafka side of any of them, and it does not read Spark jobs or Iceberg tables.

FAQ

How do I write a Kafka topic to Iceberg with Spark Structured Streaming?

Read the topic with the Kafka source, then write with writeStream.format("iceberg"), append output mode, a trigger of one minute or more, a checkpoint location and toTable. Create the Iceberg table before starting the query.

What trigger interval should an Iceberg streaming write use?

The Iceberg documentation recommends one minute at the minimum, and increasing it if needed, because each commit produces files, manifests and a snapshot.

Do I need fanout-enabled for a partitioned Iceberg table?

Not required, but without it Iceberg requires data sorted by partition in each task, which adds latency in a streaming query. The fanout writer removes the requirement and holds a file open for each partition value until the task finishes.

Does Spark Structured Streaming commit Kafka offsets when writing to Iceberg?

No. The Kafka source commits no offsets. The position is saved in the query’s checkpoint location.

Does Factor House support Spark?

No. Factor House makes Kpow for Kafka and Flex for Flink, and publishes Spark only in its open source Factor House Local labs.

Related reading