ClickHouse sharding is decided by one expression: the sharding key in the Distributed table. Everything that follows, how evenly writes spread, whether a tenant’s query touches one shard or all of them, how much the initiator has to merge, and how painful it is to add a shard later, is a consequence of that expression. It is also the one part of a cluster design that cannot be changed without moving data, which is why it deserves a workbook rather than a default.
This page is that workbook. It compares five shard-key strategies on the same workload, with the arithmetic for skew, the query-routing behaviour, the rebalancing cost and the failure mode of each, then covers the contention and troubleshooting patterns specific to a sharded cluster. The horizontal scaling hub covers when to shard at all; this page covers how.
The archive under this category holds the sharding introduction and resharding strategies, the Distributed-table scaling guide, shard-contention and sharding troubleshooting, shard rebalancing methods, custom partitioning keys, multi-tenant cluster design, the six-node setup, compression for scale, and the payments reference design.
The workload the five ClickHouse sharding strategies are compared on
An event table of 2 billion rows a day from 4,000 tenants, where the top tenant produces 9 percent of rows and the top ten together 38 percent (an illustrative distribution, but a common shape). Ninety percent of queries filter on one tenant and a time range; the rest are cross-tenant reports. The cluster is four shards with two replicas each. The sharding in ClickHouse, part 1 post sets out the vocabulary used below; the scaling horizontally with sharding and Distributed tables post covers the mechanics of the Distributed engine itself.
-- the number every strategy is judged against: how the rows are distributed by the candidate key
SELECT
tenant_id,
count() AS rows_30d,
round(100 * rows_30d / sum(rows_30d) OVER (), 2) AS pct
FROM events
WHERE ts > now() - INTERVAL 30 DAY
GROUP BY tenant_id
ORDER BY rows_30d DESC
LIMIT 10;
-- rule: with N shards, any single key value above 100 / N percent will overfill one shard under a co-locating keyStrategy 1: rand(), the ClickHouse sharding key with perfect spread and no routing
rand() as the ClickHouse sharding key spreads every insert evenly across shards and makes skew impossible. The cost is that no query can be routed: every tenant query fans out to all four shards, each shard reads its quarter of the tenant’s rows, and the initiator merges four partial results. For small clusters serving cross-tenant reports this is fine; for the tenant-scoped 90 percent it multiplies the work by the shard count and puts the merge on one node. Adding a shard later needs no rebalancing for correctness, only for capacity.
CREATE TABLE events_dist ON CLUSTER 'ch_prod' AS events
ENGINE = Distributed('ch_prod', currentDatabase(), 'events', rand());
-- every query: 4 shards read, 4 partial results merged on the initiator
-- EXPLAIN shows ReadFromRemote on every shard regardless of the WHERE clauseStrategy 2: cityHash64(tenant_id), co-location with a skew risk
Hashing the tenant identifier is the default ClickHouse sharding choice: it puts all of a tenant’s rows on one shard, so a tenant-scoped query with optimize_skip_unused_shards = 1 runs on one shard with no merge. On the workload above, the 9 percent tenant lands on one shard and that shard carries 9 percent plus its share of the rest, roughly 32 percent of all rows against 25 percent for an even split, which is tolerable at four shards and becomes a hot shard at sixteen.
The shard contention troubleshooting post covers what a hot shard looks like in system.metrics and the query log.
CREATE TABLE events_dist ON CLUSTER 'ch_prod' AS events
ENGINE = Distributed('ch_prod', currentDatabase(), 'events', cityHash64(tenant_id));
SET optimize_skip_unused_shards = 1;
SELECT count() FROM events_dist WHERE tenant_id = 42 AND ts >= today() - 7;
-- EXPLAIN: ReadFromRemote on one shard only
-- projected load per shard before committing to the key
SELECT cityHash64(tenant_id) % 4 AS shard, count() AS rows, round(100 * rows / sum(rows) OVER (), 1) AS pct
FROM events WHERE ts > now() - INTERVAL 30 DAY
GROUP BY shard ORDER BY shard;Strategy 3: composite key, co-location within a time window
cityHash64(tenant_id, toYYYYMMDD(ts)) co-locates a tenant’s rows per day rather than forever, so a large tenant’s history spreads across shards while any one day of it sits on one. Tenant queries over a week fan out to at most seven shard-days, which on four shards is usually all four, so routing is lost for most queries; what is gained is that no single tenant can overfill a shard.
It suits workloads where the cross-tenant reports matter as much as the tenant queries, and where skew is severe. The custom partitioning keys post covers the interaction between the shard key and the partition key, which are chosen together.
ENGINE = Distributed('ch_prod', currentDatabase(), 'events', cityHash64(tenant_id, toYYYYMMDD(ts)))
-- a one-day tenant query routes to one shard; a 30-day one fans out; skew bounded by the largest tenant-dayStrategy 4: explicit ClickHouse sharding with a tenant-to-shard map
For the heaviest tenants, an explicit map decides the shard: a dictionary from tenant_id to shard number, consulted in the sharding expression, with unmapped tenants falling back to a hash. The operator places the 9 percent tenant alone on shard 1 and packs smaller tenants around it, and can move a tenant by changing the map and copying its partitions. Routing is preserved for every tenant query, skew is managed by hand, and the cost is that someone has to manage it; the building multi-tenant ClickHouse clusters post sets out the map, the quotas per tenant and the isolation options.
CREATE DICTIONARY tenant_shard_map
(
tenant_id UInt32,
shard UInt8
)
PRIMARY KEY tenant_id
SOURCE(CLICKHOUSE(table 'tenant_shard_assignments'))
LAYOUT(FLAT())
LIFETIME(MIN 60 MAX 120);
CREATE TABLE events_dist ON CLUSTER 'ch_prod' AS events
ENGINE = Distributed('ch_prod', currentDatabase(), 'events',
if(dictHas('tenant_shard_map', tenant_id),
toUInt64(dictGet('tenant_shard_map', 'shard', tenant_id)),
cityHash64(tenant_id)));
-- optimize_skip_unused_shards needs a deterministic key: the dictionary must be loaded on every nodeStrategy 5: one shard per tenant class, isolation over balance
Some estates apply ClickHouse sharding by contract rather than by data: enterprise tenants on their own shards with dedicated replicas, everyone else on a shared pool, with the Distributed table’s cluster definition carrying different weights or a separate cluster per class. Balance is deliberately uneven; what is bought is isolation, so that one class’s load cannot affect another’s latency, and a per-class SLO that can be signed. The real-time payments analytics post shows this shape for a payments estate where the regulated tenants sit apart.
| Strategy | Tenant query | Skew on the workload (4 shards) | Add a shard | Fails when |
|---|---|---|---|---|
| 1. rand() | all shards, merge on initiator | none: 25 / 25 / 25 / 25 | no data move needed | tenant queries dominate |
| 2. hash(tenant) | one shard, no merge | 32 / 23 / 23 / 22 | rebalance affected tenants | one tenant > 100 / N percent |
| 3. hash(tenant, day) | one shard per day queried | ≈ 26 / 25 / 25 / 24 | rebalance recent days only | multi-day tenant queries are the norm |
| 4. explicit map | one shard, no merge | as placed, e.g. 28 / 24 / 24 / 24 | move mapped tenants by partition | no one owns the map |
| 5. per-class shards | one shard or one pool | uneven by design | per class | classes are not stable |

