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:
# 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}'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:
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())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.