ChistaDATA Inc.

Enterprise-class 24*7 ClickHouse Consultative Support and Managed Services

  • ChistaDATA
    • ClickHouse®
    • ClickHouse MergeTree
    • Why is ClickHouse So Fast
    • Columnar Stores
    • Vectorized Query
    • For CTOs
  • Engineering
    • Real-Time Analytics
    • Break Fix Engineering
    • Data Foundation
    • Data Archiving
    • Cloud Native ClickHouse
    • ClickHouse Consulting
      • Performance Audit
        • Pre- Engagement Questionnaire
    • ClickHouse Strategy
    • Online Ticketing System
  • Support
    • ClickHouse Migration
    • ClickHouse Audit
    • Data Warehousing Support
    • Data Analytics
    • Gen AI
    • Online Ticketing System
  • ClickHouse Managed Services
    • ClickHouse DBA
    • ClickHouse Performance
    • Data Strategy
    • ClickHouse Analytics
    • Data Archiving
    • DBaaS Optimization
    • Data SRE
    • Online Ticketing System
  • Blog
    • ChistaDATA Blog
  • University
  • Careers
  • Contact
  • Twitter
  • Facebook
  • LinkedIn
    • Shiv Iyer
  • GitHub
    • @ShivIyer
HomeClickHouse Horizontal Scaling

ClickHouse Horizontal Scaling

ClickHouse horizontal scaling is a sequence of decisions, not a single switch. A cluster grows from one node to replicas, from replicas to shards, and from shards to a tiered, object-storage-backed estate, and each step changes the failure modes, the query routing and the operational work. Adding shards too early costs months of rebalancing; adding them too late costs a cluster that cannot absorb the next quarter’s ingestion.

This page is the capacity-planning playbook behind those decisions: four stages, the measurement that says when to move to the next one, the schema and configuration that each stage needs, and the archive posts that cover the details. It is written for a team that has to sign off on the next twelve months of cluster growth.

The archive under this category covers sharding strategies and troubleshooting, shard rebalancing, read-write splits, the six-node reference setup, capacity planning, compression for scale, ingestion at velocity, and the 26.8 LTS changes to performance and high availability.

The ClickHouse horizontal scaling stages at a glance

Stage 0 is one well-sized node: enough for a surprising share of workloads, because a single ClickHouse server reads several gigabytes per second per core group and compresses event data ten-fold. Stage 1 adds replicas of that node for availability and read capacity, with no change to the data layout. Stage 2 shards the largest tables across replica sets when a single node’s storage, write rate or query CPU is exhausted. Stage 3 separates hot and cold storage and adds parallel replicas, at which point the estate is measured in hundreds of terabytes to petabytes.

The gigabytes to petabytes playbook walks the same sequence with sizing examples; this page adds the trigger metrics and the SQL.

StageTopologyWhat it scalesTrigger to move onCost of moving late
01 nodenothing yet: vertical headroomno HA; any single-node limit reachedoutage on hardware failure
11 shard × 2–3 replicasreads, availabilitydisk > 70 %, merges lag, insert CPU saturatedwrite stalls, TOO_MANY_PARTS
2N shards × 2–3 replicaswrites, storage, query CPUhot partition on one shard; storage > 200 TBrebalancing under load
3shards + object storage + parallel replicascold retention, large scansretention cost dominateslocal NVMe bought for cold data

Stage 0: exhaust vertical scaling before ClickHouse horizontal scaling

The cheapest ClickHouse horizontal scaling step is the one not taken. Before adding nodes, confirm that the single node is actually limited by hardware rather than by schema: a sort key that does not match the predicates, partitions that are too fine, String columns that should be LowCardinality, and dashboards scanning raw rows instead of rollups each produce symptoms that look like a capacity problem. The performance hub is the review to run first.

When the node is limited, the limits show as specific numbers: sustained disk utilisation above 70 percent, merge backlog growing across a day, insert threads pinned at 100 percent CPU during peaks, or query p95 that scales linearly with data volume. Any of those justifies Stage 1.

-- is it hardware or schema? three numbers before any node is added
SELECT
    formatReadableSize(sum(bytes_on_disk))                          AS on_disk,
    round(sum(data_uncompressed_bytes) / sum(data_compressed_bytes), 1) AS compression_ratio,
    count()                                                         AS active_parts
FROM system.parts
WHERE active;

SELECT metric, value
FROM system.asynchronous_metrics
WHERE metric IN ('DiskUsed_default', 'DiskTotal_default', 'OSCPUVirtualTimeMicroseconds')
   OR metric LIKE 'MaxPartCountForPartition%';

