DDIA · 分片与路由

Ch.7: 把大数据集拆到多节点 — 分片策略决定写入均匀性, 路由方式决定请求走多远, 二级索引决定查询的代价

扩容 1 台的迁移量: mod-N vs 一致性哈希 hash(key) mod N 3 节点: key→node = k mod 3 4 节点: = k mod 4 几乎全部 key 换主人 迁移 ≈ 75% N 变了 → 余数全变 → 全量重分布 期间缓存失效 + 迁移风暴 一致性哈希 (哈希环) 节点落在环上, key 顺时针找节点 A B C 新D 迁移 ≈ 1/(N+1) ≈ 20%: 只动环上两 点之间的 keys 请求路由: 这个 key 在哪个节点? ① 节点转发 随便打到一台, 不对就转发 (多一跳) ② 独立路由层 LBS/ZooKeeper 查表, 集中管理, 多一跳 ③ 分片感知客户端 客户端自己算路由, 零中转, 依赖驱动 集群元数据从哪来 gossip: 节点间流言传播 (Riak) — 去中心、最终一致、可能脑裂 coord 服务: ZooKeeper/etcd 集中维护路由表 — 强管理、依赖共识 (见 Ch.10) cell-based: 服务+存储自包含单元, 爆炸半径 = 1 个 cell skew / hot key: 负载不均是常态 时间戳 range 分片: 全部写入砸在"今天"的分片 2026-09 分片被写爆, 2025-12 分片在睡觉 — 相邻写集中是 range 的天性 解: key 前缀打散 — sensor_id + ts, 牺牲一点范围扫描效率 名人热点 key: 一条微博 1 亿粉丝 hash 再均匀也没用 — 单 key 就是全部负载 (名人效应) 解: 读侧本地缓存 + 写侧随机后缀拆 100 份 (读时合并) pre-splitting: 固定分片数 ≫ 节点数, key→shard 不变, 扩容只搬整 shard HBase 建表预切 split points / Kafka topic 分区数预留 — 分片是逻辑的, 节点是物理的 分片上的二级索引: 写便宜还是读便宜 本地 (文档) 索引 每分片只索引自己的文档 写: 便宜 (只碰本分片) 读: scatter/gather — 发到所有分片再合并 按颜色查车: 红/蓝/黑 各在 不同分片 → 全分片扇出 全局 (词项) 索引 索引本身按词项再分片 读: 单条件只打 1 个分片 写: 一条数据触多个分片 (车的所有属性词项都要更新) 异步更新 → 读到稍旧索引 (Elasticsearch 实际混合策略) scatter/gather 的尾部放大: 30 分片单分片 99 分位 = 100ms 整体至少一个分片慢的概率 = 1−(0.99)³⁰ = 26% — 再平衡/慢分片随时踩中

分片策略

  • • range: 范围扫描友好, 相邻写集中
  • • hash: 负载均匀, 丧失 range 能力
  • • 一致性哈希: 扩容只迁 ~1/N
  • • pre-split: 分片逻辑化, 扩容只搬整片

热点是常态

  • • 时间序列 range 分片必砸"今天"
  • • 名人 key: hash 救不了, 要拆+缓存
  • • 随机后缀写扩散, 读时合并
  • • per-shard 指标必须单独看

索引的代价转移

  • • 本地索引: 写便宜, 读 scatter/gather
  • • 全局索引: 读精准, 写触多分片
  • • 扇出 30 分片 = 尾延迟放大 26%
  • • 路由: 转发 / 路由层 / 感知客户端

💡 一句话理解

分片像图书馆分馆: 书太多一家放不下, 按类别区间分 (range: 科技馆/文学馆, 找一个范围的书只跑一家, 但新书全挤在"科技馆") 或按书名哈希分 (hash: 负载均匀, 但"找一个作者全部作品"要跑遍全城)。hash mod N的坑在于"每开一家新馆, 全城书重新分一遍"; 一致性哈希把书架摆成一个环, 新馆只接管环上相邻的一段。而热点是"金庸的书全城只有一份"的问题 — 不是分馆策略能解决的, 得把金庸的书印 100 份副本 (拆 key)。

🧠 必知必会 必考 & 必会

