Ch.7: 把大数据集拆到多节点 — 分片策略决定写入均匀性, 路由方式决定请求走多远, 二级索引决定查询的代价
分片像图书馆分馆: 书太多一家放不下, 按类别区间分 (range: 科技馆/文学馆, 找一个范围的书只跑一家, 但新书全挤在"科技馆") 或按书名哈希分 (hash: 负载均匀, 但"找一个作者全部作品"要跑遍全城)。hash mod N的坑在于"每开一家新馆, 全城书重新分一遍"; 一致性哈希把书架摆成一个环, 新馆只接管环上相邻的一段。而热点是"金庸的书全城只有一份"的问题 — 不是分馆策略能解决的, 得把金庸的书印 100 份副本 (拆 key)。
# 何时分片: 数据量 > 单盘 / 写入 > 单机 / 网络带宽打满 # 何时不分: 一台顶得住 — 先垂直扩展再谈分片
# 倾斜检测: per-shard QPS 方差 > 10× → 有热点 # hot key 排行: 客户端采样 top-k 命中 key
-- 按月分片: 2026-09 分片接走全部新写 # SELECT WHERE ts BETWEEN ... → 只扫 2 个分片 ✓ # INSERT now() → 100% 砸一个片 ✗
# hash(user_id): 写入均匀 ✓ # WHERE user_id BETWEEN 100 AND 200 → 全分片扫 ✗
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 换节点
# 环上加节点 D: 只搬 (C→D) 弧段的 keys # 迁移量 ≈ 1/(N+1); 3→4 台 ≈ 20% (vs mod-N 75%)
# 无 vnode: 4 节点弧段可能差 3 倍 # 100 vnode/节点: 负载方差按 1/√100 收敛 → 均衡
# score(k,n) = hash(k, n); owner = argmax # 加节点 D: 只影响 score(D) 成为最大值的 keys
# 感知客户端: 周期拉路由表 + 本地计算 node = route_table.pick(hash(key)) # 零中转 # 路由表过期 → 打错 → 节点回 "moved" → 刷新
# 每秒随机挑 3 个邻居交换状态 → 指数扩散 # N 节点收敛 ≈ O(log N) 轮 — 但不是事务级一致
# 本地: 查 color=red → 全分片扇出 (读放大) # 全局: 查 color=red → 只打 red 所在分片 (写放大)
p = 0.01; shards = 30 # 1-(1-p)**30 → 0.26: 四分之一请求踩慢分片
# rebalance 带宽限速 50MB/s, 低峰窗口执行 # cell 故障: 只损 1/N 用户 — 隔离是设计出来的
车联网埋点按时间分片, "今天"分片 CPU 90% 其他 2% — 前缀打散把写摊到全集群。
# 错: key = "2026-09-26T10:33:21|dev8842" → 全砸今天片 # 对: 设备 ID 前缀 + 时间, 写入均匀分布 key = "dev8842|2026-09-26T10:33:21" # 查询代价: "查某设备一天轨迹" 变成 # 前缀扫 dev8842|2026-09-26* — 仍走该设备所在少数分片 # "全车队时间范围" 变 scatter — 交给分析库 (离线)
原则: 在线写路径优先均匀, 大范围扫描甩给列存/分析系统。
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 压力平滑过渡
事故教训: 分片函数是"永不变更"级别的决定, 上线前推演扩容。
明星发一条微博, 单 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 排行监控 (客户端采样), 触发自动加副本/拆分。
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 别变。
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 语言时优先路由层
节点转发最省事但最慢; 转发节点是"随便打"语义, 热点节点多一跳负担。
按颜色查车, 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/邻居干扰)。
一条车记录 8 个属性都要进索引 — 每次更新写 8 个索引分片。
# 车辆更新一次: # color=red, brand=BMW, city=BJ, price=30w ... 8 个词项分片 write_amplification = len(indexed_fields) # = 8× # 取舍表: # 高频查询字段 (city/brand) → 全局索引, 认写放大 # 低频字段 (color) → 留在本地索引 # 索引异步更新: 读到 1s 前的索引 — 业务可接受吗? 先问业务
经验: 全局索引字段数 ≤ 5, 且每个都有明确查询场景背书。
扩容迁移 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% 自动暂停迁移。
500 节点集群 gossip 收敛 30 秒 — 这 30 秒内路由表是"旧世界"。
# gossip 语义: 最终一致, 不是事务一致 # 新节点上线 → 30s 后全网才知道 → 期间请求被拒/转发 # 缓解: # client: 收到 Moved/Missing 就刷新路由表 # hinted handoff 期间写暂存 (呼应复制页) # 关键决策 (选主/迁移) 不走 gossip → 走共识 (etcd)
分层: 状态传播用 gossip, 一致性决策用共识 — 别混。
支付系统按用户哈希切 16 个 cell, 每个 cell 自带全套服务+存储。
# 路由: user_id % 16 → cell-0..15 (互不依赖) # cell-7 数据库故障: # 影响 = 1/16 用户, 其余 15 个 cell 无感知 # 恢复 = cell 内自愈, 不跨 cell 迁移 (隔离优先) # 代价: 跨 cell 功能 (全局排行榜) 需要 cell 间异步汇聚 # 全局服务单独部署, 读各 cell 汇总流 (类似 reverse ETL)
检验题: "一个 cell 拔电源, 用户影响面多大?" 答不上来就是没单元化。
# 错: server = nodes[hash(k) % len(nodes)] # 对: 哈希环 / 固定 1024 逻辑分片
# 错: rowkey = ts + dev_id # 对: rowkey = dev_id + ts
# 错: "hash 均匀所以没事" # 对: key#rep0..99 写扩散读聚合
# 错: create table 无 SPLITS # 对: SPLITS 按字母/数值分布预切
# 错: num_tokens=3 # 对: 100+ vnode, 异构机按算力加权
# 错: results = [s.query(r) for s in shards] 直接返回 # 对: heapq.merge(*results, key=sort_key)
# 错: 每分片各等 5s # 对: asyncio.wait(timeout=0.2) + partial 标记
# 错: 全字段建全局索引 # 对: 高频查询字段才进全局索引
# 错: 白天 full-speed 迁移 2TB # 对: 50MB/s 限速 + 02:00 窗口
# 错: 迁完直接切 # 对: md5 抽样比对一致 → 灰度切流
# 错: 路由表启动加载一次 # 对: 定期刷新 + moved 即刷
# 错: 用 gossip 状态做选主 # 对: 选主走 Raft/ZK
# 错: 按 user_id 分片却按 city 查 # 对: city 查询走全局索引/宽表
# 错: 10 台机器切 10 片 # 对: 1024 逻辑片映射 10 台
# 错: 跨 3 分片事务天天 2PC # 对: 同用户数据同分片, 跨片补偿
# 错: 所有 cell 连同一个 DB # 对: cell 内 DB + 跨 cell 异步汇总流
# 错: RF=3 就"分成 3 份" # 对: 12 分片 × 每片 3 副本 = 36 副本位
# 错: 分裂上限没配, 单片无限长 # 对: max_region_size=20GB 自动切
# 错: 迁移测试永远顺利 # 对: 迁移 50% 时 kill -9 一个节点
# 错: 集群均值 CPU 30% "健康" # 对: shard-7 CPU 95% 告警 + 拆分