-- merges that cannot keep up show as a rising active part count over 24 h
SELECT toStartOfHour(event_time) AS h, max(CurrentMetric_PartsActive) AS parts
FROM system.metric_log
WHERE event_time > now() - INTERVAL 1 DAY
GROUP BY h ORDER BY h;

Stage 1: replicas, and the read-write split

A replica set is a group of nodes holding the same data through ReplicatedMergeTree and ClickHouse Keeper. It provides availability, because any replica can serve reads and accept writes, and it multiplies read capacity by the replica count. Writes are not multiplied: every insert is replicated to every node in the set, so the replica set has the write capacity of one node. The high availability and replication post covers the replication protocol and Keeper.

The read-write split routes inserts to one replica (or round-robin) and heavy reads to the others, usually through a load balancer or a per-role Distributed table with load_balancing and prefer_localhost_replica tuned per client. The optimal read-write split configuration post gives the routing table and settings.

-- one shard, three replicas: the Stage 1 table
CREATE TABLE events ON CLUSTER 'ch_prod'
(
    ts         DateTime64(3),
    tenant_id  UInt32,
    event_type LowCardinality(String),
    user_id    UInt64,
    payload    String
)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/events', '{replica}')
PARTITION BY toYYYYMM(ts)
ORDER BY (tenant_id, event_type, ts);

-- read routing for a reporting client: prefer the least-loaded replica, never the insert node
SET load_balancing = 'nearest_hostname';       -- or 'random', 'in_order', 'first_or_random'
SET prefer_localhost_replica = 0;
SET max_replica_delay_for_distributed_queries = 30;   -- skip replicas more than 30 s behind
SET fallback_to_stale_replicas_for_distributed_queries = 0;

Stage 2: sharding, the ClickHouse horizontal scaling step that changes everything

Sharding is the ClickHouse horizontal scaling step most teams mean by the phrase: it splits a table’s rows across replica sets so that each shard holds a fraction of the data and a fraction of the writes. Reads fan out through a Distributed table, run on every shard in parallel, and merge on the initiator. Storage and write throughput scale with the shard count; query latency for a well-pruned query stays flat; query latency for a query that touches every shard improves by the shard count minus the merge cost.

The decision that decides whether Stage 2 works is the sharding key. A key that co-locates the rows a query needs (tenant, account, device) lets optimize_skip_unused_shards send the query to one shard; a random key spreads writes perfectly and makes every query a full fan-out. The sharding strategies for high-growth environments post sets out the trade-off, and sharding in ClickHouse, part 1 covers the mechanics.

# remote_servers: 3 shards × 2 replicas (config.d/clusters.yaml; YAML config is supported since 22.x)
remote_servers:
  ch_prod:
    shard:
      - internal_replication: true
        replica:
          - {host: ch-s1-r1, port: 9000}
          - {host: ch-s1-r2, port: 9000}
      - internal_replication: true
        replica:
          - {host: ch-s2-r1, port: 9000}
          - {host: ch-s2-r2, port: 9000}
      - internal_replication: true
        replica:
          - {host: ch-s3-r1, port: 9000}
          - {host: ch-s3-r2, port: 9000}

-- the distributed table, sharded by tenant so tenant-scoped queries hit one shard
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;           -- prune shards from the WHERE clause
SET distributed_product_mode = 'global';       -- IN / JOIN subqueries evaluated once

Choosing the shard key without creating a hot shard

The failure mode of a co-locating key is skew: one tenant that produces 40 percent of all events lands on one shard, and that shard’s disk, merges and CPU run at three times the others. The hot spot detection and remediation post covers the diagnosis; the remedies are a composite key (cityHash64(tenant_id, toYYYYMMDD(ts))) for the largest tenants, a per-tenant override table that maps heavy tenants to an explicit shard, or a weight change in remote_servers so that new writes favour the emptier shards.

Skew is measured before sharding, from the single-node table, because the distribution of rows per key value is already known there.

-- skew check before choosing the shard key: share of rows held by the top keys
SELECT
    tenant_id,
    count()                                       AS rows,
    round(100 * rows / sum(rows) OVER (), 2)      AS pct_of_total
FROM events
WHERE ts > now() - INTERVAL 30 DAY
GROUP BY tenant_id
ORDER BY rows DESC
LIMIT 10;
-- rule of thumb: if the top key exceeds 100 / shard_count percent, it will overfill one shard

