第六章 分区 - 把大数据拆小
导读
当数据量增长到单台机器无法容纳时,我们需要将数据分区(Partitioning)(也称为分片,Sharding)——将一个大集合拆分为多个小集合,分布在不同节点上。每个节点只负责一部分数据,从而突破单机存储和计算的限制。
分区是水平扩展的核心手段。但分区不是简单地把数据"随便分分"——不合理的分区策略可能导致数据倾斜(某些节点负载过高)、跨分区查询效率低下、甚至一致性问题。
本章将深入探讨分区的策略、分区与复制的关系、分区再平衡(Rebalancing)的挑战,以及跨分区查询的困难。这些知识对于构建大规模分布式数据库至关重要。
核心概念详解
6.1 分区策略
6.1.1 按键范围分区(Range-Based Partitioning)
按键范围分区将键空间划分为连续的区间,每个分区负责一个区间。
例如,将用户 ID 从 0 到 999999 分为 10 个分区:
- 分区 0:ID 0-99999
- 分区 1:ID 100000-199999
- ...
- 分区 9:ID 900000-999999
每个分区可以独立存储在不同的节点上。
优势:
- 范围查询高效:查询
WHERE id BETWEEN 100 AND 200只需要访问一个分区。 - 自然有序:数据按键顺序存储,适合需要排序和范围扫描的场景。
劣势:
- 数据倾斜(Hot Spot):如果某些键范围的数据量远大于其他范围(如按时间分区,最新时间的数据量最大),对应的分区会成为热点,负载集中。
- 边界确定困难:需要预先知道数据的分布,才能合理设置分区边界。
典型应用:
- HBase:按 RowKey 的范围分区,每个 Region 负责一个连续的范围。
- MongoDB(分片集群):支持按范围分区(使用 Shard Key)。
- Cassandra:按 Partition Key 的哈希分区,但同一 Partition 内的数据按 Clustering Key 范围排序。
6.1.2 按哈希分区(Hash-Based Partitioning)
按哈希分区使用一个哈希函数将键映射到分区编号。
例如,使用 hash(key) % N 将键均匀分配到 N 个分区。
优势:
- 数据分布均匀:好的哈希函数可以将数据均匀分布到所有分区,避免数据倾斜。
- 简单:不需要预先知道数据分布。
劣势:
- 范围查询低效:哈希函数打乱了键的顺序,
WHERE id BETWEEN 100 AND 200可能需要访问所有分区。 - 热点仍然存在:如果某些键的访问量特别高(如某个热门用户的 ID),哈希分区不能解决这个问题——它只解决数据量的倾斜,不解决访问频率的倾斜。
典型应用:
- Cassandra:按 Partition Key 的 MD5 哈希分区。
- Kafka:按消息 Key 的哈希分区到不同的 Partition。
- Elasticsearch:按文档 ID 的哈希分区到不同的 Shard。
6.1.3 一致性哈希(Consistent Hashing)
简单的 hash(key) % N 分区有一个严重问题:当分区数 N 变化时(增加或减少分区),几乎所有键的分区归属都会改变,导致大量数据需要迁移。
一致性哈希解决了这个问题。其核心思想:
将哈希空间组织成一个环(0 到 2^32-1)。
每个节点在环上占据一个或多个位置(通过 hash(node_id) 确定)。
每个键通过 hash(key) 映射到环上,然后顺时针找到最近的节点,该节点就是键的归属节点。
当增加一个节点时:
- 只有新节点和前一个节点之间的键需要迁移。
- 其他键的归属不变。
- 数据迁移量约为
总数据量 / N。
虚拟节点(Virtual Node):
- 物理节点在环上可能分布不均匀,导致数据倾斜。
- 解决方案:每个物理节点在环上占据多个位置(虚拟节点),使得数据分布更均匀。
- 通常每个物理节点对应 100-200 个虚拟节点。
一致性哈希的典型应用:
- Amazon Dynamo:使用一致性哈希进行数据分区。
- Memcached:分布式缓存,使用一致性哈希(或变体)分配数据。
- Apache Cassandra:使用一致性哈希的变体(Partitioner)。
6.1.4 按属性分区
除了按键值分区,还可以按数据的某些属性进行分区:
- 按地理位置分区:将不同地区的数据存储在不同节点。例如,亚洲用户的数据存储在亚洲节点,欧洲用户的数据存储在欧洲节点。
- 按租户分区:在多租户系统中,将不同租户的数据存储在不同分区。
- 按时间分区:将不同时间段的数据存储在不同分区。例如,按月分区,每月数据一个分区。
这些分区策略通常与范围分区或哈希分区结合使用。
6.2 分区与复制
分区和复制是两个正交的维度:分区决定数据分布在哪些节点上,复制决定每个分区有多少副本。
典型的架构是:先分区,再复制。
例如,一个 9 节点集群,3 个分区,每个分区 3 个副本:
- 分区 1 的 Leader 在节点 A,Follower 在节点 B 和 C。
- 分区 2 的 Leader 在节点 B,Follower 在节点 C 和 A。
- 分区 3 的 Leader 在节点 C,Follower 在节点 A 和 B。
这种安排使得每个节点都是一个分区的 Leader 和两个分区的 Follower,负载均衡。
分区 Leader 的分布:
- 需要确保同一分区的多个副本不在同一个物理机架或数据中心,以防止机架级故障。
- 不同分区的 Leader 应该均匀分布在不同节点上,避免某个节点成为写入瓶颈。
6.3 分区再平衡
6.3.1 为什么需要再平衡
随着数据量增长或节点增减,需要重新分配分区到节点上,这个过程称为再平衡(Rebalancing)。
再平衡的要求:
- 数据不丢失:迁移过程中不能丢失数据。
- 服务不中断:迁移过程中系统仍然可以读写。
- 最小化数据迁移:只迁移必要的分区,减少网络带宽消耗。
6.3.2 固定分区数
一种常见的策略是预先创建固定数量的分区(远大于预期的节点数),然后在节点增减时重新分配分区。
例如,预先创建 100 个分区,初始 3 个节点,每个节点约 33 个分区。增加到 4 个节点时,重新分配为每个节点 25 个分区。
优势:
- 分区数不变,分区边界不变,不需要重新哈希。
- 数据迁移以分区为单位,粒度可控。
劣势:
- 分区数在创建时确定,后续不能轻易改变。
- 如果数据量增长远超预期,分区数可能不够,需要重新分区(代价很高)。
典型应用:Kafka、Elasticsearch、Solr。
6.3.3 动态分区
另一种策略是动态创建和合并分区。
- 分裂(Split):当某个分区的数据量超过阈值,将其分裂为两个更小的分区。
- 合并(Merge):当某个分区的数据量低于阈值,将其与相邻分区合并。
优势:
- 分区数可以随数据量自动调整。
- 不需要预先估计数据量。
劣势:
- 分裂和合并操作复杂,需要保证数据不丢失和服务不中断。
- 范围分区的分裂/合并相对简单(只需在中间点切开),哈希分区的分裂/合并更复杂。
典型应用:HBase(Region 的分裂和合并)。
6.4 跨分区查询
分区的一个重大限制是跨分区查询的效率低下。
6.4.1 二级索引的挑战
主键索引通常与分区键一致——数据按主键分区,主键查询可以直接路由到正确的分区。但二级索引(非分区键的索引)则面临挑战:
方案一:每个分区维护自己的二级索引(Local Index)
- 每个分区独立维护自己的二级索引。
- 查询时需要扇出(Fan-out)到所有分区,每个分区在本地索引中查找,然后汇总结果。
- 优势:写入高效,只需要更新本地索引。
- 劣势:读取需要访问所有分区,延迟高。
方案二:全局二级索引(Global Index)
- 维护一个覆盖所有分区的全局索引。
- 查询时只需要访问全局索引,找到数据所在的分区,然后直接读取。
- 优势:读取高效,只需要两次访问(索引 + 数据)。
- 劣势:写入需要更新全局索引,可能涉及跨分区通信,一致性保证复杂。
全局索引的实现方式:
- 集中式索引:一个专门的节点维护全局索引。简单但可能成为瓶颈。
- 分布式索引:全局索引本身也是分区的。查询需要先路由到正确的索引分区。
典型应用:
- Cassandra:只支持 Local Index(2.x)和 SASI Index(自定义)。
- CockroachDB / TiDB:支持 Global Index(通过分布式事务保证一致性)。
- Elasticsearch:每个 Shard 维护自己的索引,查询时扇出到所有 Shard。
6.4.2 关联查询(JOIN)
分布式数据库中的 JOIN 操作极其困难,因为关联的数据可能分布在不同的分区上。
方案一:应用层 JOIN
- 分别查询各个表的数据,在应用层进行关联。
- 优势:数据库实现简单。
- 劣势:需要传输大量数据到应用层,网络开销大。
方案二:嵌套循环 JOIN(Broadcast JOIN)
- 将小表广播到所有分区,在每个分区上与大表的本地数据进行 JOIN。
- 优势:减少跨分区数据传输。
- 劣势:小表必须能放入每个分区的内存。
方案三:分区 JOIN(Co-located JOIN)
- 确保关联的表使用相同的分区键和分区策略,使得关联的数据在同一个分区上。
- 优势:JOIN 在本地执行,无跨分区通信。
- 劣势:限制了分区策略的选择,不是所有 JOIN 都能满足。
方案四:分布式 JOIN
- 使用 Shuffle 操作,将关联的数据重新分区到同一个节点上,然后执行 JOIN。
- 优势:适用于任意 JOIN。
- 劣势:需要大量数据 Shuffle,网络开销大。
这就是为什么大多数分布式数据库(如 Cassandra、HBase)不支持 JOIN,而分布式 SQL 数据库(如 CockroachDB、TiDB)的 JOIN 性能通常不如集中式数据库。
6.5 分区倾斜
6.5.1 问题描述
分区倾斜(Partition Skew)是指某些分区的数据量或访问量远大于其他分区,导致负载不均衡。
常见原因:
- 自然数据分布不均:社交网络中,某些用户的数据量远大于普通用户。
- 时间序列数据:最新时间的数据访问量最大。
- 哈希冲突:不同的键被哈希到同一个分区。
6.5.2 解决方案
引入随机前缀:在键前添加随机前缀(如 0-99),将一个热点键分散到 100 个分区。读取时读取所有前缀的数据并汇总。
复合分区键:使用复合键(如 user_id + timestamp)代替单一键,使得同一用户的数据按时间分散到不同分区。
读写分离:对于读热点,增加副本数;对于写热点,增加分区数。
应用层缓存:将热点数据缓存到 Redis 等缓存系统中,减少数据库压力。
重要知识点
知识点 1:分区键的选择
分区键的选择是分区设计中最关键的决策。好的分区键应该满足:
数据分布均匀:避免某些分区过大或过小。
查询模式匹配:常用的查询应该能路由到单个分区,避免跨分区查询。
避免热点:高频访问的数据不应该集中在同一个分区。
实践中,分区键的选择往往需要在"数据分布均匀"和"查询效率"之间权衡。例如:
- 按
user_id分区:查询单个用户的数据高效,但如果某些用户数据量大,会倾斜。 - 按
user_id + date分区:数据分布更均匀,但查询单个用户的所有数据需要跨分区。
知识点 2:分区与事务
分布式事务是分区系统的一个重大挑战。当事务涉及多个分区时,需要跨分区的协调来保证原子性。
两阶段提交(Two-Phase Commit, 2PC):
准备阶段:协调者向所有参与者发送准备请求,参与者执行事务但不提交,锁定相关数据。
提交阶段:如果所有参与者都准备成功,协调者发送提交请求;否则发送回滚请求。
2PC 的问题:
- 阻塞:如果协调者在准备阶段后崩溃,参与者会一直等待(持有锁),直到协调者恢复。
- 性能开销:需要多轮网络通信和磁盘日志。
替代方案:
- Saga:将长事务拆分为一系列本地事务,每个本地事务有对应的补偿操作。如果某个步骤失败,执行之前步骤的补偿操作。
- Try-Confirm-Cancel(TCC):类似 2PC 但更灵活,每个操作分为 Try(预留资源)、Confirm(确认提交)、Cancel(取消释放)三个阶段。
知识点 3:跨分区操作的局限性
分区系统对跨分区操作有天然的限制:
- 唯一约束:保证跨分区的唯一性(如全局唯一的用户名)需要跨分区通信,性能差。很多分布式数据库只保证分区内的唯一性。
- 自增 ID:自增 ID 在分区系统中很难实现,因为每个分区独立分配 ID 会导致 ID 不连续或冲突。替代方案:UUID、雪花算法(Snowflake)。
- 排序:全局排序需要所有分区的数据汇总后排序,无法利用分区的局部性。
这些限制是分区系统"水平扩展"的代价。在设计数据模型时,需要充分考虑这些限制,避免依赖跨分区操作。
知识点 4:分区策略对比
| 策略 | 数据分布 | 范围查询 | 再平衡 | 热点处理 |
|---|---|---|---|---|
| 范围分区 | 可能倾斜 | 高效 | 需要分裂/合并 | 困难 |
| 哈希分区 | 均匀 | 低效 | 需要一致性哈希 | 需要额外策略 |
| 一致性哈希 | 较均匀 | 低效 | 高效(最小迁移) | 需要虚拟节点 |
| 目录分区 | 可控 | 取决于策略 | 需要重新分配 | 取决于策略 |
常见误区
误区 1:"分区越多越好"
纠正:过多的分区会导致:
- 每个分区的元数据(如分区边界、Leader 信息)需要维护,元数据量过大。
- 跨分区查询需要扇出到更多分区,延迟增加。
- 分区之间的协调(如分布式事务)开销增加。
- 副本数 × 分区数 = 总副本数,存储和网络开销增加。
通常建议分区数在几十到几千之间,具体取决于数据量和集群规模。
误区 2:"哈希分区可以解决所有数据倾斜问题"
纠正:哈希分区只解决数据量的倾斜(使得每个分区的数据量大致相等),但不能解决访问频率的倾斜。如果某个键的访问量特别高(如热门商品的 ID),无论怎么哈希,这个键的所有访问都会路由到同一个分区。解决访问频率的倾斜需要应用层缓存、读写分离等策略。
误区 3:"分区后就不需要复制了"
纠正:分区和复制解决的是不同的问题。分区解决数据量和负载的扩展问题,复制解决可用性和容错问题。一个分区如果没有副本,当该分区的节点故障时,该分区的数据就不可用了。生产环境中,每个分区至少应该有 2-3 个副本。
误区 4:"分布式数据库的 JOIN 和集中式数据库一样快"
纠正:分布式数据库的 JOIN 通常比集中式数据库慢得多,因为关联的数据可能分布在不同节点上,需要跨网络传输数据。如果 JOIN 的性能是核心需求,应该考虑:
- 使用 Co-located JOIN(确保关联数据在同一分区)。
- 在应用层进行 JOIN。
- 使用反规范化(将关联数据冗余存储在同一表中)。
- 使用专门的 OLAP 系统(如 Spark SQL、Presto)进行复杂查询。
误区 5:"分区策略一旦确定就不能改变"
纠正:虽然改变分区策略的代价很高(需要重新分配所有数据),但并非不可能。一些数据库支持在线重新分区:
- HBase:支持 Region 的分裂和合并。
- CockroachDB:支持动态调整分区范围。
- Elasticsearch:不支持在线重新分区,但支持 Reindex API(代价较高)。
在设计初期就选择合适的分区策略,可以避免未来的重新分区成本。
实践应用
实践 1:设计分区策略
设计分区策略的步骤:
分析数据量:预估总数据量和增长趋势。
分析访问模式:哪些查询最频繁?是点查还是范围查?是否涉及 JOIN?
选择分区键:根据访问模式选择分区键,确保常用查询能路由到单个分区。
选择分区策略:根据数据分布选择范围分区或哈希分区。
确定分区数:根据数据量和集群规模确定初始分区数。
模拟验证:使用真实数据模拟分区效果,检查是否存在倾斜。
实践 2:监控分区均衡
定期监控分区的均衡性:
数据量分布:每个分区的数据量是否大致相等?最大分区和最小分区的比值是多少?
请求分布:每个分区的读写请求量是否均衡?是否存在热点分区?
副本分布:同一分区的副本是否分布在不同机架/数据中心?
Leader 分布:各节点的 Leader 数量是否均衡?
当发现不均衡时,及时采取措施:
- 数据倾斜:调整分区策略或分裂/合并分区。
- 访问热点:增加副本、引入缓存、或使用随机前缀分散请求。
实践 3:处理跨分区查询
如果应用需要频繁的跨分区查询,考虑以下优化:
重新设计数据模型:将经常关联的数据放在同一个分区(Co-location)。
使用全局二级索引:如果数据库支持,为常用查询创建全局索引。
异步聚合:将跨分区查询的结果异步聚合到专门的表中。
使用搜索引擎:将需要复杂查询的数据同步到 Elasticsearch,利用其分布式查询能力。
OLAP 分离:将分析查询路由到专门的 OLAP 系统(如 ClickHouse、Presto)。
实践 4:分区再平衡的最佳实践
执行分区再平衡时:
提前规划:预估再平衡的数据量和时间,选择业务低峰期执行。
限速控制:限制再平衡的网络带宽和 I/O 使用,避免影响正常业务。
增量迁移:分批迁移分区,每批完成后验证数据一致性。
监控影响:监控再平衡期间的系统性能,必要时暂停或回滚。
验证完整性:迁移完成后,验证源分区和目标分区的数据一致性。
本章小结
本章深入探讨了数据分区的核心概念和实践:
分区策略:
- 范围分区:适合范围查询,但可能数据倾斜。
- 哈希分区:数据分布均匀,但范围查询低效。
- 一致性哈希:支持高效的再平衡,适合动态集群。
分区与复制:
- 先分区再复制,每个分区维护多个副本。
- 分区 Leader 需要均匀分布在不同节点和机架上。
分区再平衡:
- 固定分区数策略:预先创建足够多的分区,增减节点时重新分配。
- 动态分区策略:根据数据量自动分裂和合并分区。
跨分区查询:
- 二级索引分为 Local Index(每个分区独立维护)和 Global Index(全局统一维护)。
- JOIN 操作在分布式系统中代价高昂,应尽量通过数据模型设计避免。
分区倾斜:
- 数据量倾斜可以通过合理的分区策略解决。
- 访问频率倾斜需要应用层缓存、随机前缀等策略。
分区是大规模分布式系统的基础。合理的分区设计可以使得系统高效扩展,不合理的分区设计则会导致热点、倾斜和性能问题。在设计分区策略时,需要综合考虑数据分布、访问模式、查询需求和运维复杂度。