ClickHouse 聚合查询在 100 亿行场景 RT 10 秒,如何优化分区和索引

解读

  1. 场景定位:100 亿行≈TB 级,单表日增量 5~10 亿行,常见于实时数仓、行为日志、IoT 时序。
  2. 性能目标:国内互联网生产普遍要求 P99<2 s,10 s 明显超标,需降一个数量级。
  3. 瓶颈假设:
    • 分区剪枝失效 → 扫描全表 100 亿行
    • 索引跳数不足 → Granule 读取过多
    • 聚合现场计算 → CPU 成为短板
    • 单节点线程/磁盘打满 → 横向扩展未到位
  4. 面试考点:能否把“分区设计、一级索引、二级索引、物化视图、集群拓扑”串成闭环,并给出可量化的验证手段。

知识点

  1. ClickHouse 存储层:
    • Partition Directory → Part → Column.bin/Mrk 文件
    • 一级索引 primary.idx(稀疏,每 8192 行一个 mark)
    • 二级索引(跳数索引:minmax、set、bloom_filter、tokenbf_v1)
  2. 分区剪枝:WHERE 条件必须包含分区键,且类型与分区表达式一致,否则无法下推。
  3. 索引选择性:一级索引键顺序决定 marks 过滤率;二级索引弥补低基数列。
  4. 聚合加速:
    • 预聚合(物化视图 / AggregatingMergeTree)
    • 采样(SAMPLE)
    • 增量合并(GROUP BY INCREMENTAL)
  5. 资源并行:max_threads、max_bytes_before_external_group_by、max_memory_usage 需与 CPU 核数、内存匹配;分布式表需避免“单点 reduce”。
  6. 国内常用工具:
    • 系统表 system.parts、system.query_log、system.metric_log
    • 阿里云 EMR-CK、腾讯云 TCHouse,均支持弹性扩容与 SSD 本地盘,面试需提到“云盘 IOPS 上限”。

答案

“我会按‘先剪枝、再索引、后预聚合、最后扩容’四步落地,每一步都给出量化验证。”

  1. 分区剪枝
    a) 选分区键:业务查询 90% 带 stat_date + product_id,采用 toYYYYMMDD(event_time) 作为分区键,保证单次查询只扫最近 7 天,分区目录从 365 个降到 7 个,扫描行数直接降 98%。
    b) 分区粒度:单分区 1~2 亿行(约 10 GB),Part 数量 < 300,避免 Too many parts 报错;通过 max_bytes_to_merge_at_max_space_in_pool=150G 控制后台合并。
    c) 验证:执行 EXPLAIN SYNTAX SELECT … WHERE stat_date=20240620,确认 Selected Marks 从 12 M 降到 240 k。

  2. 一级索引优化
    a) 排序键顺序:(stat_date, product_id, user_id),保证高基数在前,marks 过滤率 > 95%;通过 system.query_log 查看 ReadRows/TotalRows≈5%
    b) 避免函数:查询中禁止 toDate(event_time)=‘2024-06-20’,改为 event_time>=‘2024-06-20 00:00:00’ AND event_time<‘2024-06-21 00:00:00’,让索引走范围扫描。

  3. 二级索引补充
    a) 对低基数列 status(仅 8 个枚举值)建 set(0) 跳数索引,使 WHERE status=2 额外过滤 70% marks。
    b) 对高基数字符串 device_idbloom_filter(0.01),假阳性 <1%,在 ad-hoc 查询中 marks 过滤率从 0 提升到 88%。

  4. 预聚合 / 物化视图
    a) 创建 AggregatingMergeTree 物化视图,按 (stat_date, product_id, city) 预计算 sumMerge(uv)、avgMerge(duration),原 100 亿行聚合变 2 亿行,查询 RT 从 10 s 降到 0.8 s。
    b) 写入侧采用 INSERT SELECT 实时同步,延迟 <1 min;通过 system.mutations 监控合并压力。

  5. 集群层横向扩展
    a) 采用 8 分片 × 2 副本,分布式表按 rand() 分片,避免热点;每个分片只扫 12.5 亿行,并行度 8×16 核,CPU 利用率从 60% 提升到 85%。
    b) 调整 max_threads=32max_memory_usage=0(让 CK 自动感知容器 limit),max_bytes_before_external_group_by=200G,防止内存溢出。

  6. 量化结果
    同一 SQL SELECT product_id, sum(uv) FROM log WHERE stat_date=20240620 GROUP BY product_id

    • 优化前:扫描 100 亿行,RT 10.2 s,P99 12 s
    • 优化后:扫描 1.2 亿行(物化视图),RT 0.7 s,P99 1.1 s
      满足国内 SLA P99<2 s 要求,且磁盘 I/O 下降 90%,CPU 下降 65%。
  7. 回退方案
    若物化视图延迟敏感,可降级为 MergeTree + 索引优化,RT 也能降到 2.5 s,仍优于 10 s。

拓展思考

  1. 分区键与生命周期:国内合规要求 90 天冷存、180 天删除,可采用 TTL toDate(event_time) + INTERVAL 90 DAY TO VOLUME 'cold', 180 DAY DELETE,面试可追问“冷存用 HDD 还是对象存储?”。
  2. 索引膨胀:跳数索引也会占用磁盘,每列增加 2~5%,需在 system.parts 监控 secondary_index_bytes;当索引命中率 <80% 时及时下线。
  3. 多租户隔离:同一集群不同业务线共用,为防止大查询打满 CPU,可开启 max_concurrent_queries_for_userpriority,并结合 K8s cgroup 硬限。
  4. 版本差异:国内主流 21.8→23.8 升级后,analyzer=1 改写 SQL 计划,可能导致索引失效,上线前必须在影子集群回放 7 天流量。
  5. 真实面试陷阱:考官可能追问“如果分区键选错,如何在线重建?”——答案是用 ALTER TABLE ATTACH PARTITION 新建临时表双写,再原子替换,避免阻塞写入。