06

分区

把大数据拆小

阅读量:2 · 预计 18 分钟读完

哈希分区范围分区二次分区再平衡
关联层级:L7 应用抽象
阅读进度5%

第六章 分区 - 把大数据拆小

导读

当数据量增长到单台机器无法容纳时,我们需要将数据分区(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 操作在分布式系统中代价高昂,应尽量通过数据模型设计避免。

分区倾斜:

- 数据量倾斜可以通过合理的分区策略解决。

- 访问频率倾斜需要应用层缓存、随机前缀等策略。

分区是大规模分布式系统的基础。合理的分区设计可以使得系统高效扩展,不合理的分区设计则会导致热点、倾斜和性能问题。在设计分区策略时,需要综合考虑数据分布、访问模式、查询需求和运维复杂度。