Small Commits and Compaction

Small Commits, Compaction and Exactly-Once Landing

A live stream lands in small files. live_events.py generates sample events for 1 July 2026 (an order placed and paid every 0.1 seconds); the script feeds them for 90 seconds and kills the Connect worker midway:

live.sh: a live trickle of events, with the Connect worker killed midwayShell
# Stream 20 events/s for 90 s into booknest.order-events, killing the Connect worker midway.
python3 "$DEMOS/streaming/live_events.py" 20 90 |
  docker exec -i l1-kafka /opt/kafka/bin/kafka-console-producer.sh \
    --bootstrap-server l1-kafka:9092 --topic booknest.order-events \
    --reader-property parse.key=true --reader-property key.separator='|' &
sleep 40; docker kill l1-connect >/dev/null; echo "$(date -u +%T) worker killed (SIGKILL)"
sleep 5; docker start l1-connect >/dev/null; echo "$(date -u +%T) worker restarted"
wait; sleep 45                                     # let the last commit happen
bash "$DEMOS/streaming/commits.sh"
docker exec l1-kafka /opt/kafka/bin/kafka-get-offsets.sh --bootstrap-server l1-kafka:9092 \
  --topic booknest.order-events | awk -F: '{n += $3} END {print "messages in topic:", n}'
Output
02:10:43 worker killed (SIGKILL)
02:10:48 worker restarted
time      op      rows    files  total_rows  total_files
02:10:01  append  390737  54     390737      54
02:10:22  append  370     3      391107      57
02:10:36  append  288     3      391395      60
02:11:53  append  1142    3      392537      63
messages in topic: 392537

Exactly-once held: the table has as many rows as the topic has messages, and count(DISTINCT event_id) in DuckDB 61,228 agreed, although the worker died holding written but uncommitted files. No commit happened while the worker was down, so the one after the restart (02:11:53) is larger: the new tasks re-read everything after the last committed offsets, and the dead worker's files stayed orphans. Small files pile up: each commit adds a file per task and partition, here three of about 300 rows each; at 15 seconds that is 17,280 files a day in one partition. Compaction cures it, and optimistic concurrency lets it run beside the sink:

compact_events.py: compact the sink's small files while it keeps writingPython
from lake import spark
files = "SELECT count(*) AS files, round(avg(file_size_in_bytes) / 1024) AS avg_kib " \
        "FROM booknest.order_events.files"
print("before:", spark.sql(files).first())
r = spark.sql("""CALL system.rewrite_data_files(table => 'booknest.order_events',
    options => map('min-input-files', '2', 'target-file-size-bytes', '134217728'))""")
print(r.select("rewritten_data_files_count", "added_data_files_count").first())
print("after: ", spark.sql(files).first())
Output
before: Row(files=63, avg_kib=55.0)
Row(rewritten_data_files_count=63, added_data_files_count=19)
after:  Row(files=19, avg_kib=181.0)

Nineteen files remain, one per month. In production, compact per partition on a schedule (Compaction and Expiry), keep the interval at minutes unless readers need fresher data, and expire old snapshots to delete the small files.