In Delta Lake, what does OPTIMIZE change about a table's data files?
answer
- fewer, bigger files; same rows
- the log records the swap, storage keeps both
- one commit: removes plus adds
- flagged so streams ignore the rewrite
- bin-packing up to a target file size
basics
~20 sOPTIMIZE bin-packs many small Parquet files into fewer large ones (roughly 1 GB by default), committing remove and add actions in a new table version. Row content is unchanged, and the replaced files stay on storage until VACUUM deletes them.
solid answer
~50 s`OPTIMIZE tbl` reads the data files that sit **below the target file size** and rewrites their rows into a smaller number of larger Parquet files — bin-packing, not sorting. The result is one new commit in `_delta_log` containing a `remove` action per old file and an `add` action per new file, both flagged `dataChange: false` so streaming readers treat the rewrite as a no-op rather than re-emitting rows. Files that are already at or above the target are skipped, which makes repeated runs cheap and effectively idempotent. An optional `WHERE` clause restricts the work to given partitions. OPTIMIZE does not free storage: the old files remain on object storage, referenced only by historical versions, until `VACUUM` passes the retention threshold and deletes them. Verify what it did with `DESCRIBE HISTORY`, whose `operationMetrics` report `numRemovedFiles`, `numAddedFiles` and `totalFilesSkipped`.
go deeper
Know that a Delta table is many Parquet files and that lots of tiny ones make queries slow. Be able to say that OPTIMIZE combines small files into bigger ones without changing the data.
Explain the mechanics: which files are eligible, the target file size, that the result is one atomic commit with remove and add actions, and that storage is only reclaimed later by VACUUM.
Show operating judgment — scoping OPTIMIZE to recently written partitions, reading operationMetrics to justify the cost, and knowing the dataChange flag is what keeps downstream streaming jobs from re-processing everything.
Own the tradeoff across a platform: which tables earn compaction at all, whether auto-compaction at write time beats scheduled runs, and how the write amplification and compute bill compare with the query savings you can actually measure.
## What OPTIMIZE is for A Delta Lake table is a directory of immutable Parquet data files plus a `_delta_log` directory that records, commit by commit, which files are part of the table. Nothing rewrites a data file in place. Every append, every streaming micro-batch, every `MERGE` that touches a handful of rows adds new files. A table fed by a one-minute streaming job accumulates on the order of a thousand files a day, most of them a few megabytes or less. That hurts reads in ways that have nothing to do with how much data there is: the engine must list the files, open each footer, plan a task per file, and pay per-request latency on object storage. It also inflates the log — every one of those files is an `add` entry that must be replayed to compute the current table state. `OPTIMIZE` is the command that collapses that sprawl. ## What the command actually does ```sql OPTIMIZE sales.events WHERE event_date >= '2026-08-01'; ``` OPTIMIZE selects data files **below the target file size**, reads their rows, and writes them out as a smaller number of larger files. This is *bin-packing*: choose sets of small files whose combined size lands near the target, and concatenate them. Rows are not sorted, filtered, or changed in any way — unless you add `ZORDER BY`, which is a different, more expensive mode. The target file size defaults to about 1 GB. On Databricks it is tunable per table with the `delta.targetFileSize` property (and can be auto-tuned by table size); in open-source Delta the corresponding knob is a Spark configuration on the optimize command. Files already at or above the target are left untouched and counted in `totalFilesSkipped`, which is why running OPTIMIZE twice in a row is nearly free the second time — bin-packing is effectively idempotent. The optional `WHERE` predicate must filter on partition columns. That matters operationally: on a large historical table you almost never want to optimize everything, only the partitions that received writes since the last run. ## What lands in the transaction log OPTIMIZE produces exactly one new commit — a new numbered JSON file in `_delta_log`. Inside it: - a `commitInfo` action with `"operation": "OPTIMIZE"` and its `operationParameters` / `operationMetrics`, - one `remove` action per rewritten input file, - one `add` action per newly written file, carrying the new file's size, partition values and column statistics. Because the commit is atomic, readers see either the old file set or the new one, never a mixture. A reader that started before the commit keeps reading the old files happily — they still exist on storage. ## The dataChange flag Every `add` and `remove` written by OPTIMIZE sets `"dataChange": false`. This is the flag that says "this commit reshuffles bytes without changing the table's rows." Structured Streaming readers use it to skip the version entirely; without it, a compaction job would look like a giant batch of new rows and every downstream stream would re-process the whole table. If you ever compact by hand — read a partition and overwrite it — you lose that property, which is one good reason to use the command rather than roll your own. ## What OPTIMIZE does not do - It does **not** delete anything from storage. Freeing space is `VACUUM`'s job, and only after the deleted-file retention window has passed. - It does **not** change logical content: no dedupe, no delete, no schema change. - It does **not** fix a bad physical layout by itself. If queries filter on a column the data is not clustered by, bin-packing gives you fewer files but not better skipping; that requires `ZORDER BY` or liquid clustering. - It does **not** repair over-partitioning. A table partitioned by a high-cardinality column still has thousands of tiny partitions after OPTIMIZE, each holding one small file, because bin-packing never merges across partition boundaries. ## Cost, scheduling and the automatic variants OPTIMIZE is pure write amplification: you pay to rewrite data that was already correct, in exchange for cheaper reads later. That trade is worth it for tables read many times and written in small increments, and not worth it for write-once tables already produced at a good file size. Schedule it against recently written partitions rather than the whole table. On Databricks two table properties reduce the need for a scheduled run: `delta.autoOptimize.optimizeWrite`, which repartitions data before writing so fewer small files are produced in the first place, and `delta.autoOptimize.autoCompact`, which runs a small compaction inline after a write when a partition has accumulated too many files (targeting a smaller file size than a full OPTIMIZE). They reduce the backlog; they do not replace periodic OPTIMIZE on a heavily written table. ## Verifying the result `DESCRIBE HISTORY tbl` shows one row per commit. For an OPTIMIZE row, `operationMetrics` carries `numRemovedFiles`, `numAddedFiles`, `totalFilesSkipped` and file-size statistics — the honest measure of whether the run did anything. A run that removed 800 files and added 6 was worth it; a run that skipped everything means the table was already in shape and your schedule is too aggressive.
- Why does OPTIMIZE not shrink the table's storage footprint?Because it only rewrites data logically: the new files are added and the old ones marked `remove` in the log, but the old Parquet files stay on object storage so historical versions and in-flight readers can still resolve them. Storage is reclaimed by `VACUUM`, which deletes unreferenced files older than the deleted-file retention window (7 days by default). Until then, an optimized table temporarily occupies more space, not less.
- A table is partitioned by user_id and OPTIMIZE leaves thousands of tiny files — why?Bin-packing never merges files across partition boundaries, and a high-cardinality partition column gives you one directory per user with only a few kilobytes in it. OPTIMIZE cannot fix over-partitioning; the fix is to repartition the table on a coarse column (or drop partitioning in favour of clustering keys) and rewrite it once.
- How do you tell whether an OPTIMIZE run actually helped?Read the `operationMetrics` on that commit in `DESCRIBE HISTORY`: `numRemovedFiles` versus `numAddedFiles` and the resulting file-size distribution, plus `totalFilesSkipped`. Pair it with query-side evidence — files scanned and task counts in the query plan before and after. A run that mostly skipped files means the table was already well shaped.
saying these in an interview costs you the question
- Says OPTIMIZE deletes the old files from storage
- Claims OPTIMIZE sorts or clusters data by itself
- Thinks OPTIMIZE removes duplicate rows or applies pending deletes
- Believes compaction rewrites files in place under the same names
- Assumes OPTIMIZE merges small files across partition directories