sharding 动机
数据集/写入超单机: 拆到多节点扩展写与存储。每条记录恰属一个分片; 小表别分片 — 复杂度白付。
# 何时分片: 数据量 > 单盘 / 写入 > 单机 / 网络带宽打满
# 何时不分: 一台顶得住 — 先垂直扩展再谈分片
skew / hot spot
负载不均叫偏斜, 被猛打的分片叫热点 (名人/今日时间戳)。分片策略解决"平均", 解决不了"单 key 巨热"。
# 倾斜检测: per-shard QPS 方差 > 10× → 有热点
# hot key 排行: 客户端采样 top-k 命中 key
range partitioning
按 key 区间切分: 范围扫描只打相关分片; 相邻 key 挤同片, 时间序列场景"今天"分片独热。
-- 按月分片: 2026-09 分片接走全部新写
# SELECT WHERE ts BETWEEN ... → 只扫 2 个分片 ✓
# INSERT now()                  → 100% 砸一个片 ✗
hash partitioning
hash(key) 决定分片: 负载均匀打散, 但丧失范围能力 — 相邻 key 四散, BETWEEN 变全分片扇出。
# hash(user_id): 写入均匀 ✓
# WHERE user_id BETWEEN 100 AND 200 → 全分片扫 ✗
hash mod N 灾难
节点数变化 → 几乎所有 key 的归属都变 → 全量迁移。分片数必须与节点数解耦。
def migrate_ratio(n_old, n_new):
    same = sum(1 for k in range(1000)
               if k % n_old == k % n_new)
    return 1 - same/1000
# 3→4: 75% 的 key 换节点
consistent hashing
节点放哈希环上, key 顺时针找最近节点: 加/减 1 节点只迁移 ~1/(N+1) 的 key — Dynamo/Cassandra 的底座。
# 环上加节点 D: 只搬 (C→D) 弧段的 keys
# 迁移量 ≈ 1/(N+1);  3→4 台 ≈ 20% (vs mod-N 75%)
virtual node
每物理节点在环上放 100+ 虚拟点: 热点自动摊开、异构机器按 vnode 数加权、下线平滑。
# 无 vnode: 4 节点弧段可能差 3 倍
# 100 vnode/节点: 负载方差按 1/√100 收敛 → 均衡
rendezvous hashing
会合哈希: 对每个 key, 所有节点按 hash(key,node) 排序取最高者 — 逐 key 分配, 不切 range, 增删也只迁 1/N。
# score(k,n) = hash(k, n);  owner = argmax
# 加节点 D: 只影响 score(D) 成为最大值的 keys
request routing 三式
节点转发 (随便打, 不对就转)、独立路由层 (集中查表)、分片感知客户端 (自己算) — 分别是"多一跳/集中依赖/驱动复杂"的取舍。
# 感知客户端: 周期拉路由表 + 本地计算
node = route_table.pick(hash(key))   # 零中转
# 路由表过期 → 打错 → 节点回 "moved" → 刷新
gossip protocol
流言式传播集群状态: 去中心、最终一致、抗单点; 代价是收敛延迟与潜在脑裂 — 元数据"够新"而非"最新"。
# 每秒随机挑 3 个邻居交换状态 → 指数扩散
# N 节点收敛 ≈ O(log N) 轮 — 但不是事务级一致
本地 vs 全局二级索引
本地 (文档) 索引随数据同片: 写便宜, 查询 scatter/gather; 全局 (词项) 索引按索引键再分片: 读精准, 写触多分片。
# 本地: 查 color=red → 全分片扇出 (读放大)
# 全局: 查 color=red → 只打 red 所在分片 (写放大)
scatter/gather 尾部放大
30 个分片并发查, 单分片 1% 慢 → 整体 26% 慢 (1−0.99³⁰); 必须并发扇出 + 整体超时 + 部分结果降级。
p = 0.01; shards = 30
# 1-(1-p)**30 → 0.26: 四分之一请求踩慢分片
rebalance 与 cell
再平衡 = 分片在节点间迁移, 重负载操作宜限速/人工窗口; cell-based 把服务+存储组成自包含单元, 故障爆炸半径=1 cell。
# rebalance 带宽限速 50MB/s, 低峰窗口执行
# cell 故障: 只损 1/N 用户 — 隔离是设计出来的

