Compaction and Expiry

Compacting Small Files and Expiring Snapshots

Streaming and frequent loads produce small files, each costing a request, a footer read and a manifest entry. demos/ch08/production/small_files.py writes the last week of June into a scratch table orders_stream (partitioned by day) as 12 micro-batches of four tasks, as a streaming job would (Streaming into the Lakehouse). Compaction then rewrites the files and expiry deletes what no snapshot needs:

compact.py: bin-pack the small files, then expire old snapshotsPython
"""Compaction and snapshot expiry on a table written by 12 tiny commits (small_files.py)."""
from lake import spark
S = "booknest.orders_stream"
def stats(when):
    r = spark.sql(f"""SELECT (SELECT count(*) FROM {S}.snapshots) AS snaps,
        (SELECT count(*) FROM {S}.manifests) AS mans, count(*) AS files,
        round(avg(file_size_in_bytes) / 1024, 1) AS kb FROM {S}.files""").first()
    print(f"{when}: {r.snaps} snapshots, {r.mans} manifests, {r.files} files of {r.kb} KB")
stats("before")
r = spark.sql(f"CALL lake.system.rewrite_data_files('{S}')").first()     # bin-packing
print(f"rewrite: {r.rewritten_data_files_count} -> {r.added_data_files_count} files")
r = spark.sql(f"""CALL lake.system.expire_snapshots(table => '{S}', older_than => now(),
                  retain_last => 1)""").first()                     # keep only the newest
print(f"expire: deleted {r.deleted_data_files_count} data files, "
      f"{r.deleted_manifest_files_count} manifests, {r.deleted_manifest_lists_count} lists")
stats("after")
Output
before: 13 snapshots, 12 manifests, 335 files of 2.4 KB
rewrite: 335 -> 7 files
expire: deleted 335 data files, 12 manifests, 13 lists
after: 1 snapshots, 13 manifests, 7 files of 4.8 KB

rewrite_data_files bin-packed 335 files into one per day toward a 512 MB target (sort and zorder strategies also cluster rows), and folds delete files in on merge-on-read tables. The old files stay until expire_snapshots removes the snapshots that reference them; only then is storage freed and time travel to those snapshots gone. Trino 403,499 has the same steps as ALTER TABLE ... EXECUTE optimize and EXECUTE expire_snapshots. Schedule both from the orchestrator (Orchestration and Pipelines), with a retention longer than your longest job.