-- after sharding: bytes per shard from every node
SELECT hostName() AS host, formatReadableSize(sum(bytes_on_disk)) AS on_disk, sum(rows) AS rows
FROM clusterAllReplicas('ch_prod', system.parts)
WHERE active AND table = 'events'
GROUP BY host ORDER BY host;

Rebalancing: the cost of ClickHouse horizontal scaling done late

Adding a shard to a cluster does not move existing data; new inserts spread across the new shard count and old data stays where it was. Rebalancing is a manual operation, and the shard rebalancing methods post compares the four: partition-level ALTER TABLE … MOVE PARTITION TO SHARD (experimental), a copy through INSERT … SELECT FROM remote() per partition followed by a verified DROP PARTITION, the clickhouse-copier utility (deprecated since 24.x), and a full reload from the upstream source.

The per-partition copy is the method used in production because it is resumable, verifiable and reversible. Each partition is copied, row counts and checksums compared on both sides, and only then dropped from the source, behind an explicit confirmation gate in the runbook.

-- per-partition rebalance: copy, verify, then (gated) drop
INSERT INTO events SELECT *
FROM remote('ch-s1-r1:9000', currentDatabase(), 'events', '${CH_USER}', '${CH_PASSWORD}')
WHERE toYYYYMM(ts) = 202606 AND cityHash64(tenant_id) % 4 = 3;     -- rows that now belong to shard 4

-- verification on both sides before any drop
SELECT count(), sum(cityHash64(user_id, ts)) AS checksum
FROM events WHERE toYYYYMM(ts) = 202606 AND cityHash64(tenant_id) % 4 = 3;

-- CONFIRMATION GATE: counts and checksums equal on source and target, recorded in the change ticket
-- ALTER TABLE events ON CLUSTER 'ch_prod' DELETE WHERE toYYYYMM(ts) = 202606 AND cityHash64(tenant_id) % 4 = 3;

Keeper sizing: the coordination ceiling on ClickHouse horizontal scaling

Every replicated insert, merge and mutation is a Keeper transaction, so a cluster’s write rate is bounded by Keeper’s throughput long before it is bounded by disk. Three Keeper nodes on dedicated hosts (or five for multi-region) with fast local storage for the log handle tens of thousands of transactions per second; the number to watch is system.zookeeper latency and the ZooKeeperTransactions profile event per insert. Small frequent inserts are the usual cause of Keeper saturation, and they are fixed at the ingestion layer, not by adding Keeper nodes.

The six-node cluster setup post gives the full Keeper and macro configuration for a three-shard, two-replica reference cluster.

-- Keeper pressure per insert: transactions and wait time
SELECT
    quantile(0.95)(ProfileEvents['ZooKeeperTransactions'])        AS p95_zk_txn_per_insert,
    quantile(0.95)(ProfileEvents['ZooKeeperWaitMicroseconds'])/1e3 AS p95_zk_wait_ms,
    count()                                                        AS inserts_last_hour
FROM system.query_log
WHERE type = 'QueryFinish' AND query_kind = 'Insert'
  AND event_time > now() - INTERVAL 1 HOUR;

-- replication health across the cluster
SELECT hostName(), table, queue_size, inserts_in_queue, merges_in_queue, absolute_delay
FROM clusterAllReplicas('ch_prod', system.replicas)
WHERE queue_size > 50 OR absolute_delay > 60;

Stage 3: object storage, tiering and parallel replicas

Past a few hundred terabytes, most of the data is cold and the cost of keeping it on local NVMe dominates. Storage policies move parts to an S3-compatible volume after the hot window with a TTL clause, keeping the same table and the same queries; the hot tier holds the last 30 to 90 days on local disk. The storage policies and load balancing post covers the configuration.

Parallel replicas (production-ready since 24.x with enable_parallel_replicas) let a single large query use every replica of a shard rather than one, which turns replica count into scan bandwidth. Together with tiering it is the ClickHouse horizontal scaling shape for a petabyte estate: few shards, several replicas, object storage behind them. The 26.8 LTS performance and HA changes post lists what changed in this area in the latest LTS.

-- tiered storage: hot on NVMe, cold on S3 after 60 days
ALTER TABLE events ON CLUSTER 'ch_prod'
    MODIFY TTL toDateTime(ts) + INTERVAL 60 DAY TO VOLUME 'cold',
               toDateTime(ts) + INTERVAL 25 MONTH DELETE;