Reading the plan: did ClickHouse sharding route the query or fan it out?
The evidence for every strategy is the same: EXPLAIN on the Distributed table shows one ReadFromRemote per shard that will be read, and the per-shard query log, joined on initial_query_id, shows where the time went. A query that should route to one shard and shows four remote reads has a non-deterministic key, a predicate the planner cannot evaluate against the key, or optimize_skip_unused_shards off. The sharding troubleshooting and performance optimisation post is the reference for the cases where the plan disagrees with the intent.
EXPLAIN SELECT count() FROM events_dist WHERE tenant_id = 42 AND ts >= today() - 7;
-- routed: one "ReadFromRemote (Read from remote replica)" step
-- fanned out: one per shard, followed by a MergingAggregated step on the initiator
-- where the time went, per shard, for one query
SELECT hostName() AS host, is_initial_query, query_duration_ms, read_rows, formatReadableSize(memory_usage) AS mem
FROM clusterAllReplicas('ch_prod', system.query_log)
WHERE initial_query_id = '${QUERY_ID}' AND type = 'QueryFinish'
ORDER BY is_initial_query DESC, host;
-- subqueries: evaluated once with GLOBAL, or once per shard without
SET distributed_product_mode = 'global';
SELECT count() FROM events_dist WHERE user_id GLOBAL IN (SELECT user_id FROM flagged_users);Contention on one shard: the four causes in ClickHouse sharding
In ClickHouse sharding, a shard that is slower than its peers has one of four problems: it holds more data (skew, strategies 2 and 4), it receives more inserts (a producer pinned to one node rather than round-robin), it runs more merges (small inserts landing there), or it is answering the initiator role for every distributed query because the load balancer sends all client connections to it.
The shard contention troubleshooting guide works through each with the metric that separates them; the fourth is the one most often missed, and the fix is to spread the initiator role across replicas.
-- one row per node: the four contention signals side by side
SELECT
hostName() AS host,
(SELECT sum(bytes_on_disk) FROM system.parts WHERE active AND table = 'events') AS bytes_events,
(SELECT count() FROM system.query_log WHERE query_kind = 'Insert' AND event_time > now() - INTERVAL 1 HOUR) AS inserts_1h,
(SELECT count() FROM system.part_log WHERE event_type = 'MergeParts' AND event_time > now() - INTERVAL 1 HOUR) AS merges_1h,
(SELECT count() FROM system.query_log WHERE is_initial_query AND query_kind = 'Select' AND event_time > now() - INTERVAL 1 HOUR) AS initiator_1h
FROM clusterAllReplicas('ch_prod', system.one)
ORDER BY host;Resharding: what it costs to change a ClickHouse sharding key
Changing the key, or adding a shard under strategies 2 to 4, means moving rows, and there is no online operation that does it. The production method is per partition: copy the rows that now belong elsewhere with INSERT … SELECT FROM remote(), verify counts and a checksum on both sides, and only then delete from the source behind a confirmation gate.
The shard rebalancing methods post compares this with the experimental MOVE PARTITION TO SHARD, the deprecated clickhouse-copier, and a full reload; the sharding and resharding strategies post covers planning the change so that it runs partition by partition over days rather than as one weekend.
-- resharding one partition from 4 to 5 shards, rows that now hash to the new shard
INSERT INTO events
SELECT * FROM remote('ch-s1-r1:9000', currentDatabase(), 'events', '${CH_USER}', '${CH_PASSWORD}')
WHERE toYYYYMM(ts) = 202608 AND cityHash64(tenant_id) % 5 = 4;
-- verify both sides before any delete
SELECT count(), sum(cityHash64(user_id, ts)) FROM events WHERE toYYYYMM(ts) = 202608 AND cityHash64(tenant_id) % 5 = 4;
-- CONFIRMATION GATE: counts and checksums match and are recorded in the change ticket
-- ALTER TABLE events DELETE WHERE toYYYYMM(ts) = 202608 AND cityHash64(tenant_id) % 5 = 4; -- on the source shard onlyClickHouse sharding key and partition key: two decisions made together
The shard key decides which node holds a row; the partition key decides which part on that node. They interact in two places. Rebalancing moves partitions, so a partition expression aligned with the shard key’s time component (daily partitions under strategy 3) makes moves cheap.
A partition key that includes the tenant under strategy 2 produces one partition per tenant per month on the shard, which for a shard holding a thousand tenants is a part explosion. The partition hub covers the design rules; the short version for a sharded cluster is monthly or daily partitions by time and nothing else.
Writing through the Distributed table, or directly to shards
The sharding key is applied on insert, and there are two ways to apply it. Inserting through the Distributed table lets the initiator split each batch by key and ship the pieces asynchronously, which is simple and doubles network traffic; inserting directly into the local table on the right shard, with the producer computing the key itself, halves the traffic and removes the distribution queue as a failure point, at the cost of producers that know the topology.
Large estates use direct inserts for the firehose and the Distributed table for everything else. Either way, batch size still decides part count: a Distributed insert of 1,000 rows across four shards produces four parts of 250.
-- through the Distributed table: split by key, shipped async (default) or synchronously
SET distributed_foreground_insert = 1; -- fail fast instead of queueing when a shard is down
INSERT INTO events_dist SELECT * FROM input('ts DateTime64(3), tenant_id UInt32, event_type String, user_id UInt64') FORMAT JSONEachRow;
-- direct: the producer computes cityHash64(tenant_id) % 4 and connects to that shard's replica
INSERT INTO events FORMAT JSONEachRow -- on ch-s3-r1, for rows whose key maps to shard 3
-- the queue that exists only for the first path
SELECT database, table, data_files, formatReadableSize(data_compressed_bytes) AS pending, is_blocked
FROM system.distribution_queue;Compression and the shard count
The shard count in the capacity plan is a function of compressed bytes, so codec choices move it directly: a table that compresses 12× rather than 6× needs half the shards for the same retention. The data compression for performance and scalability post measures the common cases, and the fast data loops post covers the transformation tables that grow alongside the sharded fact table. Run the compression review before the shard-count decision, not after.
The strategy table above is the deliverable of the shard-key decision: one row per candidate, filled in from the customer’s own distribution query, with the chosen row and the reason recorded in the design document before any Distributed table is created.
Version notes
optimize_skip_unused_shards and distributed_product_mode have been available since 19.x. MOVE PARTITION TO SHARD remains experimental; clickhouse-copier was deprecated in 24.x. Parallel replicas (production-ready from 24.x) change the calculus for large scans by letting replicas, not only shards, contribute bandwidth, which is why some estates stop at fewer shards than they once would have. The performance settings in 26.8 LTS post lists the current defaults; the Distributed engine documentation is the source for the sharding-expression rules. Confirm on the running version.
Reading the archive
Start with the ClickHouse sharding introduction (part 1) and the Distributed-table scaling guide for the mechanism, then the multi-tenant clusters post for strategies 2, 4 and 5 in practice. Read shard contention troubleshooting and sharding troubleshooting before the first incident rather than during it. The rebalancing and resharding posts cover the change procedure; the six-node setup gives the reference configuration; custom partitioning keys and compression cover the two decisions made alongside the key.
ChistaDATA runs the shard-key decision as a short ClickHouse consulting engagement from the customer’s own row distribution and query log, and carries out resharding under managed services partition by partition with the verification gates above. Every strategy on this page should be tested on staging with production-shaped data, and no resharding step runs without a tested backup of the source shard.