先说结论
Distributed 表本身通常不存业务数据,而是把写入路由到分片的本地 MergeTree 表,并把查询下发到相关分片后汇总结果。分片键决定数据分布,副本配置决定同一分片内的冗余与读取选择。
随机分片分布简单但难以裁剪;按租户或业务键分片可让相关查询落在少数分片,却要防止大租户热点。
查询流程
协调节点改写并下发查询,各分片执行局部过滤和聚合,再返回中间结果进行最终合并。应尽量把过滤和预聚合下推,避免跨网络传输大量明细。
本地表与 Distributed 表的职责
客户端 -> Distributed 表(路由/聚合)
↓
shard 1 本地 MergeTree + replicas
shard 2 本地 MergeTree + replicas
shard 3 本地 MergeTree + replicas
本地表负责真实数据存储、排序、分区和副本复制;Distributed 表是逻辑入口,包含集群拓扑、目标表和分片表达式。运维时要确认每个节点本地表定义一致,Distributed 表不会替你自动修复底层 schema 漂移。
写入模式对比
直接由应用按路由写本地表,路径短、失败可见,但客户端要掌握拓扑和重试;写 Distributed 表可由服务端根据分片键转发,某些配置下支持本地队列异步投递,客户端简单,但要监控队列积压、目标不可达与重复投递语义。
无论哪种,重试都可能造成重复数据。事实表应携带事件 ID、版本或采用可最终去重的表设计;不能因为“插入成功超时”就盲目无限重发。
分片键与查询裁剪
如果按 cityHash64(tenant_id) 分片,大多数带 tenant_id 的点查可以用 optimize_skip_unused_shards 等能力减少扇出,前提是查询条件和分片表达式可推导。若 SQL 没有分片键,协调节点通常要广播到所有分片。
分片键同时决定均衡性和局部性。按时间分片便于冷热存储但当前时间写热点明显;按租户哈希均匀但跨租户报表需要全分片聚合。不要只看数据量均匀,还要看 CPU、QPS 和磁盘 I/O 是否均匀。
分布式聚合与 JOIN
分片先完成局部 group by、filter 和 partial aggregate,协调节点再汇总。尽量让局部阶段减少数据量;在 SELECT * 大明细、全局排序、大 offset、精确 distinct 上,网络和协调节点内存可能成为瓶颈。
跨分片 JOIN 的语义和代价要明确。小维表可复制到每个分片或使用字典;大表广播会爆炸。GLOBAL 相关行为会把右表结果传到各分片,不可在大数据集上盲用。
故障与可观测性
查询失败时区分单副本不可用、单分片无副本、协调节点过载和网络超时。监控每分片扫描行数、返回字节、队列、replication lag 和热点,平均集群 QPS 健康不代表单一分片没有倾斜。
落地时要注意什么
写入可直接写本地表并由应用路由,也可写 Distributed 表异步转发;两种方式在失败重试、落盘队列和可观测性上不同。副本表还需正确配置一致的拓扑宏和复制路径。
容易踩坑的地方
增加分片不一定让查询更快,扇出、网络和最终聚合可能增加延迟。GLOBAL JOIN 等操作可能广播数据,应控制右表大小并优先通过建模减少跨分片连接。
常见问题
追问:分片键怎么选?
在均匀分布、查询局部性和扩展性之间权衡,用真实租户规模和查询模式验证,必要时为超大租户单独路由。
追问:Distributed 表可以当作高可用代理吗?
它可按拓扑选择副本并路由,但高可用还依赖本地表副本配置、网络、协调节点容量和客户端重试。不能只部署一台协调节点。
追问:为什么一个查询比单机更慢?
分布式查询增加扇出、网络、局部/全局归并和慢分片拖尾。若数据量不大或过滤无法下推,单机可能更快。
追问:如何避免跨分片 JOIN?
用小维表复制、字典、预聚合宽表、物化视图或在 ETL 阶段完成关联。核心原则是让在线分析查询尽量局部完成。