ClickHouse Cloud 索引分片:PB级数据,7.7倍分析加速!
本文字数:13915;估计阅读时间:35 分钟
作者:James Cunningham
TL;DR
索引分片 (Index Sharding) 将主索引和次级索引的分析工作分散到各个副本上,从而释放工作内存以供查询执行,并在我们对一个500亿行表的测试中,将索引分析速度提升高达7.7倍。
ClickHouse 旨在实现大规模下的高速性能。
然而,“规模”对不同的人来说有着截然不同的含义。对于 ClickHouse 工程团队而言,这是一个循环:先关注一个被认为是巨大的指标,然后努力优化,使其成为常态,接着再寻找下一个挑战。我们近期关注的重点便是索引大小和索引分析时间。
并行副本 (Parallel replicas) 是 ClickHouse 23.3 版本中引入的新特性,它通过在整个集群中分布式地读取数据,使查询执行实现了水平扩展。
索引分片 (Index Sharding) 将相同的原理向前推进了一步:即在实际读取任何数据之前进行的索引分析阶段。
截至目前,ClickHouse 一直基于一个假设运行:每个副本都会将整个索引从对象存储 (Object Storage) 加载到工作内存 (Working Memory) 中。
对于大多数工作负载而言,这几乎不会构成问题。索引通常较小,工作内存充裕,且 ClickHouse 稀疏索引 (Sparse Index) 的设计理念使其天生保持精简。
但在真正大规模的场景下,动辄数百亿行数据、PB 级存储、数十个甚至更多副本时,仅主键索引 (Primary Key Index) 就可能在每个副本上占用 100GB 甚至更多的空间。如果再叠加向量搜索索引 (Vector Search Indexes)、布隆过滤器 (Bloom Filters) 和全文搜索索引 (Full-Text Search Indexes) 等,你就会面临一个非常高昂的代价:同一个庞大的索引在每个副本上都被冗余加载。
索引分片 (Index Sharding) 改变了这一机制。不再是每个副本都加载所有索引,而是将索引在整个集群中进行分区。每个副本只负责其分到的部分,所有副本协作共同覆盖全部索引。
这一改变带来了三方面的显著优势:
横向扩展时,每个副本分配给索引的工作内存更少。
由于分析工作可以并行地分布到更多机器上,索引分析具备横向可扩展性。
由于引用局部性 (Locality of Reference)显著提升,单次处理速度更快。
当数据写入 ClickHouse MergeTree 表时,会以“parts”的形式进行组织,每个 part 都是不可变、自包含的列数据单元。ClickHouse 会在写入每个 part 时,同步构建其索引,并将其与列数据一同存储在对象存储中。
在 ClickHouse Cloud 中,这意味着 主键索引 (Primary Key Index)、标记文件 (Mark Files) 以及所有 二级索引 (Secondary Indexes) 都与 part 数据本身一同存储。
当涉及索引列的查询到来时,ClickHouse 必须将这些索引文件从对象存储加载到工作内存中,然后才能确定哪些数据粒度 (granules) 值得读取。
索引无法直接在磁盘或对象存储中进行查询;它必须驻留在内存中才能进行分析。
这使得索引内存成为每个副本的固定开销:每个参与查询分析的副本,在执行任何分析前都必须加载相关索引。ClickHouse 的主键索引被设计为稀疏索引。它不是追踪每一行数据,而是为每个数据粒度 (granule) 存储一个条目(称为“标记 (mark)”),默认情况下,一个数据粒度包含 8,192 行数据。
这使得索引足够小,可以完全加载到内存中,同时仍能提供 快速二分查找 功能,从而跳过索引判断为不相关的数据。
当您运行一个 WHERE 子句涉及索引列的查询时,ClickHouse 会执行索引分析:这是一个两步过程。首先,如果表配置了分区键 (partitioning key),它会从查询中移除所有不相关的分区。然后,它扫描每个数据分片 (data part) 的标记,以定位可能包含匹配行的粒度块 (granules),完全跳过其余部分,最后仅从磁盘流式传输选定的粒度块。正是这种机制赋予了 ClickHouse 独特的查询速度。在索引良好的表上,它仅读取极少一部分数据。
同样的机制也适用于二级索引 (secondary indexes):包括用于集合成员测试的布隆过滤器 (bloom filters)、用于全文搜索的文本索引 (text indexes),以及用于近似近邻搜索 (approximate nearest neighbor) 的向量索引 (vector indexes)。所有这些索引都会从各自的索引文件加载到工作内存中,以实现粒度块跳跃。
对于拥有数十亿行的典型表来说,这种机制完全可以应对。然而,当规模扩展到上限时,瓶颈便会浮现。
在 ClickHouse Cloud 中,所有副本 (replicas) 共享存储。一个在对象存储 (object storage) 中存储五 PB 数据的表,无需在副本之间重复存储这些数据;只有计算资源会被复制。数据存储于一处,副本仅在需要时进行流式传输。
然而,到目前为止,索引的情况却并非如此。尽管每个索引都与数据一同存储在对象存储中,但每个处理查询的副本都必须将活动的主键索引 (primary key index) 从对象存储加载到各自的工作内存中。
每个需要评估布隆过滤器或文本索引的副本,同样会加载相应的索引;或者通过配置 use_skip_indexes_on_data_read,将索引分析中的跳过索引评估下推到数据读取阶段执行。如果您为了应对更高的查询并发或吞吐量而增加副本,那么每个新增副本都会加载一份完整的索引副本。
在拍字节级别的数据规模下,仅主键在内存中,每个副本就能占用 100-400 GiB。若再考虑二级索引标记,以及向量搜索和全文搜索的日益普及,增加副本带来的内存开销会占据大量工作内存;而这部分内存本可以用于处理查询。
症结在于,为了获得更好的横向扩展性能而增加副本数量,反而会使情况恶化。以一个 100 GB 的索引为例,整个集群消耗的累计工作内存量与副本数量成正比:
每个副本都持有完全相同的索引。该索引的每一个字节都在集群中被完整复制,因此,每增加一个节点,成本便会随之攀升。
索引分片的核心理念很简单:
假设您有 N 个副本,那么每个副本仅需负责整个索引的 1/N。
其工作原理如下:当查询到来时,查询发起者不再在本地加载并分析完整的索引,而是将任务分配给所有可用的副本。每个副本通过一个虚拟哈希环 (virtual hash ring) 接收要分析的数据部分子集;该技术被称为一致性哈希 (consistent hashing)。每个副本从对象存储加载其被分配的索引部分,执行分析,并返回匹配的数据粒范围。发起者将这些结果合并成一个需要读取的完整视图,在此过程中,无需任何单个节点接触或加载整个索引。
随后,实际数据读取将以常规方式,通过并行副本 (parallel replicas) 完成,每个副本负责读取其已分析的数据部分。
您可以通过 EXPLAIN indexes=1 命令查看这种分布的实际效果:
EXPLAIN indexes=1
SELECT UserID FROM hits
WHERE UserID = 1
SETTINGS distributed_index_analysis = 1
FORMAT LineAsString;
Indexes:
PrimaryKey
Keys: UserID
Condition: (UserID in [...])
Parts: 208/208
Granules: 247702/143169495
Distributed:
Address: replica-1:9000 Parts received: 35 Granules received: 45094
Address: replica-2:9000 Parts received: 47 Granules received: 53988
Address: replica-3:9000 Parts received: 43 Granules received: 47387
Address: replica-4:9000 Parts received: 43 Granules received: 56130
Address: replica-5:9000 Parts received: 40 Granules received: 45103
以图表形式描绘的相同分布:
这种分布适用于所有索引类型。
ClickHouse Cloud 均支持主键索引 (primary key indexes)、布隆过滤器 (bloom filters)、全文搜索索引 (full-text search indexes) 和向量搜索索引 (vector search indexes)。这对于二级索引 (secondary indexes) 尤为重要,因为二级索引相对于表而言,其规模可能更大。一个包含多个文本或向量索引的表,其索引数据量很容易达到数百 GB,而这些数据此前必须在每个节点上进行复制。
当为服务添加副本 (replica) 时,分区分配 (part assignments) 会根据新的副本数量重新进行均衡。新副本加入服务时,可以选择通过启用 prewarm_primary_key_cache 和 prewarm_mark_cache 参数,预填充其主键缓存 (primary key cache) 和标记缓存 (mark cache)。如果这些参数未启用,新副本启动时会占用较少内存,并根据分析请求 (analysis requests) 的到来按需从对象存储 (object storage) 中加载其被分配的索引。当 use_primary_key_cache 启用时,现有副本会检测到某些分区不再是它们的责任,并在后台卸载这些分区,自动回收工作内存。
这正是 ClickHouse Cloud 计算存储分离 (compute-storage separation) 架构优势的集中体现。
在传统的无共享架构 (shared-nothing architecture) 中,添加一个副本意味着必须先将数据移动或复制到新节点,然后该节点才能参与查询执行 (query execution)。而在 ClickHouse Cloud 中,数据不会发生移动。所有副本都从同一个共享对象存储中读取数据,因此新副本一旦加载了其被分配的索引切片 (index slice),即可立即开始处理索引分析请求。因此,横向扩展 (scaling out) 的成本完全取决于索引加载时间 (index loading time),而非数据传输 (data transfer) 开销。
因此,索引分片 (Index Sharding) 带来的内存优势,恰好在您最需要的时候——即横向扩展 (scale out) 时——得到充分体现:每个副本的索引占用空间 (index footprint) 会随之缩小,而添加新副本的成本,仅是预热总索引中一小部分所需的时间。
在分布式系统中,瞬态故障处理 (Transient failure handling) 是一个需要重点考虑的关键组件。对特定数据分片进行索引分析的请求,可能会因多种原因发生瞬态故障。常见情况包括发起节点 (initiator) 与负责副本 (responsible replica) 之间请求过程中的网络故障,以及负责副本尚未加载请求的数据分片。然而,分布式系统中可能出现的问题层出不穷,每天都有新的挑战,因此,让我们来探讨一下我们如何处理已知的情况。
当发起节点将索引分析任务分发给所有副本时,每个副本会返回它们被要求分析的每个数据分片的分析结果。如果特定数据分片发生故障,我们的解决方案很简单:回退到本地分析该数据分片。发起节点将不再重试导致最初故障的副本,而是会将该数据分片的索引加载到本地内存中并运行分析。未来的请求将继续尝试联系负责节点,任何故障都会导致回退到发起节点的内存中进行处理。
在引入索引分片之前,副本与索引内存之间的关系是线性且不可避免的。一个 100GB 的索引部署在三个副本上,总共会消耗 300GB 的运行内存。如果扩展到九个副本,这个数字将变为 900GB,每个副本都会承载完整的索引数据,无论有多少其他副本在做同样的工作。
随着索引分片的引入,整个集群中索引所消耗的累计运行内存现在静态地限制为索引本身的尺寸,无论你运行多少个副本。增加更多副本时,累计总量保持不变,而每个独立副本的内存份额则会按比例缩小。
让我们通过一个例子来直观地展示索引分析任务如何在 ClickHouse 集群中分配。在一张拥有 16GB 主键 (primary key)、包含 25 个数据分片的表中,当 10 个副本上启用了分布式索引分析 (Distributed Index Analysis) 后,内存分配情况如下所示:
每个副本仅存储其所需数据。16GB 的索引在整个集群中以聚合形式存在一次,而非在集群内的每个节点上都单独存在。这彻底改变了横向扩展的经济模式。您无需被迫配置越来越大的实例来仅仅承载索引内存开销,即可增加副本数量以提升并发性和吞吐量。
每个副本上释放的工作内存,可用于其最初设计的用途:处理查询。
还有一个与第一个优势相辅相成的益处:当索引分析进行分布式处理时,其速度也会显著提升。
若无索引分片,从单个查询的角度来看,索引分析从根本上受限于单个节点。一个节点需要完成所有工作,在开始读取任何数据之前,扫描数亿个数据粒(granule)的标记。对于那些拥有负载繁重且选择性极高的二级索引的表,例如 Vector Search (向量搜索) 和 Full-text Search (全文搜索),这一分析阶段往往是查询的主要开销。
将分析工作分发到各个副本上,将单一瓶颈转化为分布式瓶颈。每个副本同时处理其数据切片,而协调器(initiator)则合并紧凑的范围结果,而不是由其自身执行完整的评估。更多的副本意味着更高的并行性,更高的并行性意味着更快的分析。
在一个包含 500 亿行数据(17,000 个数据分区,600 万个标记)的表上,以 10 个副本进行基准测试结果如下:
随着副本数量的增加,性能提升会进一步放大。在启用索引分片的情况下,将副本数量从 10 个扩展到 20 个,结果如下:
如果没有索引分片,增加副本虽然能提升数据读取吞吐量,但对索引分析毫无帮助,索引分析仍局限于单节点执行。而有了索引分片,你新增的每个副本都能提升分析吞吐量并有助于内存分布。在索引分析成为查询瓶颈时(这在包含二级索引的大型表上,特别是全文搜索和向量搜索场景中非常普遍),索引分片的作用就体现在:它能将索引分析从每次查询的固定成本负担,转变为从副本投资中获取的强大性能倍增器。
索引分片最适用于索引分析在查询成本中占据显著比例的工作负载:即具有多个二级索引、高副本数量,以及那些在数据读取开始前,严重依赖这些索引进行数据粒度(granules)筛选的选择性过滤条件。全文搜索、向量相似性和布隆过滤器索引是最典型的例子。在大型表上,它们每个副本可能占用数 GB 的工作内存,一旦表达到这种规模,随着副本数量的增加,内存节省和分析并行性都将呈倍数增长。
为确保分布式分析的协调开销始终物有所值,索引分片会在满足以下两个表级阈值时自动激活。第一个是 distributed_index_analysis_min_parts_to_activate (默认值: 10),它要求在尝试进行分布式处理前,数据分片数达到最小阈值。第二个,也是更重要的,是 distributed_index_analysis_min_indexes_bytes_to_activate (默认值: 1073741824,即 1GB),它要求磁盘上所有索引的未压缩总大小超过 1GB。低于此阈值时,本地加载索引既快速又经济。一旦超过此阈值,分析成本便开始显著影响查询延迟和每个副本的工作内存,而分布式处理恰能有效解决这些问题。
这两个阈值都是表级设置,可以根据你的工作负载进行调整:
ALTER TABLE my_favorite_table MODIFY SETTING
distributed_index_analysis_min_parts_to_activate = 20,
distributed_index_analysis_min_indexes_bytes_to_activate = 21474836480; -- 20 GB
当 distributed_index_analysis 被启用时,只有当这两个条件都满足时,分析才能从本地模式升级为分布式模式,从而确保较小规模的分析任务仍能保持高性能。
在我们一个内部数据库中,一个合并表 (merge table) 聚合了八个子表 (sub-tables)。通过执行 SELECT * ... LIMIT 1 查询,我们能够确保选择所有数据颗粒 (granules),并让查询处理 (query processing) 来限制查询结果。在对索引分析 (index analysis) 的精简解释中,我们发现索引分析被分发到了十个副本 (replicas) 上:
EXPLAIN indexes=1 select * from merge_table.merge_table LIMIT 1
SETTINGS distributed_index_analysis = 1
FORMAT LineAsString;
Expression ((Project names + (Projection + Change column names to column identifiers)))
Limit (preliminary LIMIT)
ReadFromMerge
Expression
ReadFromMergeTree (merge_table.merge_table)
Indexes:
MinMax
Condition: true
Parts: 1877/1877
Granules: 292646425/292646425
Partition
Condition: true
Parts: 1877/1877
Granules: 292646425/292646425
PrimaryKey
Condition: true
Parts: 1877/1877
Granules: 292646425/292646425
Distributed:
Replicas: 10
Parts send: 1074
Parts received: 1074
Granules send: 279134908
Granules received: 279134908
Ranges: 1877
Tables: 8
如果我们移除 compact=1 参数来展开分析并细化观察,可以看到有两张表由于其自身规模较大,从分布式分析中被排除在外(即在本地进行分析):
EXPLAIN indexes=1 select * from merge_table.merge_table LIMIT 1
SETTINGS distributed_index_analysis = 1
FORMAT LineAsString;
Expression ((Project names + (Projection + Change column names to column identifiers)))
...
Expression
ReadFromMergeTree (merge_table.table-1)
Indexes:
...
PrimaryKey
Condition: true
Parts: 11/11
Granules: 1314708/1314708
Ranges: 11
Expression
ReadFromMergeTree (merge_table.table-2)
Indexes:
...
PrimaryKey
Condition: true
Parts: 112/112
Granules: 76119629/76119629
Distributed:
Address: replica-1:9000
Parts send: 24
Parts received: 24
Granules send: 16902011
Granules received: 16902011
Address: replica-2:9000
Parts send: 21
Parts received: 21
Granules send: 13104084
Granules received: 13104084
...
Ranges: 112
Expression
ReadFromMergeTree (merge_table.table-3)
Indexes:
...
PrimaryKey
Condition: true
Parts: 14/14
Granules: 41/41
Ranges: 14
Expression
ReadFromMergeTree (merge_table.table-4)
Indexes:
...
PrimaryKey
Condition: true
Parts: 21/21
Granules: 1955/1955
Ranges: 21
Expression
ReadFromMergeTree (merge_table.table-5)
Indexes:
...
PrimaryKey
Condition: true
Parts: 401/401
Granules: 654988/654988
Ranges: 401
Expression
ReadFromMergeTree (merge_table.table-6)
Indexes:
...
PrimaryKey
Condition: true
Parts: 962/962
Granules: 203015279/203015279
Distributed:
Address: replica-1:9000
Parts send: 195
Parts received: 195
Granules send: 43951345
Granules received: 43951345
Address: replica-2:9000
Parts send: 172
Parts received: 172
Granules send: 37951746
Granules received: 37951746
...
Expression
ReadFromMergeTree (merge_table.table-7)
...
PrimaryKey
Condition: true
Parts: 253/253
Granules: 11413730/11413730
Ranges: 253
Expression
ReadFromMergeTree (merge_table.table-8)
Indexes:
...
PrimaryKey
Condition: true
Parts: 103/103
Granules: 125095/125095
Ranges: 103
我们得到的结果是,数据颗粒在两次分布式分析中分布得异常均匀,而其余的表则显示出:
分布式到 5 个副本 (table-2 和 table-6):
在发起者上进行本地分析(所有低于阈值的表):
分布式索引分析的完整 EXPLAIN 输出
EXPLAIN indexes=1 select * from merge_table.merge_table LIMIT 1
SETTINGS distributed_index_analysis = 1
FORMAT LineAsString;
Expression ((Project names + (Projection + Change column names to column identifiers)))
Limit (preliminary LIMIT)
ReadFromMerge
Expression
ReadFromMergeTree (merge_table.table-1)
Indexes:
MinMax
Condition: true
Parts: 11/11
Granules: 1314708/1314708
Partition
Condition: true
Parts: 11/11
Granules: 1314708/1314708
PrimaryKey
Condition: true
Parts: 11/11
Granules: 1314708/1314708
Ranges: 11
Expression
ReadFromMergeTree (merge_table.table-2)
Indexes:
MinMax
Condition: true
Parts: 112/112
Granules: 76119629/76119629
Partition
Condition: true
Parts: 112/112
Granules: 76119629/76119629
PrimaryKey
Condition: true
Parts: 112/112
Granules: 76119629/76119629
Distributed:
Address: replica-1:9000
Parts send: 24
Parts received: 24
Granules send: 16902011
Granules received: 16902011
Address: replica-2:9000
Parts send: 21
Parts received: 21
Granules send: 13104084
Granules received: 13104084
Address: replica-3:9000
Parts send: 20
Parts received: 20
Granules send: 12832693
Granules received: 12832693
Address: replica-4:9000
Parts send: 23
Parts received: 23
Granules send: 16373147
Granules received: 16373147
Address: replica-5:9000
Parts send: 24
Parts received: 24
Granules send: 16907694
Granules received: 16907694
Ranges: 112
Expression
ReadFromMergeTree (merge_table.table-3)
Indexes:
MinMax
Condition: true
Parts: 14/14
Granules: 41/41
Partition
Condition: true
Parts: 14/14
Granules: 41/41
PrimaryKey
Condition: true
Parts: 14/14
Granules: 41/41
Ranges: 14
Expression
ReadFromMergeTree (merge_table.table-4)
Indexes:
MinMax
Condition: true
Parts: 21/21
Granules: 1955/1955
Partition
Condition: true
Parts: 21/21
Granules: 1955/1955
PrimaryKey
Condition: true
Parts: 21/21
Granules: 1955/1955
Ranges: 21
Expression
ReadFromMergeTree (merge_table.table-5)
Indexes:
MinMax
Condition: true
Parts: 401/401
Granules: 654988/654988
Partition
Condition: true
Parts: 401/401
Granules: 654988/654988
PrimaryKey
Condition: true
Parts: 401/401
Granules: 654988/654988
Ranges: 401
Expression
ReadFromMergeTree (merge_table.table-6)
Indexes:
MinMax
Condition: true
Parts: 962/962
Granules: 203015279/203015279
Partition
Condition: true
Parts: 962/962
Granules: 203015279/203015279
PrimaryKey
Condition: true
Parts: 962/962
Granules: 203015279/203015279
Distributed:
Address: replica-1:9000
Parts send: 195
Parts received: 195
Granules send: 43951345
Granules received: 43951345
Address: replica-2:9000
Parts send: 172
Parts received: 172
Granules send: 37951746
Granules received: 37951746
Address: replica-3:9000
Parts send: 190
Parts received: 190
Granules send: 38465035
Granules received: 38465035
Address: replica-4:9000
Parts send: 203
Parts received: 203
Granules send: 41708027
Granules received: 41708027
Address: replica-5:9000
Parts send: 202
Parts received: 202
Granules send: 40939126
Granules received: 40939126
Ranges: 962
Expression
ReadFromMergeTree (merge_table.table-7)
Indexes:
MinMax
Condition: true
Parts: 253/253
Granules: 11413730/11413730
Partition
Condition: true
Parts: 253/253
Granules: 11413730/11413730
PrimaryKey
Condition: true
Parts: 253/253
Granules: 11413730/11413730
Ranges: 253
Expression
ReadFromMergeTree (merge_table.table-8)
Indexes:
MinMax
Condition: true
Parts: 103/103
Granules: 125095/125095
Partition
Condition: true
Parts: 103/103
Granules: 125095/125095
PrimaryKey
Condition: true
Parts: 103/103
Granules: 125095/125095
Ranges: 103
目前,ClickHouse Cloud 已为 SharedMergeTree 表提供 Index Sharding (索引分片) 的私有预览版。如果您正在运行的业务工作负载面临索引内存限制,或者索引分析时间是查询延迟 (query latency) 的重要组成部分,我们非常希望听取您的反馈。
如需请求访问权限,请联系您的 ClickHouse 客户团队,或访问 clickhouse.com/contact 与我们联系。
/END/
试用阿里云 ClickHouse企业版
轻松节省30%云资源成本?阿里云数据库ClickHouse 云原生架构全新升级,首次购买ClickHouse企业版计算和存储资源组合,首月消费不超过99.58元(包含最大16CCU+450G OSS用量)了解详情:https://t.aliyun.com/Kz5Z0q9G
征稿启示
面向社区长期正文,文章内容包括但不限于关于 ClickHouse 的技术研究、项目实践和创新做法等。建议行文风格干货输出&图文并茂。质量合格的文章将会发布在本公众号,优秀者也有机会推荐到 ClickHouse 官网。请将文章稿件的 WORD 版本发邮件至:[email protected]