-- a large scan across all replicas of each shard
SET enable_parallel_replicas = 1,
    max_parallel_replicas = 3,
    cluster_for_parallel_replicas = 'ch_prod';
ClickHouse horizontal scaling diagram: four stages from a single node through replicas and shards to tiered object storage with parallel replicas, with the trigger metric and the cost of moving late at each stage
Four stages of ClickHouse horizontal scaling, the metric that triggers each move, and what it costs to move late. Thresholds are illustrative.

Capacity planning: the arithmetic behind the ClickHouse horizontal scaling plan

The plan is a spreadsheet with five inputs: daily ingested rows, bytes per row after compression, retention in days, peak query concurrency, and the growth rate. Storage per shard is rows per day multiplied by compressed bytes per row multiplied by retention, divided by shard count, with 30 percent headroom for merges and 2× for the replica. Write capacity per shard is measured, not assumed, by loading a day of production data into one node and recording rows per second at 70 percent CPU. The comprehensive guide to horizontal scaling and capacity planning post has the worked model.

Illustrative example for a payments platform: 2 billion rows per day at 40 compressed bytes per row is 80 GB per day, 29 TB per year per copy; with 13 months’ retention, two replicas and 30 percent headroom, three shards of 2 × 30 TB NVMe hold it with room for a year of 40 percent growth. The real-time payments analytics post applies the same arithmetic to a Southeast Asian payments estate.

-- the two measured inputs: compressed bytes per row, and rows per day
SELECT
    table,
    round(sum(data_compressed_bytes) / sum(rows), 1)              AS compressed_bytes_per_row,
    formatReadableQuantity(sumIf(rows, modification_time > now() - INTERVAL 1 DAY)) AS rows_last_day
FROM system.parts
WHERE active AND database = currentDatabase()
GROUP BY table
ORDER BY sum(bytes_on_disk) DESC;

Compression and materialised columns as scaling levers

Before a shard is added for storage, the compression review is run: codec choices per column (Delta and DoubleDelta for timestamps and counters, ZSTD levels for strings, LowCardinality for enumerations) routinely halve on-disk size, which halves the shard count the plan needs. The data compression for performance and scalability post and the compression hub cover the choices and their measurement.

Materialised columns work the other way: they add bytes to save CPU, by precomputing the expressions that every query evaluates. The materialized columns post shows when the trade is worth it, and the high-velocity ingestion post covers the insert-side cost.

Troubleshooting a sharded cluster

The sharding troubleshooting and performance optimisation post is the reference for the failures specific to Stage 2 and beyond: a Distributed table whose async inserts are queued on disk under /var/lib/clickhouse/data/…/events_dist/ because a shard is unreachable, initiator memory blown by a final merge of large GROUP BY states, queries that hit every shard because the predicate is not on the sharding key, and replica delay that turns a read into stale data.

Each has a system table that shows it: system.distribution_queue, system.query_log on the initiator versus initial_query_id on the shards, and system.replicas. The troubleshooting hub routes symptoms to fixes.

-- queued distributed inserts: a shard is behind or unreachable
SELECT database, table, data_files, formatReadableSize(data_compressed_bytes) AS pending, last_exception
FROM system.distribution_queue
WHERE is_blocked OR data_files > 0;

-- initiator vs shard time for one query: where does the fan-out spend it?
SELECT hostName() AS host, is_initial_query, query_duration_ms, formatReadableSize(memory_usage) AS mem, read_rows
FROM clusterAllReplicas('ch_prod', system.query_log)
WHERE initial_query_id = '${QUERY_ID}' AND type = 'QueryFinish'
ORDER BY is_initial_query DESC, host;

Version notes

ClickHouse Keeper has been the recommended coordinator since 22.x. clickhouse-copier was deprecated in 24.x; use the per-partition copy pattern above. Parallel replicas are production-ready from 24.x and further improved in 25.x and 26.x; MOVE PARTITION TO SHARD remains experimental. The Distributed table engine documentation is the source for the settings named on this page; confirm each on the running version before it goes into a change plan.

Reading the archive

Read in stage order; the ClickHouse horizontal scaling posts build on each other. For Stage 1: high availability and replication, and the read-write split guide. For Stage 2: sharding strategies, sharding part 1, the six-node setup, hot spot detection, shard rebalancing methods, and sharding troubleshooting. For Stage 3 and planning: the gigabytes-to-petabytes playbook, capacity planning, compression for scalability, materialised columns, the 26.8 LTS changes, and the payments reference design. The dbt on managed ClickHouse and fast data loops posts cover the transformation layer that grows alongside the cluster.

