When the stream needs work on the way in (joins, deduplication, upserts by key), a stream processor writes the table. Apache Flink 129 (2.3.0, June 2026; Iceberg 1.12.0 129 ships iceberg-flink-runtime for Flink 1.20 and 2.1-2.3) commits once per checkpoint: the checkpoint flushes the files and stores the Kafka 129 offsets in Flink's state, and the files join the table only when it completes, exactly once without a control topic. With upsert-enabled on a v2 table with identifier fields, Flink writes equality deletes and keeps one row per key; the Dynamic Iceberg Sink routes records and evolves schemas at runtime. Flink was not run here (its JobManager and TaskManager would crowd this host); Spark 129 runs the same pattern with Structured Streaming (Structured Streaming):
"""Spark Structured Streaming: Kafka topic booknest.order-events into an Iceberg table."""
import os
os.environ["PYSPARK_SUBMIT_ARGS"] = \
"--driver-class-path '/home/dev/v7-l1/ch08/jars/kafka/*' pyspark-shell" # Kafka source
from pyspark.sql import functions as F
from lake import spark
schema = "event_id BIGINT, ts TIMESTAMP, type STRING, order_id BIGINT, customer_id BIGINT, " \
"total DECIMAL(10,2)"
spark.sql(f"""CREATE TABLE IF NOT EXISTS booknest.order_events_spark ({schema})
USING iceberg PARTITIONED BY (months(ts))""")
events = (spark.readStream.format("kafka")
.option("kafka.bootstrap.servers", "localhost:31092")
.option("subscribe", "booknest.order-events")
.option("startingOffsets", "earliest")
.option("maxOffsetsPerTrigger", 150000).load() # at most 150,000 per batch
.select(F.from_json(F.col("value").cast("string"), schema).alias("e")).select("e.*"))
for run in (1, 2): # run 2 resumes from the checkpoint: nothing new to read
q = (events.writeStream.format("iceberg").outputMode("append")
.trigger(availableNow=True) # drain what is there, then stop
.option("checkpointLocation", "/home/dev/v7-l1/ch08/checkpoints/order_events_spark")
.toTable("booknest.order_events_spark"))
q.awaitTermination()
print(f"run {run}: batches", [(p["batchId"], p["numInputRows"]) for p in q.recentProgress])
print(spark.sql("""SELECT count(*) AS events, count(DISTINCT event_id) AS ids,
count_if(type = 'order_placed') AS orders FROM booknest.order_events_spark""").first())run 1: batches [(0, 149999), (1, 149999), (2, 92539)] run 2: batches [(3, 0)] Row(events=392537, ids=392537, orders=100900)
The Kafka source jars (spark-sql-kafka-0-10_2.13 4.1.3, kafka-clients 3.9.1, commons-pool2) came from Maven 129 . Each micro-batch is one commit whose snapshot summary records the query ID and spark.sql.streaming.epochId, and the checkpoint holds the offsets, so a replayed batch is a no-op and the second run read nothing. availableNow suits a scheduled job (Orchestration and Pipelines); processingTime triggers keep it running.