cancel
Showing results for 
Search instead for 
Did you mean: 
Data Engineering
Join discussions on data engineering best practices, architectures, and optimization strategies within the Databricks Community. Exchange insights and solutions with fellow data engineers.
cancel
Showing results for 
Search instead for 
Did you mean: 

Liquid clustering vs. partitioning for a 5 TB Silver table with frequent MERGEs

Islam_hoti
New Contributor III

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.

1 REPLY 1

Louis_Frolio
Databricks Employee
Databricks Employee

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.