🏭 生产实战 real world

场景 1 · 按天 range 分片写热点: 前缀打散改造

车联网埋点按时间分片, "今天"分片 CPU 90% 其他 2% — 前缀打散把写摊到全集群。

# 错: key = "2026-09-26T10:33:21|dev8842"  → 全砸今天片

# 对: 设备 ID 前缀 + 时间, 写入均匀分布
key = "dev8842|2026-09-26T10:33:21"

# 查询代价: "查某设备一天轨迹" 变成
#   前缀扫 dev8842|2026-09-26* — 仍走该设备所在少数分片
# "全车队时间范围" 变 scatter — 交给分析库 (离线)

原则: 在线写路径优先均匀, 大范围扫描甩给列存/分析系统。

场景 2 · mod-N 扩容灾难: 缓存集体失效事故

memcached 从 10 台扩到 11 台, 90%+ key 换节点, 当晚缓存命中率归零打爆 DB。

def hit_ratio_after_scale(n_old, n_new, keys=10_000):
    hit = sum(1 for k in keys if k % n_old == k % n_new)
    return hit / len(keys)

# 10→11: 命中率 ≈ 9% — 相当于缓存被清空

# 修复: 一致性哈希 (ketama), 迁移 ≈ 1/11 ≈ 9% 的 key
# 且迁移期间双读 (旧+新), DB 压力平滑过渡

事故教训: 分片函数是"永不变更"级别的决定, 上线前推演扩容。

场景 3 · 名人热点 key: 拆分 + 本地缓存组合拳

明星发一条微博, 单 key 80 万 QPS, 所在分片 CPU 100% — 分片策略无解, 要拆 key。

# 写侧: 拆 100 份随机后缀
for i in range(100):
    write(f"post:12345:rep:{i}", counters[i])

# 读侧: 先本地缓存 (1s TTL), 未命中并发拉多份求和
cached = local_cache.get("post:12345")
if cached is None:
    parts = fanout_read([f"post:12345:rep:{i}" for i in range(100)])
    cached = sum(parts)
# 单 key 80万 qps → 每分片 8千 qps, 分片 CPU 回落 12%

识别: 热点 key 排行监控 (客户端采样), 触发自动加副本/拆分。

场景 4 · 预分裂: 新集群上线不倾斜

HBase 新表没预切, 全部数据先挤一个 RegionServer, 再慢慢自动分裂。

# 建表时按 rowkey 分布预切 split points
# 用户 ID 分布: 按字母频率切 (A~C/D~H/...)
create 'user_events', 'cf',
  SPLITS => ['d', 'h', 'm', 's', 'x']

# 或者: 固定 1024 逻辑分片 > 20 台机器
# key→shard 映射永不变, 扩容只搬整 shard (rebalance)

口诀: 分片数按终态规划, 不按当前机器数 — 节点会加, key→shard 别变。

场景 5 · 路由三式选型: 感知客户端零中转

P99 敏感的 KV 集群, 每多一跳多 0.5ms — 客户端感知路由最省。

# 路由层方案: client → LBS → shard (2 跳)
# 感知客户端:  client → shard  (1 跳)

class ShardClient:
    def get(self, key):
        node = self.table.pick(key)      # 本地路由表
        try:
            return node.get(key)
        except Moved:
            self.table.refresh()         # 路由表过期
            return self.table.pick(key).get(key)
# 代价: 每种语言都要驱动; 团队 < 3 语言时优先路由层

节点转发最省事但最慢; 转发节点是"随便打"语义, 热点节点多一跳负担。

场景 6 · 本地二级索引扇出: 并发 + 超时 + 部分结果

按颜色查车, 30 分片扇出 — 必须并发、限时、允许部分结果。

async def query_by_color(color):
    tasks = [shard.query(color) for shard in shards]
    done, pending = await asyncio.wait(
        tasks, timeout=0.2)                # 整体 200ms 封顶
    rows = [r for t in done for r in t.result()]
    if pending:
        metrics.incr("scatter_partial")    # 部分结果标记
    return rows
# 无超时的扇出 = 最慢分片决定 P99 (呼应尾延迟放大)

监控: partial 命中率 > 1% → 排查慢分片 (大分片/GC/邻居干扰)。

