In Iceberg, what does ALTER TABLE ... WRITE ORDERED BY change about later writes?
answer
- a policy the table carries, not a rewrite
- writers are asked, readers are not told
- tight per-file bounds are the payoff
- local order is cheap, global order costs a shuffle
basics
~10 sIt records a sort order in table metadata and makes it the table's default, so engines writing afterwards sort rows to match it. Existing files are untouched and readers do not enforce it.
solid answer
~40 sThe statement adds an entry to the table's `sort-orders` list and moves `default-sort-order-id` to it; the order is a list of columns with direction and null ordering. It is a declaration for writers, not a guarantee readers enforce: an engine with the Iceberg extensions will sort the rows it writes to match, so each new data file holds a contiguous slice of the sort columns and its recorded lower and upper bounds become tight, which is what makes later queries skip files. `WRITE LOCALLY ORDERED BY` sorts within each write task without changing how rows are distributed across tasks, and the distribution itself is controlled separately by `write.distribution-mode`. Nothing existing is reorganised — applying the order to data already written requires rewriting those files through the table's maintenance procedures.
code
sql · 12 lines-- becomes the table's default order for future writes
ALTER TABLE prod.db.events WRITE ORDERED BY user_id, event_ts;
-- sort within each write task only, no extra shuffle
ALTER TABLE prod.db.events WRITE LOCALLY ORDERED BY user_id;
-- set distribution and local order together
ALTER TABLE prod.db.events
WRITE DISTRIBUTED BY PARTITION LOCALLY ORDERED BY user_id;
-- clear it
ALTER TABLE prod.db.events WRITE UNORDERED;go deeper
Know that Iceberg lets a table declare how new data should be sorted, that the statement is instant, and that it changes future writes rather than reorganising files that already exist.
Explain the payoff: sorted files get narrow lower and upper bounds per column in their manifest entries, and those bounds are what let a query skip files that partition pruning cannot eliminate.
Show the cost side. A global order needs a range shuffle and can introduce write skew; a local order is cheap but leaves overlapping files. Be ready to explain why a declared order changed nothing when the writer ignored it or the files predate it.
Treat write order and distribution as platform policy carried by the table rather than settings copied into every job, and be prepared to justify the ongoing rewrite budget needed to keep historical data matching that policy.
## What is actually stored Iceberg table metadata carries a list of sort orders alongside the schemas and partition specs, plus a `default-sort-order-id` naming the current one. A sort order is an ordered list of fields, each with a source column id, an optional transform, a direction and a null ordering. `ALTER TABLE prod.db.events WRITE ORDERED BY user_id, event_ts` appends such an order and makes it the default; `ALTER TABLE ... WRITE UNORDERED` sets the table back to no order. Like partition spec evolution, this is metadata-only and returns immediately on a table of any size. ```sql ALTER TABLE prod.db.events WRITE ORDERED BY user_id, event_ts; ALTER TABLE prod.db.events WRITE ORDERED BY region ASC NULLS LAST, event_ts DESC; ALTER TABLE prod.db.events WRITE LOCALLY ORDERED BY user_id; ALTER TABLE prod.db.events WRITE UNORDERED; ``` ## Declaration, not enforcement The crucial framing is that a sort order is advisory. Nothing in the format checks that a data file is sorted, and readers do not depend on it for correctness — a file written out of order is still a valid file with valid bounds. What the order does is tell any writer that cooperates what layout the table wants, so the table's owner sets the policy once instead of every job setting it in its own code. Engines using the Iceberg SQL extensions apply it automatically on subsequent writes; a job using a raw writer API may ignore it entirely, which is exactly the situation to watch for when the layout does not improve after the DDL. ## Local order versus distribution Two separate decisions govern how rows land in files. *Distribution* decides which write task a row goes to; *local order* decides how one task sorts the rows it has before writing them out. `WRITE LOCALLY ORDERED BY` sets only the second, so each output file is internally sorted but two files from different tasks may cover overlapping ranges. `WRITE ORDERED BY` asks for a globally meaningful order, which requires the writer to range-distribute rows so that files do not overlap. The distribution side is also exposed as the table property `write.distribution-mode`, whose values are `none`, `hash` and `range`, and the combined DDL form `WRITE DISTRIBUTED BY PARTITION LOCALLY ORDERED BY ...` sets both explicitly. The tradeoff is plain: a global order costs a shuffle and can make writes slower and more skew-prone; a local order is nearly free but leaves files overlapping, so file-level pruning is weaker. ## Why sorting helps a query at all Every manifest entry records lower and upper bounds per column for its data file. When files are sorted on a column, those bounds are narrow and disjoint, so a predicate on that column eliminates almost every file before any data is read. When files are unsorted, each one spans nearly the full range of values and the bounds exclude nothing. This is the mechanism behind the common advice to sort by the column you filter on but do not partition by — the bucketed-key case is the canonical example, where partition pruning is impossible and file bounds are the only lever left. Sort order also interacts with partitioning: within a partition, ordering by a secondary column recovers pruning that the partition field alone cannot provide. And because rows arriving in sorted order compress and encode better, sorted tables are often smaller as well as faster, though that is a secondary benefit. ## What it does not do It does not sort existing data. A table with a year of unsorted files gains nothing until those files are rewritten, and that rewrite is a maintenance operation with full read-and-write IO for the range you choose — a separate decision from setting the order, and usually applied to the hot range only. It does not sort across commits either: two writes, each internally ordered, still produce files whose ranges overlap unless the writes were range-distributed against the same boundaries. And it does not make an arbitrary query fast; it helps only predicates on the leading columns of the order, which is why the column list should be chosen from actual query patterns rather than from the schema's natural key. ## Interview framing Say what is stored (an order in table metadata plus a default id), that it is advisory to writers, that it affects future writes only, and that the payoff is tighter per-file bounds and therefore file-level pruning. Then show judgment about the cost: a global order buys pruning with a shuffle, a local order is cheap and weaker, and retrofitting old data is a rewrite you schedule rather than a side effect of the DDL.
- A team sets a write order and sees no query improvement. What are the likely causes?Either the existing files predate the change and were never rewritten, or the writing job does not use the Iceberg SQL extensions and quietly ignores the declared order, or the queries filter on a column that is not a leading field of the order. Check when the current files were written and whether their recorded bounds on the sort column are actually narrow.
- When is WRITE LOCALLY ORDERED BY the better choice than a global order?When write latency and shuffle cost matter more than maximum pruning, or when each task already receives a natural slice of the key so files barely overlap. Local ordering avoids the range shuffle and still compresses well; you accept that files from different tasks overlap, so file-level pruning is partial rather than near-perfect.
- How does a declared sort order relate to partitioning?They are complementary. Partitioning removes whole groups of files by projecting predicates through a transform; sorting narrows each file's column bounds so the remaining files can be skipped too. Sorting is the main lever for columns you filter on but deliberately did not make partition fields, such as a bucketed or high-cardinality key.
saying these in an interview costs you the question
- Says the statement sorts the table's existing files
- Claims readers verify or depend on the declared sort order
- Treats local order and global order as the same thing
- Expects a sort order to help predicates on any column
- Confuses the declared order with running a data rewrite