Rebuilding a Table
Some changes ClickHouse cannot make to the table it has. The sorting key is the on-disk order and, in a ReplacingMergeTree, the dedup identity, so changing it means writing every row again into a new table. The classifier puts these in the rebuild class, and both chant sql plan and the applier refuse them in place. ClickHouseRebuildOp is what runs them.
Start from the plan
Section titled “Start from the plan”Change the declaration in a pull request. chant sql plan <env> dist/schema.json (or chant sql diff between the base and head builds) refuses the change, exits 2, and names the Op to run for each refused table:
shop.events (shop.events) [REBUILD] orderBy: ( ts , user_id ) -> ( user_id , ts ) SQLCH220 Change the sorting key. ...
Refused: 1 change(s) need a rebuild, which ClickHouse cannot make to the existing table. ...
Run it as the rebuild migration Op, declared in an *.op.ts file (import { ClickHouseRebuildOp } from "@intentius/chant-lexicon-sql/clickhouse"), then `chant run <name>` until it is done: export const { op } = ClickHouseRebuildOp({ name: "rebuild-shop-events", env: "prod", table: "shop.events", dualWrite: { mode: "materialized-view", cutoverColumn: "ts" } });--json carries the same as rebuildOps, and the applier’s not-attempted detail names ClickHouseRebuildOp({ table: "shop.events", ... }). The suggested dualWrite uses the table’s first time column when it has one, and app mode otherwise. Put the declaration in the pull request beside the schema change:
import { ClickHouseRebuildOp } from "@intentius/chant-lexicon-sql/clickhouse";
export const { op } = ClickHouseRebuildOp({ name: "rebuild-shop-events", env: "prod", // sql.profiles.prod table: "shop.events", dualWrite: { mode: "materialized-view", cutoverColumn: "ts" }, retain: "7d",});The Op’s own Plan phase classifies the table’s change again, against the server, when it runs. A table whose change is not a rebuild is refused with the rules that apply in place, and the applier is the way to make them.
The phases
Section titled “The phases”| Phase | What it does |
|---|---|
| Build | chant build (the project’s build script); build: false skips it |
| Plan | classifies the change, reports how far the rebuild has got |
| Create | the new table, events__chant_new, from the declaration, once the old table’s unfinished mutations are done |
| Dual write | a materialized view from the old table to the new one, or a gate for the application to stop writing |
| Backfill | INSERT ... SELECT per partition of the old table, each with a receipt |
| Verify | row counts and checksums per partition, every row of the old table against the new |
| Approve | a gate bound to the plan and the verification |
| Swap | compare the tables again, then EXCHANGE TABLES and recreate the materialized views that read the table |
| Retain | the old table stays as events__chant_old until retain (default 7d) has passed |
| Approve drop | a gate bound to that old table |
| Drop | the old table, once its date has passed |
| onFailure | drops the new table and the dual-write view |
Run it with chant run rebuild-shop-events. Each run goes as far as the next gate and exits 3 there. Every step reads the server again before it acts, so the next run, after an approval or a crash, carries on from where the last one stopped and repeats nothing that is done.
The objects the rebuild makes carry chant’s ownership marker with two more pairs, rebuild=shop.events and the object’s role (new, dual or old). A plan, an import and a prune leave them out, so an ApplyOp running beside a rebuild neither proposes dropping them nor prunes them.
Keeping up with writes
Section titled “Keeping up with writes”dualWrite decides how rows written during the rebuild reach the new table.
{ mode: "materialized-view", cutoverColumn: "ts" } keeps writes flowing. The Op creates events__chant_dual, a materialized view on the old table that writes every row whose ts is at or after a cut-over into the new table. The cut-over is the server’s time when the view is created plus cutoverDelay (default 5s), rounded up to a whole second. The backfill copies the rows before it once no more of them can arrive: the server’s clock has passed the cut-over, and every INSERT into the old table that began before it has finished, asynchronous inserts still waiting in their buffer included. It polls both, so on a quiet table it starts a few seconds after the view is made. A write still running cutoverTimeout (default 10m) after the cut-over stops the backfill with its query id. The cut-over column has to be a time that rows arrive in order of, give or take the delay, so a writer that holds rows in batches before inserting them needs a cutoverDelay longer than its batches. A row that arrives later than that, with a time before the cut-over, is in neither half; the verification finds it and the run fails rather than swap a table missing it. ClickHouse’s backfilling guide describes the same split.
The view sees only rows inserted after it was created. Rows the table already held with a ts at or after the cut-over (a booking next month, an expiry date) are the backfill’s too: for each partition it copies the rows at or after the cut-over that the new table does not have yet, row for row, so those the view wrote are not copied twice. A row the view was still committing when the backfill read both tables can end up copied by both; the backfill looks for rows the new table holds more often than the old one, deletes them and copies them again.
{ mode: "app" } is for a table with no such column. The Dual write phase is a gate, <name>-writes-stopped, and approving it says the application has stopped writing to the old table (paused, or buffering). The backfill then copies every row, and writes resume against the table’s own name after the swap, where the new table now is.
The gates
Section titled “The gates”The Verify phase compares every row of the old table with the new one, per partition: the rows before the cut-over, those the old table held after it, and those the view has written since. The old table is read first, and a difference is read again up to five times, a second apart, before it fails the run, so a row the view is still committing is not taken for a missing one. The Swap phase runs the same comparison once more just before the EXCHANGE, and fails rather than swap when they differ, so a row that reached the old table alone after the verification is not left behind in events__chant_old.
A table whose engine collapses rows that share a sorting key when parts merge (SummingMergeTree, ReplacingMergeTree, AggregatingMergeTree, CollapsingMergeTree, VersionedCollapsingMergeTree, CoalescingMergeTree, GraphiteMergeTree, and their Replicated forms), rebuilt into one of them, is compared under FINAL. The old table can hold unmerged parts with several rows for one key; the copy inserts them in one block, and the new table collapses them as it writes, so their raw counts differ while the merged rows are the same. Under FINAL both tables are read as merged, and the Verification outcome says so. FINAL reads cost more than a plain count on a large table. If the tables still differ, the new sorting key collapses the copied rows differently from how the old table merges them; the error names the engine. Run OPTIMIZE TABLE <db>.<table> FINAL and run the rebuild again, so the copy reads merged rows. The plan’s hand-off to the Op adds the same note for such a table.
A table whose engine keeps every row (a plain MergeTree) rebuilt into one that collapses them is compared by what the new engine keeps per sorting key. The new table holds one row per key where the old one holds them all, so their row counts differ by design. Both tables are grouped by the new table’s sorting key and partition, and each key is compared on what a merge of the new engine leaves unchanged:
| New engine | Compared per key |
|---|---|
SummingMergeTree | the sums of its summed columns (those listed, else every numeric column outside the sorting and partition keys), in their own types; a key whose sums are all zero is left out |
ReplacingMergeTree | the highest version, or the key alone without a version column |
CollapsingMergeTree, VersionedCollapsingMergeTree | the sum of the sign column; a key whose signs cancel out is left out |
AggregatingMergeTree | sum, min, max and the groupBit functions of its SimpleAggregateFunction columns |
CoalescingMergeTree | which of its Nullable columns hold a value |
Each side is read without FINAL, so it does not matter which parts have merged. A key missing from the new table, or one whose sums differ, still fails the verification, and the Verification outcome counts the old table’s rows and the keys they make. A GraphiteMergeTree, whose rows depend on their age, is still compared row for row.
The swap gate’s approval is for one plan and one verification. Its digest covers the definitions on both sides, the classified changes, the columns copied, the dual-write mode and, per partition, the row count and checksum of the rows before the cut-over (every row in app mode). Rows at or after the cut-over keep arriving through the view, so they are compared but not bound: an approval stays valid while writes go on, and a run that verifies different numbers before the cut-over needs a fresh approval. The run record carries the verification as outcomes (Verification, VerifiedPartitions, VerifiedRows), which chant run status and chant operator log show:
chant run rebuild-shop-events # gated: approve-rebuild-shop-eventschant approve rebuild-shop-events approve-rebuild-shop-events --actor alexchant run rebuild-shop-events # swaps; gated: approve-rebuild-shop-events-dropchant approve rebuild-shop-events approve-rebuild-shop-events-drop --actor alexchant run rebuild-shop-events # drops the old table once retain has passedgate.approval takes a quorum, roles and a Cedar policy as any gate does; with a policy, the verified partition and row counts are added to its context. The drop gate binds the old table by its UUID and retention date. Approving it early is allowed: the Drop phase drops nothing until the date has passed and says so, and a run after the date drops it.
Once the swap has run, the swap gate has nothing left to approve, and later runs pass it on the approval already recorded.
Receipts and resuming
Section titled “Receipts and resuming”Each partition’s copy is an effect with a receipt, in the read, compare, run, write cycle effect() uses: a partition whose receipt matches is skipped, and a partition’s receipt is written after its copy succeeded, last. The receipts are rows on the same server:
SELECT address, expectation, run_id, written_at FROM chant_receipts.receiptsWHERE address LIKE 'shop/prod/rebuild/shop.events/%'In a Replicated database they are in shop.__chant_receipts instead, beside the tables, so every replica has them.
They live there, not on the chant/lifecycle branch, because a receipt has to go when the data it witnesses goes. onFailure drops the new table, and each receipt’s expectation is bound to the new table’s UUID, so a new table made again finds every old receipt stale and copies every partition again. A receipt in git would go on saying those partitions were copied.
A run stopped during the backfill (Ctrl-C, a killed job, a lost machine) is not a failure: it runs no onFailure, and the next run resumes at the first partition without a receipt. A run killed between a partition’s INSERT and its receipt leaves rows with no receipt. The next run finds them, deletes them (waiting for that mutation in system.mutations) and copies the partition once more, so no row is in the new table twice. Each partition’s copy runs under a query id made from the new table and the partition, and a copy the server is still running from a run whose client went away is killed first.
A step that fails is a failure. The backfill step retries three times, and the Verify and Swap steps fail on any difference. Unless the Op keeps its new table on failure (onFailure: "keep"), onFailure then drops events__chant_new and events__chant_dual, only when their comments name this rebuild and carry this project’s marker, and the next run starts again from a new table. After the EXCHANGE, onFailure touches nothing: the next run finishes the swap.
Inside a larger approved change
Section titled “Inside a larger approved change”A tool that runs the rebuild as one step of a change it has already had approved (a migration runner whose own gate covers the whole change, say) can set two options so it does not have to edit the Op’s config.
export const { op } = ClickHouseRebuildOp({ name: "rebuild-shop-events", env: "prod", table: "shop.events", dualWrite: { mode: "materialized-view", cutoverColumn: "ts" }, gates: "outer", // the caller's approval covers the swap onFailure: "keep", // a failed run resumes from its receipts});gates: "outer" leaves out the swap gate, the drop gate and the Drop phase, so a run goes from Plan through Retain without stopping. The Verify phase still compares every row, and the Swap phase still compares them again before the EXCHANGE: any difference fails the run with nothing swapped. The old table is kept as events__chant_old with its retention date. Drop it after that date, or run the Op with its own gates (gates: "own", the default): its swap gate has nothing left to swap, and its drop gate binds that old table. In app mode the <name>-writes-stopped gate stays, because it confirms that the application stopped writing and does not approve the swap.
onFailure: "keep" leaves out onFailure. A step that fails leaves events__chant_new, events__chant_dual and the partitions’ receipts, and the next run resumes the backfill at the first partition without a receipt, as it does after a run that was stopped. A failure that a rerun would hit again (a verification difference, or a new table made from an earlier declaration) says so in its error. To start again from a new table, drop events__chant_new and events__chant_dual.
Background rewrites
Section titled “Background rewrites”The Create phase waits until system.mutations lists nothing unfinished for the old table, so the copy reads what a pending type change or TTL change produces, and the Verify phase waits on both tables. The wait is ten minutes by default (mutationTimeout); a mutation still running then fails the step with its id, and one the server reports failing fails it at once with the server’s reason.
Materialized views that read the table
Section titled “Materialized views that read the table”In ClickHouse 26.8 a materialized view follows its source table by name, so after the EXCHANGE the views that read shop.events already read the new table. The Swap phase detaches and attaches each of them, which makes the server analyze its query against the new table’s columns and types: a view the new definition breaks fails the swap there, with the old table still kept, instead of on the next insert. A view with an inner table keeps its data, which dropping and creating it would not. Inserts that land while a view is detached are not passed to it, so in materialized-view mode, where writes go on through the swap, that window is one DETACH and one ATTACH long. The Swap phase reports the views as the Dependents outcome.
On a cluster of shards
Section titled “On a cluster of shards”With topology: "cluster:<name>" on the profile, the new table, the dual-write view, the swap and the drop run ON CLUSTER, so every server gets them. When the cluster has more than one shard, each shard holds its own rows, and the backfill and the verification work shard by shard. All of it goes through the profile’s server:
- A shard’s rows are read with
cluster('<name>', shop, events)and_shard_num = <n>, one replica per shard. They are written back to the same shard withINSERT INTO FUNCTION cluster('<name>', shop, events__chant_new, <key>), where the constant sharding key lands on that shard. The insert runs in the foreground, so the statement returns once the shard has the rows. - The unit of work is a partition on one shard, and each has its own receipt,
<stack>/<env>/rebuild/shop.events/shard<n>/<partition>. A rerun after a failure on one shard skips every partition that any shard already copied. - Rows that a cut-off copy left on one shard are deleted with an
ALTER TABLE ... ON CLUSTER ... DELETEwhose condition,getMacro('shard') = '<macro>', holds on that shard’s servers only. Each shard needs its own{shard}macro, which the cluster’s Keeper paths already use. The Op refuses a cluster without one. - The verification compares every row of every shard, partition by partition per shard (
shard<n>/<partition>). A shard that the copy missed is therefore a difference, and nothing is swapped.
With one shard nothing changes: its replicas replicate what is written on any of them.
On a Replicated database
Section titled “On a Replicated database”A database on the Replicated engine (the one ClickHouse’s Kubernetes operator makes) is Atomic on each replica, and it runs every DDL statement on all of them through Keeper. The rebuild runs there as it does on one server: the new table, the dual-write view, the EXCHANGE TABLES, the views’ DETACH and ATTACH and the drop each reach every replica, and chant run can talk to any one of them, or to a different one on each run. The e2e test runs it on two replicas of the pinned server with Keeper, moving between them from run to run.
The tables have to replicate their rows. Declare them ReplicatedMergeTree (or another Replicated*MergeTree) with no arguments, which a Replicated database fills in as '/clickhouse/tables/{uuid}/{shard}', '{replica}'; the plan reads that back as the same engine. Each table then gets its own Keeper path from its UUID, so the new table never shares the old one’s.
What changes on a Replicated database:
- The receipts are rows of
shop.__chant_receipts, aReplicatedReplacingMergeTreein the same database, so a backfill resumed on another replica reads the receipts the first one wrote. Its comment carries chant’s marker, and plan, import and prune leave it out. - Before the backfill lists the old table’s partitions, before it counts what a cut-off copy left in the new table, before it reads the receipts, and before the verification reads either table, the replica waits until it has fetched what the others wrote (
SYSTEM SYNC REPLICA ... LIGHTWEIGHT). - The kill of a copy still running under its query id is
KILL QUERY ON CLUSTER 'shop', the database’s own cluster, because that copy may be on another replica than the one the resumed run talks to. - Every copy runs with
insert_deduplicate = 0. A replicated table remembers the blocks it was given, and a partition copied again after its rows were cleared would otherwise be dropped as a repeat. - A dependent view is detached with
DETACH TABLE ... PERMANENTLY, the onlyDETACHa Replicated database takes, then attached again.
A replica that is down
Section titled “A replica that is down”The rebuild goes on with the replicas that are up, and the one that is down catches up from Keeper when it comes back. It stops only when going on would lose rows.
DDL goes on without the down replica. The database writes each statement to its log in Keeper, the replica chant run talks to executes it, and every other replica executes it from the log, retrying until it succeeds (Replicated database engine). The Op runs every statement (the new table, the dual-write view, the receipts table, EXCHANGE TABLES, DROP VIEW, RENAME, MODIFY COMMENT, the views’ DETACH ... PERMANENTLY and ATTACH, the drops in the Drop phase and in onFailure, and KILL QUERY ON CLUSTER) with distributed_ddl_output_mode = throw_only_active, which “doesn’t wait for inactive replicas of the Replicated database” and still fails the statement when a live replica fails it. Under the server’s default, throw, each statement would wait the whole distributed_ddl_task_timeout (180 seconds) for the down replica and then fail the step, after the live replicas had run it. A replica counts as inactive once its Keeper session has expired, about 30 seconds after it went down with the default session_timeout_ms, so the first statement after that waits up to that long and then goes on. When the replica comes back, it runs the queued statements in order and fetches the new table’s parts from the others. A copy still running on the down replica needs no kill: a stopped server runs nothing, and a server cut off from Keeper cannot commit a part to a replicated table.
Reading rows waits, then stops. A part written on a replica that went down before the others fetched it exists on that replica alone. A backfill that went on without it would copy the table less those rows, the verification would compare two tables that both lack them and match, and the swap would lose them. So each SYSTEM SYNC REPLICA is bounded: it waits up to replicaTimeout (default 2m) for this replica to fetch what is in its queue (SYSTEM SYNC REPLICA on its own waits the session’s receive_timeout). When the down replica holds nothing the others lack, which is the usual case since parts are fetched within seconds, the wait returns at once and the step goes on. Otherwise the step fails with the parts still to fetch and where they are:
shop.events: waited 120s for replica r1 to fetch what the other replicas wrote, and 1 part(s) are still to fetch: 202602_3_3_0 from r2 (inactive).Those rows are on r2 alone, so the step stops rather than read the table without them; nothing has been swapped.Bring r2 back (it rejoins through Keeper and this replica fetches the parts), then run again.The step’s retries wait again, and once they are spent the run fails and onFailure drops the new table and the dual-write view, as for any failed step. The old table is untouched. Bring the replica back and run the Op again; it starts from a new table. Keep replicaTimeout well under the step’s timeout (five minutes for the Verify step, which waits for two tables).
Two things the Op cannot go on without: the replica chant run talks to (point CLICKHOUSE_URL, or the profile, at a live one), and Keeper. Without Keeper a replicated table takes no inserts and a Replicated database runs no DDL.
The e2e test covers both cases on two replicas. With rows on the second replica alone when it stops, the backfill stops naming it and onFailure drops the new table on the first replica without waiting; restarted, the second replica hands its rows over. Stopped in the middle of a backfill, the run goes on with the first replica through the swap and the drop, and the second, restarted, holds the declaration, every row and a dependent view that works, and plans no change.
What the Op refuses
Section titled “What the Op refuses”- A database that is neither
AtomicnorReplicated.EXCHANGE TABLESneeds the Atomic engine, and swapping with twoRENAMEs would leave a moment with no table for writes to land in. Convert the database to Atomic first. - In a Replicated database, a table that is a plain
MergeTree(or any engine that is notReplicated*MergeTree), or a declaration that makes it one. Such a table keeps a different set of rows on each replica: the backfill would copy the rows of the replica it runs on, and theEXCHANGEon every replica would swap in a table without the others’ rows. Rebuild it to a replicated engine, or move it to an Atomic database first. - A declaration whose
Replicated*MergeTreenames the same explicit Keeper path as the table has, without{uuid}in it. The new table would get that path too. Put{uuid}in the path, or leave the engine’s arguments out. - A table the server does not have, or a declaration that is not a table.
- A change with no rebuild in it. The applier makes those.
- An object under one of the rebuild’s names (
__chant_new,__chant_dual,__chant_old) that does not carry this rebuild’s marker. The Op stops rather than fill or drop somebody else’s table. - A new rebuild of a table whose last rebuild still retains its old table. Drop the old table first.