场景 7 · 全局二级索引: 写放大的量化与取舍

一条车记录 8 个属性都要进索引 — 每次更新写 8 个索引分片。

# 车辆更新一次:
#   color=red, brand=BMW, city=BJ, price=30w ... 8 个词项分片
write_amplification = len(indexed_fields)   # = 8×

# 取舍表:
#   高频查询字段 (city/brand) → 全局索引, 认写放大
#   低频字段 (color)          → 留在本地索引
# 索引异步更新: 读到 1s 前的索引 — 业务可接受吗? 先问业务

经验: 全局索引字段数 ≤ 5, 且每个都有明确查询场景背书。

场景 8 · rebalance 窗口: 限速 + 低峰 + 双写校验

扩容迁移 2TB 数据, 高峰执行把在线 P99 顶爆 — 迁移是"重负载操作"要当发布管理。

# rebalance SOP:
# 1) 限速: 迁移带宽 50MB/s (在线 IO 的 10%)
# 2) 窗口: 02:00~06:00, 避开大促周
# 3) 迁移中双写: 旧分片为主, 新分片校验 md5 一致后切流
# 4) 切流后旧分片保留 24h 再删 (回滚保险)

rebalance --rate-limit 50MB --start 02:00 --drain-after-check

告警: 迁移期间在线集群 P99 基线偏移 > 20% 自动暂停迁移。

场景 9 · gossip 参数: 收敛速度与脑裂风险

500 节点集群 gossip 收敛 30 秒 — 这 30 秒内路由表是"旧世界"。

# gossip 语义: 最终一致, 不是事务一致
# 新节点上线 → 30s 后全网才知道 → 期间请求被拒/转发

# 缓解:
#   client: 收到 Moved/Missing 就刷新路由表
#   hinted handoff 期间写暂存 (呼应复制页)
# 关键决策 (选主/迁移) 不走 gossip → 走共识 (etcd)

分层: 状态传播用 gossip, 一致性决策用共识 — 别混。

场景 10 · cell-based 单元化: 爆炸半径=1/16

支付系统按用户哈希切 16 个 cell, 每个 cell 自带全套服务+存储。

# 路由: user_id % 16 → cell-0..15 (互不依赖)
# cell-7 数据库故障:
#   影响 = 1/16 用户, 其余 15 个 cell 无感知
#   恢复 = cell 内自愈, 不跨 cell 迁移 (隔离优先)

# 代价: 跨 cell 功能 (全局排行榜) 需要 cell 间异步汇聚
# 全局服务单独部署, 读各 cell 汇总流 (类似 reverse ETL)

检验题: "一个 cell 拔电源, 用户影响面多大?" 答不上来就是没单元化。

⚠️ 编码注意与常见坑 pitfalls

