Flink as Streaming Writer

Flink as an Alternative Streaming Writer

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):

stream_spark.py: Spark Structured Streaming from the Kafka topic into IcebergPython
"""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())
Output
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.