ChistaDATA runs the capacity review and the shard-key decision as a fixed-scope ClickHouse consulting engagement and operates the resulting clusters under managed services. Every step on this page changes a running cluster’s topology: rehearse it on staging with production-sized data, keep the rollback path written down, and confirm a tested backup before the first shard is added.

Real-Time Payments Analytics
ClickHouse

Real-Time Payments Analytics on ClickHouse: 6 Proven Layers for Southeast Asia

ChistaDATA Inc.
A reference model for real-time payments analytics on ClickHouse with the Apache stack (Kafka, Flink, Iceberg, Spark, Airflow, Superset) for a high-volume Southeast Asian mobile payment platform: schema, ingestion, serving, residency and operations.

[…]

ClickHouse 26.8 LTS
ChistaDATA

ClickHouse 26.8 LTS: 7 Essential Performance and HA Changes

ChistaDATA Inc.
ClickHouse 26.8 LTS read for production: adaptive aggregation, IEJoin, plan-based parallel replicas, Keeper on disk, always_fetch_mutated_part, and the 26.3 to 26.8 breaking-change checklist for self-managed and Cloud.

[…]

ClickHouse Sharding
ClickHouse

ClickHouse Sharding Troubleshooting and Performance Optimization

ChistaDATA Inc.
A field guide to ClickHouse sharding troubleshooting: Distributed send-queue backlogs, duplicate rows from internal_replication, shard skew, initiator merge bottlenecks, distributed JOIN failures and unavailable shards, each with the system.* evidence, the mechanism, the staged fix, and the alert that catches it early.

[…]

ClickHouse sharding architecture showing distributed database nodes with multiple shards for high-growth environments
ClickHouse

ClickHouse Sharding Strategies for High-Growth Environments

ChistaDATA Inc.
ClickHouse sharding is one of the most consequential architectural decisions you will make when scaling an analytical database from millions to billions of rows per day. In a high-growth environment — where query volumes double […]
Scaling ClickHouse
ClickHouse

Scaling ClickHouse from Gigabytes to Petabytes: A Practical Playbook

ChistaDATA Inc.
Most teams don’t run into ClickHouse scaling challenges when they’re managing a few hundred gigabytes. The problems surface the moment data volumes cross a threshold where naive configurations begin to crack — queries slow down, […]
complex queries in clickhouse
ChistaDATA

Transforming Your Data in a Managed ClickHouse® Cluster with dbt

Shiv Iyer
Transforming Your Data in a Managed ClickHouse® Cluster with dbt: A Complete Guide Introduction In today’s data-driven landscape, organizations are constantly seeking efficient ways to transform raw data into actionable insights. The combination of ClickHouse®, […]
ChistaDATA

Building Fast Data Loops in ClickHouse®

ChistaDATA Inc.
Building Fast Data Loops in ClickHouse: From Insert to Query Response in ClickHouse® In today’s data-driven world, the speed at which you can ingest, process, and query data determines your competitive advantage. ClickHouse® excels at […]
ChistaDATA

Data Compression in ClickHouse for Performance and Scalability

ChistaDATA Inc.
Implementing Data Compression in ClickHouse: A Complete Guide to Optimal Performance and Scalability Introduction Data compression in ClickHouse is a critical optimization technique that can dramatically improve query performance, reduce storage costs, and enhance overall […]
ClickHouse Performance

How do we implement intelligent Caching on ClickHouse with machine learning?

Shiv Iyer
Introduction: Intelligent Caching Implementing intelligent caching with machine learning in a ClickHouse environment involves predicting data access patterns and optimizing cache usage based on these predictions. This approach helps to ensure that the most frequently […]
ClickHouse

ClickHouse Data Ingestion: Built for High-Velocity, High-Volume Data

Shiv Iyer
ClickHouse Is Ideal for High-Velocity ClickHouse is particularly well-suited for projects that require high-velocity, high-volume data ingestion and real-time analytics, primarily due to its specialised architecture and distinct features. Its columnar storage model plays a […]

Posts pagination

1 2 »

ChistaDATA is committed to open source software and building high performance ColumnStores

In the spirit of freedom, independence and innovation. ChistaDATA Corporation is not affiliated with ClickHouse Corporation 

Tell us how we can help!

Loading

Search ChistaDATA Website

★READ THIS WARNING★