坑 1 · mod-N 直接扩容 — 75% key 换主人, 缓存全失效. 原因: 分片与节点数耦合. 正解: 一致性哈希/固定逻辑分片。
# 错: server = nodes[hash(k) % len(nodes)]
# 对: 哈希环 / 固定 1024 逻辑分片
坑 2 · 时间戳 range 分片 — 全部写入砸"今天"分片. 原因: 顺序 key 天然集中. 正解: 前缀打散或 hash+时间二级索引。
# 错: rowkey = ts + dev_id
# 对: rowkey = dev_id + ts
坑 3 · 热点 key 不拆 — 名人微博单 key 80万 qps 打挂分片. 原因: 以为 hash 能解决单 key. 正解: 随机后缀拆 N 份 + 读合并 + 本地缓存。
# 错: "hash 均匀所以没事"
# 对: key#rep0..99 写扩散读聚合
坑 4 · 不做预分裂 — 新表全挤一个 RegionServer 自动分裂半年. 原因: 建表图省事. 正解: 按 key 分布预切 split points。
# 错: create table 无 SPLITS
# 对: SPLITS 按字母/数值分布预切
坑 5 · vnode 太少 — 3 个 vnode/节点, 负载方差 3 倍. 原因: 默认值直接用. 正解: 每物理节点 100+ vnode 或按容量加权。
# 错: num_tokens=3
# 对: 100+ vnode, 异构机按算力加权
坑 6 · range 跨片不合并 — 各分片返回乱序, 分页错乱. 原因: 只做了 scatter 没做 merge. 正解: 归并排序 + 统一游标。
# 错: results = [s.query(r) for s in shards] 直接返回
# 对: heapq.merge(*results, key=sort_key)
坑 7 · 扇出无整体超时 — 最慢分片决定 P99. 原因: 逐分片超时. 正解: 整体 deadline + 部分结果降级。
# 错: 每分片各等 5s
# 对: asyncio.wait(timeout=0.2) + partial 标记
坑 8 · 全局索引写放大低估 — 8 字段索引 = 每次 update 写 8 片. 原因: 只算主数据. 正解: 索引字段 ≤5 且逐一背书。
# 错: 全字段建全局索引
# 对: 高频查询字段才进全局索引
坑 9 · rebalance 高峰执行 — 迁移 IO 与在线流量抢盘. 原因: 当成轻操作. 正解: 限速 + 低峰窗口 + 基线偏移熔断。
# 错: 白天 full-speed 迁移 2TB
# 对: 50MB/s 限速 + 02:00 窗口
坑 10 · 迁移期双写不一致 — 旧片新片数据分叉, 切流后数据错. 原因: 没有校验. 正解: 双写 + 校验和一致后才切流。
# 错: 迁完直接切
# 对: md5 抽样比对一致 → 灰度切流
坑 11 · 路由表不同步 — 请求打到已迁移节点. 原因: 表更新滞后. 正解: Moved 响应触发刷新 + 版本号校验。
# 错: 路由表启动加载一次
# 对: 定期刷新 + moved 即刷
坑 12 · gossip 当强一致 — 30s 收敛窗口内决策基于旧世界. 原因: 语义混淆. 正解: 状态传播 gossip, 关键决策共识 (etcd)。
# 错: 用 gossip 状态做选主
# 对: 选主走 Raft/ZK
坑 13 · 分片键选错 — 查询不带分片键 = 全分片扫. 原因: key 设计没对齐查询. 正解: 高频查询必带分片键, 否则进二级索引/分析库。
# 错: 按 user_id 分片却按 city 查
# 对: city 查询走全局索引/宽表
坑 14 · 分片数 = 节点数 — 每次扩容都重分布. 原因: 物理思维. 正解: 逻辑分片数 ≫ 节点数, 扩容只搬整片。
# 错: 10 台机器切 10 片
# 对: 1024 逻辑片映射 10 台
坑 15 · 跨分片事务硬上 — 每笔订单 2PC, 协调器成瓶颈. 原因: 分片后才想事务. 正解: 分片键=事务边界, 跨片走 saga/异步。
# 错: 跨 3 分片事务天天 2PC
# 对: 同用户数据同分片, 跨片补偿
坑 16 · cell 间强同步 — 单元化又强耦合, 爆炸半径失效. 原因: 图省事共享库. 正解: cell 自包含, 跨 cell 只异步汇聚。
# 错: 所有 cell 连同一个 DB
# 对: cell 内 DB + 跨 cell 异步汇总流
坑 17 · 分片与复制混淆 — 把"副本数"当"分片数"配置. 原因: 概念纠缠. 正解: 分片=横向切块, 复制=每片竖着拷贝, 两者正交。
# 错: RF=3 就"分成 3 份"
# 对: 12 分片 × 每片 3 副本 = 36 副本位
坑 18 · 大分片不分裂 — 单片 500GB, 迁移一次 10 小时. 原因: 分裂策略缺失. 正解: 片大小阈值自动分裂 (如 20GB)。
# 错: 分裂上限没配, 单片无限长
# 对: max_region_size=20GB 自动切
坑 19 · 迁移+故障并发没测 — 迁移中节点挂, 状态机卡死. 原因: 只测单故障. 正解: 迁移注入故障演练 (Jepsen 式)。
# 错: 迁移测试永远顺利
# 对: 迁移 50% 时 kill -9 一个节点
坑 20 · 忽略 per-shard 监控 — 全局平均掩盖单分片冒烟. 原因: 只有聚合大盘. 正解: per-shard QPS/P99/容量 + 倾斜告警。
# 错: 集群均值 CPU 30% "健康"
# 对: shard-7 CPU 95% 告警 + 拆分