Monday
Hi all,
We have a Silver Delta table (~5 TB, growing ~20 GB/day) that is updated every hour with a MERGE on order_id. Right now it's partitioned by order_date. Most MERGE batches touch the last 3 days, but late updates can reach orders up to 90 days old.
We're thinking about moving to liquid clustering (CLUSTER BY (order_id) or CLUSTER BY AUTO). Questions:
Has anyone measured the MERGE performance difference after switching from date partitioning to liquid clustering on a similar table?
Is it better to cluster by the MERGE key (order_id) or by order_date, given that most reads filter by date?
What's the safest way to migrate an existing partitioned table? CREATE TABLE ... CLONE plus ALTER TABLE ... CLUSTER BY, or a full rewrite?
Any benchmarks or lessons learned would be appreciated.
Tuesday
Hi @Islam_hoti ,
Great question! We’ve faced a similar scenario with a high-volume transactional table (though ours is roughly 3TB).
One lesson learned: Check your Z-ORDER or Liquid Clustering metadata after the first few MERGE cycles. Sometimes the auto-clustering takes a few runs to settle into the optimal layout."
Monday
Hi @Islam_hoti ,
Good question, and a common one now that liquid clustering is the default recommendation for new Delta tables. For context, the docs say tables between 1 TB and 100 TB should use liquid clustering instead of partitioning, so at 5 TB and growing you're squarely in the range. I'll take your three questions in order, starting with the one that's changed most recently.
Migration path
You no longer have to choose between CLONE and a full rewrite. On Databricks Runtime 18.1 and above, you can convert a partitioned table in place:
ALTER TABLE silver.orders
REPLACE PARTITIONED BY WITH CLUSTER BY (order_date, order_id);
OPTIMIZE silver.orders FULL;
The conversion runs as a series of REORG operations plus a protocol upgrade, and batch reads and writes keep working during it (the docs recommend DBR 15.4 LTS or above for anything touching the table while it converts). If you're on an older runtime, build a new table with CLUSTER BY via CTAS, validate it, and cut over through your normal table-name change process.
On CLONE specifically: it isn't the conversion, it's the test harness. A clone carries the partition spec with it, so plain ALTER TABLE ... CLUSTER BY won't work on it (that command only accepts unpartitioned tables). What a clone is good for is an isolated benchmark:
CREATE TABLE bench_orders SHALLOW CLONE silver.orders;
ALTER TABLE bench_orders
REPLACE PARTITIONED BY WITH CLUSTER BY (order_date, order_id);
OPTIMIZE bench_orders FULL;
Two things to plan for. Existing files aren't reclustered until you run OPTIMIZE, and OPTIMIZE FULL on 5 TB can take hours, so do it once, off-peak. That also means the shallow clone is cheap to create but not cheap to benchmark, since the OPTIMIZE FULL writes a full copy of the data into the clone's location. And if you have streaming readers on the production table, they'll need a restart after conversion; the docs have a table showing the behavior with and without schema tracking.
Which keys
Cluster by both. I wouldn't pick order_id alone just because it's the MERGE key; clustering keys should reflect both what your reads filter on and what narrows the MERGE target search. Liquid clustering supports up to four keys and data skipping works on any subset of them, so CLUSTER BY (order_date, order_id) gives your date-filtered reads their skipping and gives the MERGE a way to prune the target scan on order_id. That second part matters for you specifically: with date partitioning, a batch that touches orders 90 days old has to read into 90 partitions worth of files, and within each one there's no order on order_id, so the join reads more than it needs. With order_id as a clustering key, the target-side scan can skip files whose order_id range doesn't overlap the batch.
The conversion command also sets the old partition column as a hierarchical clustering key automatically (DBR 17.1 and above), meaning data gets laid out by order_date first and order_id within that. That's the right shape for your workload. Keep order_id as a standard key rather than a hierarchical one; the docs warn that prioritizing a high-cardinality column hurts skipping on the other keys.
On CLUSTER BY AUTO: it's a good option if the table is Unity Catalog managed and you have predictive optimization on. It starts from the current partition columns and adjusts based on the observed query workload. If your read and MERGE patterns are as stable as you describe, I'd set the keys explicitly, since you already know the answer AUTO would have to learn. One practical note if you do want to test it: predictive optimization picks keys from query history, so a fresh bench clone that nobody queries won't show you much about how AUTO behaves.
The late-arriving updates
This is the part that decides how much you'll actually gain. Don't add a three-day date predicate to the MERGE unless it's guaranteed correct; a late update that falls outside it silently becomes an insert. If the source batch carries a reliable affected-date range, add that as a constraint in the ON clause so the target search can prune on order_date as well as order_id. If it doesn't, the 90-day window is your real search space, and the order_id clustering key is doing most of the work.
MERGE performance
I don't have a like-for-like benchmark to hand you, and anyone who does will have a table shaped differently enough from yours that I'd take the numbers loosely. What I can say is where the gains come from: file pruning on the order_id side of the join, deletion vectors (enabled by default with liquid clustering) that turn updates into cheap tombstones instead of full file rewrites, and row-level concurrency, which lets hourly MERGEs and background OPTIMIZE run without conflicting. Low shuffle merge, on by default since DBR 10.4, also tries to preserve the clustered layout of unmodified rows, though updated and inserted rows still need periodic OPTIMIZE to fall back into place.
The reliable way to get your own number is the bench clone above. Replay a few days of your real MERGE batches against the current table and the converted clone, on identical compute and concurrency, and compare end-to-end duration, target files scanned, bytes read, shuffle volume, files rewritten, and output file counts. Run your main date-filtered reads against both as well. Cold-cache runs and warm runs will differ, so capture both. While you're in there, confirm Photon and dynamic file pruning are on, and watch for small-file growth on the hourly cadence. That test takes an afternoon and tells you more than any forum benchmark would.
References
If you run the comparison, please post the numbers back here. Real before-and-after results on a MERGE-heavy table are exactly what this thread will get searched for.
Cheers, Louis.
Tuesday
Hi @Islam_hoti ,
Great question! We’ve faced a similar scenario with a high-volume transactional table (though ours is roughly 3TB).
One lesson learned: Check your Z-ORDER or Liquid Clustering metadata after the first few MERGE cycles. Sometimes the auto-clustering takes a few runs to settle into the optimal layout."
Tuesday
Hi! Good candidate for liquid clustering, since the "hot 3 days plus a long tail of late updates" pattern is exactly where fixed date partitions start to hurt. A few thoughts on each question.
MERGE performance
Results vary a lot by data distribution, so I'd benchmark on your own table rather than rely on someone else's numbers. The gains usually come from two places: better file skipping when the target is matched, and smaller rewrites when deletion vectors are enabled (updated rows get marked rather than whole files rewritten). An easy way to test is to clone the table, apply the new layout to the clone, replay a day's worth of hourly batches, and compare the operation metrics in DESCRIBE HISTORY, especially execution time, target files removed, and target rows copied. The last two tell you how much data each MERGE is rewriting, which is usually where the cost is.
order_id vs. order_date
It depends on how your order IDs are generated. If they increase over time, clustering by order_id naturally keeps recent orders together, so it behaves somewhat like date clustering and helps both the MERGE and date-filtered reads. If the IDs are random (UUIDs, hashes), clustering by order_id won't help date queries much. In that case, clustering by (order_date, order_id) is a reasonable compromise, though each additional key dilutes the benefit a bit.
Either way, the biggest win for MERGE is often adding a date predicate to the merge condition, for example restricting the target to the date range present in the incoming batch. That lets Delta skip most files instead of scanning the whole table for matches. This only works if order_date never changes for an order.
CLUSTER BY AUTO is worth considering if you're on Unity Catalog managed tables with predictive optimization, since it picks keys based on your actual query patterns. For a MERGE-heavy table where you already know the access pattern, explicit keys are more predictable, at least to start.
Migration
Clone plus ALTER TABLE ... CLUSTER BY won't work here: liquid clustering can't be enabled on a partitioned table, and a clone keeps the original partitioning. You'll need a rewrite. The cleanest option is CREATE OR REPLACE TABLE ... CLUSTER BY (...) AS SELECT * FROM the same table. It keeps the table's identity, permissions, and history, so you can time travel or restore if something goes wrong. After that, run OPTIMIZE so the data is fully clustered.
A few safety tips: pause the hourly MERGE job during the rewrite, test on a clone first to estimate runtime and cost at 5 TB, and check for downstream streaming readers. Replacing the table isn't an append, so streaming consumers may fail on it unless you handle that (for example, restarting them or using the option to skip change commits).
Hope that helps, and I'd be interested to hear your before/after metrics if you do run the comparison!