* Everything changes over time – Our blogs/posts and comments changes over time, That’s how it should be! Whatever we comment from ChistaDATA Inc. Teams (including Shiv Iyer) and other stakeholders or guest bloggers posted here are never permanent, These things worked for us. But, there is no guarantee they will work for you too, When using the recommendations from ChistaDATA or MinervaDB or MinervaSQL or any other online resources / Google,  You must test the advice before applying them to your production systems, and always invest for a robust Database DR solution, Thank you for understanding. 

Recent Posts from ChistaDATA

  • Advanced ClickHouse Troubleshooting on 26.8 LTS: 9 Proven Techniques for Slow Queries, Stuck Merges and Memory Errors
  • ClickHouse 26.8 LTS: The Advanced Features That Change Real-Time Analytics Performance
  • Real-Time Analytics ClickHouse Workshop for CTOs and Data Architects
  • ClickHouse Performance Audit: 8 Best Tips for 26.8 LTS
  • Real-Time Payments Analytics on ClickHouse: 6 Proven Layers for Southeast Asia

☎ TOLL FREE PHONE (24*7)

(844)395-5717

🚩 ChistaDATA Inc. FAX

+1 (209) 314-2364

CORPORATE ADDRESS: CALIFORNIA

ChistaDATA Inc.
440 N BARRANCA AVE #9718 COVINA,
CA 91723
════════════════════════════════
Email: info@chistadata.com

CORPORATE ADDRESS: NEW CASTLE, DELAWARE

ChistaDATA Inc.,
256 Chapman Road STE 105-4,
Newark, New Castle 19702,
Delaware
════════════════════════════════
Email: info@chistadata.com

CORPORATE ADDRESS: DELAWARE

ChistaDATA Inc.,
PO Box 2093 PHILADELPHIA PIKE #3339
CLAYMONT, DE 19703
════════════════════════════════
Email: info@chistadata.com

HOW CAN WE HELP?

We are committed to building Optimal, Scalable, Highly Available, Reliable, Fault-Tolerant and Secured Database Infrastructure Operations for WebScale to our customers globally

CHISTADATA IS COMMITTED TO OPEN SOURCE SOFTWARE AND BUILDING HIGH PERFORMANCE COLUMNSTORES

In the spirit of freedom, independence and innovation. ChistaDATA Corporation is not affiliated with ClickHouse Corporation 

ChistaDATA Inc. Knowledge base is licensed under the Apache License, Version 2.0 (the “License”)

Copyright 2022 ChistaDATA Inc

Licensed under the Apache License, Version 2.0 (the “License”); you may not use this file except in compliance with the License. You may obtain a copy of the License at

http://www.apache.org/licenses/LICENSE-2.0

Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an “AS IS” BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions and limitations under the License.

PostgreSQL is a registered trademark of the PostgreSQL Community Association. ClickHouse is a registered trademark of ClickHouse, Inc. MongoDB is a registered trademark of MongoDB, Inc. Couchbase is a registered trademark of Couchbase, Inc. Redis is a registered trademark of Redis Ltd. Apache Cassandra is a registered trademark of the Apache Software Foundation. Milvus is a registered trademark of Zilliz. MinIO is a registered trademark of MinIO, Inc. Amazon Redshift and Amazon Aurora are registered trademarks of [Amazon.com](http://amazon.com/), Inc. Google Cloud is a registered trademark of Google LLC. Snowflake is a registered trademark of Snowflake Inc. Databricks is a registered trademark of Databricks, Inc. MySQL and InnoDB are registered trademarks of Oracle Corporation. MariaDB is a trademark of MariaDB Corporation Ab. All other trademarks are the property of their respective owners. Any other product or company names mentioned may be trademarks or trade names of their respective owners. Copyright © 2010–2026. All Rights Reserved by ChistaDATA®.

Contents

×
  • The ClickHouse horizontal scaling stages at a glance
  • Stage 0: exhaust vertical scaling before ClickHouse horizontal scaling
  • Stage 1: replicas, and the read-write split
  • Stage 2: sharding, the ClickHouse horizontal scaling step that changes everything
  • Choosing the shard key without creating a hot shard
  • Rebalancing: the cost of ClickHouse horizontal scaling done late
  • Keeper sizing: the coordination ceiling on ClickHouse horizontal scaling
  • Stage 3: object storage, tiering and parallel replicas
  • Capacity planning: the arithmetic behind the ClickHouse horizontal scaling plan
  • Compression and materialised columns as scaling levers
  • Troubleshooting a sharded cluster
  • Version notes
  • Reading the archive
→ Index