设计一个分布式 KV 存储
提出问题
分布式 KV 存储是后端面试中出现频率最高的系统设计题之一,从 DynamoDB 到 Redis Cluster 再到 Cassandra,其底层架构思想贯穿了大部分现代分布式存储系统。面试官出这道题,考的不是你知道某个具体产品怎么用,而是你能否在数十到数百台机器上,把数据均衡分片、保证高可用、容忍节点故障这三个核心问题想清楚。生产中也一样——选型时你需要理解一致性哈希为什么比取模好,quorum 为什么能折中一致性和可用性,Hinted Handoff 和 Merkle 树分别解决什么场景的问题。
业务场景与需求推导
拿到题目先别急着画环,先把需求定下来:
场景假设:设计一个类似 DynamoDB 的 KV 存储,支撑电商系统的用户会话数据。数据量约 10 亿 key,单条 value 平均 2KB,QPS 约 50 万写 / 200 万读,99.9% 延迟 < 20ms。
推导出的需求:
| 维度 | 需求 | 影响设计 |
|---|---|---|
| 存储容量 | ~20TB 数据 | 需要分片,单机存不下 |
| 写入吞吐 | 50 万 QPS | 单机 IO 瓶颈,需要多副本异步写入 |
| 读取延迟 | P99.9 < 20ms | 不能跨多节点聚合读,读 quorum 要小 |
| 可用性 | 容忍 2 台机器故障 | 副本数 ≥ 3,自动故障转移 |
| 一致性 | 最终一致性可接受 | 不用强一致,选 AP 方向 |
需求定完了才开始选技术方案。
分片方案:一致性哈希 vs 取模
取模哈希的问题
最简单方案是 hash(key) % N,N 是节点数。扩容到 N+1 时,几乎全部 key 的映射都变了,要迁移的数据量是 (N / (N+1)) * 总量。N=10 时迁移 90% 数据。这在生产上不可接受——你不可能为了加一台机器把 20TB 数据全部重排一遍。
一致性哈希
一致性哈希把哈希值空间 [0, 2^32-1) 组织成一个环,每个节点落在一个哈希点上,每个 key 沿顺时针找到最近节点。新增节点时,只有该节点与前驱节点之间的数据需要迁移,迁移量约 1/N。
但标准一致性哈希有经典的数据倾斜问题:节点少时,环上节点分布不均匀,导致部分节点负载是其他节点的几十倍。Cassandra 早期版本(1.0 以前)在生产中就遇到过这个问题,某次扩缩容后 3 个节点承载了 80% 的流量。
解法:虚拟节点(Virtual Nodes)。每个物理节点在环上映射 100~200 个虚拟节点,大大降低负载不均的概率。Dynamo 论文推荐 100~200 个虚拟节点,Cassandra 默认 256 个虚拟节点(num_tokens = 256)。
// 一致性哈希 + 虚拟节点
class ConsistentHashRing {
private final TreeMap<Integer, String> ring = new TreeMap<>();
private final int virtualNodeCount;
private final HashFunction hashFn;
public ConsistentHashRing(int virtualNodeCount, HashFunction hashFn) {
this.virtualNodeCount = virtualNodeCount;
this.hashFn = hashFn;
}
public void addNode(String nodeId) {
for (int i = 0; i < virtualNodeCount; i++) {
int hash = hashFn.hash(nodeId + "#" + i);
ring.put(hash, nodeId);
}
}
public void removeNode(String nodeId) {
for (int i = 0; i < virtualNodeCount; i++) {
int hash = hashFn.hash(nodeId + "#" + i);
ring.remove(hash);
}
}
public String getNode(String key) {
int hash = hashFn.hash(key);
Map.Entry<Integer, String> entry = ring.ceilingEntry(hash);
if (entry == null) {
entry = ring.firstEntry();
}
return entry.getValue();
}
// 获取 key 的 N 个副本节点(顺时针取 N 个不同物理节点)
public List<String> getNodes(String key, int replicationFactor) {
Set<String> nodes = new LinkedHashSet<>();
int hash = hashFn.hash(key);
Map.Entry<Integer, String> entry = ring.ceilingEntry(hash);
if (entry == null) {
entry = ring.firstEntry();
}
NavigableMap<Integer, String> tailMap = ring.tailMap(entry.getKey(), true);
for (Map.Entry<Integer, String> e : tailMap.entrySet()) {
nodes.add(e.getValue());
if (nodes.size() >= replicationFactor) break;
}
// 如果环上节点不够,从头补
if (nodes.size() < replicationFactor) {
for (Map.Entry<Integer, String> e : ring.entrySet()) {
nodes.add(e.getValue());
if (nodes.size() >= replicationFactor) break;
}
}
return new ArrayList<>(nodes);
}
}注意:getNodes 里要跳过相同物理节点,否则同一个物理节点可能在环上被多次选中,导致副本数不够。这个问题在虚拟节点数多时尤其容易踩坑。
副本与 Quorum 机制
写入流程时序
客户端 → 协调节点(Coordinator)
│
├─ 第一步:一致性哈希定位副本节点 [A, B, C]
│
├─ 第二步:并行写入 A、B、C
│ ├─ A: 写入成功 ✓
│ ├─ B: 写入成功 ✓
│ └─ C: 超时 ✗(网络抖动)
│
└─ 第三步:协调节点收到 W 个确认(W=2),返回客户端成功每个 key 写入顺时针方向 N 个物理节点(N=3 典型值)。协调节点向所有 N 个副本发送写入请求,只要收到 W 个确认就算成功。
Quorum 配置策略
W + R > N 是读写重叠的条件,保证强一致。面试常问的配置组合:
| W | R | N | 特点 | 适用场景 | 延迟典型值 |
|---|---|---|---|---|---|
| 3 | 1 | 3 | 写入慢(3次确认),读取快(1次) | 电商商品详情,写少读多 | 写 5ms, 读 1ms |
| 1 | 3 | 3 | 写入快,读取慢(3次聚合) | 日志收集,写多读少 | 写 1ms, 读 5ms |
| 2 | 2 | 3 | 读写折中,Dynamo 默认 | 用户会话,读写均衡 | 各约 2ms |
| 1 | 1 | 3 | 弱一致性,W+R ≤ N | 关注数/阅读量,允许丢失 | 写 1ms, 读 1ms |
生产真实数据:DynamoDB 默认配置是 N=3, W=2, R=2,但在某些场景(如购物车会话)会降级到 W=1 以降低写入失败率。Cassandra 的 QUORUM 一致性级别对应 W = N/2 + 1。
一个真实踩坑:Quorum 不够 w 导致的脏读
某次线上事故,团队把 W 降到 1 来扛写入洪峰,R 保持 2,N=3,W+R=3=3,表面看满足强一致条件。但实际出现了脏读:原因是某个副本写入后立即崩溃,读请求落到另外两个还没来得及同步的副本上,读到了旧数据。
根因:W + R > N 保证的是只要写操作成功,后续读一定能读到最新数据。但如果写操作成功(W=1 确认)但唯一写入的副本在返回前崩溃了,后续读就永远读不到这次写入。这不是 W+R>N 的条件问题,而是持久化保证问题——写入返回成功前,数据必须 fsync 到磁盘,不能只写内存。
向量时钟与冲突检测
什么场景会产生冲突
时间线 t1: 客户端 A 写入 key=K, value=V1(节点 A 和 B 各存一份)
时间线 t2: 网络分区,节点 A 和 B 无法通信
时间线 t3: 客户端 A 写入 K=V2(只写到了 A)
时间线 t3: 客户端 B 写入 K=V3(只写到了 B)
时间线 t4: 网络恢复,A 和 B 发现 K 有两个不同版本此时 V2 和 V3 是并发写入,没有先后关系,无法自动合并。
向量时钟的版本比较
每个节点维护 {nodeId: version} 映射。写入时递增本节点版本号。
比较规则:假设 v1 = {A:2, B:1},v2 = {A:2, B:2}:
- v1 的每个分量 ≤ v2 的对应分量,且至少一个严格小于 → v1 是 v2 的祖先(可自动合并)
- 如果
v1[A]=2, v1[B]=1和v2[A]=1, v2[B]=2,两者互有大小 → 并发冲突,返回客户端解决
生产注意:向量时钟会无限增长。Dynamo 的做法是每 N 个版本合并一次(截断)。Cassandra 的解决方式是:如果节点数固定且不大,用时间戳+节点 ID 的简单方式替换。实践中,大多数场景不需要向量时钟,用时间戳选最大值(Last Write Wins, LWW)就够用,但代价是可能丢数据。
// 向量时钟的完整实现
class VectorClock implements Comparable<VectorClock> {
private final Map<String, Long> clocks = new HashMap<>();
public VectorClock increment(String nodeId) {
VectorClock vc = new VectorClock();
vc.clocks.putAll(this.clocks);
vc.clocks.merge(nodeId, 1L, Long::sum);
return vc;
}
public enum Relation { ANCESTOR, DESCENDANT, EQUAL, CONCURRENT }
public static Relation compare(VectorClock v1, VectorClock v2) {
boolean v1Older = true, v2Older = true;
Set<String> allKeys = new HashSet<>(v1.clocks.keySet());
allKeys.addAll(v2.clocks.keySet());
for (String k : allKeys) {
long c1 = v1.clocks.getOrDefault(k, 0L);
long c2 = v2.clocks.getOrDefault(k, 0L);
if (c1 > c2) v2Older = false;
if (c2 > c1) v1Older = false;
}
if (v1Older && !v2Older) return Relation.ANCESTOR;
if (v2Older && !v1Older) return Relation.DESCENDANT;
if (v1Older && v2Older) return Relation.EQUAL;
return Relation.CONCURRENT;
}
@Override
public int compareTo(VectorClock o) {
return switch (compare(this, o)) {
case ANCESTOR -> -1;
case DESCENDANT -> 1;
default -> 0;
};
}
}故障处理:Hinted Handoff 与 Merkle 树
Hinted Handoff 的流程
时间线: 写入 key=K 到节点 A,但 A 宕机
1. 协调节点检测到向 A 写入超时(比如 500ms 无响应)
2. 协调节点选择另一个节点 D 作为临时存储
3. 写入 D,并在 D 的元数据中记录 hint: "这条数据属于 A"
4. 客户端收到写入成功确认
5. 后台线程持续探测 A 是否恢复(每 30 秒一次)
6. A 恢复后,D 将 hint 数据回传给 A 并删除 hintHinted Handoff 的坑:
- D 在回传前也宕机了怎么办?→ hint 数据丢失,使用 Merkle 树做全量同步
- hint 数据在 D 上堆积太多 → 需要限制 hint 数量,超过阈值直接拒绝写入(返回不可用)
- 回传时机不对导致数据被覆盖 → 需要结合向量时钟或时间戳判断版本
Merkle 树同步
Hinted Handoff 只能处理几分钟级别的临时故障。节点离线数小时甚至数天后重新上线,必须做全量数据一致性校验。
Merkle 树是这样工作的:每个节点将数据按 key 范围分成固定大小的段(比如 256 个 key 一个段),每个段计算哈希值作为叶子节点,内部节点是子节点哈希的拼接后哈希。两节点交换根哈希:
节点 A: Merkle 根 = 0x3F8A... → 与节点 B 交换
节点 B: Merkle 根 = 0x3F8A... → 一致,跳过
节点 B: Merkle 根 = 0x7B22... → 不一致,递归向下找差异叶子
叶子 1: 0xA1 → 一致
叶子 2: 0xE3 → 一不一致 → 只同步这个段的数据复杂度:Merkle 树构建和比较的复杂度是 O(N log N),但只需传输差异数据,网络开销远小于全量对账。Dynamo 论文中每节点每 10 分钟触发一次 Merkle 树同步。
一个真实生产问题:某公司用 Cassandra 集群,节点宕机 6 小时后恢复,Merkle 树同步时发现差异段数太多,导致恢复期间该分区写入延迟飙到 200ms+。原因是该节点承载了热点数据,离线期间积累了 50 万条差异。解法:先把该节点的流量切走,等 Merkle 树同步完成后再恢复读写。
面试追问清单
面试官可能会追问这些,提前准备好:
Q1:为什么不用 Paxos/Raft 做一致性,而用 Quorum + 向量时钟? → 因为场景是 AP(高可用优先),Paxos/Raft 在领导者宕机时会有不可用窗口(选举时间),不适合低延迟写入场景。Dynamo 的设计哲学是"永远可写"。
Q2:Hinted Handoff 和 Merkle 树有什么区别? → Hinted Handoff 处理临时故障(秒-分钟级),Merkle 树处理长期离线(小时-天级)。Hinted Handoff 是"把数据暂时存到别处",Merkle 树是"定期对比一致性"。
Q3:一致性哈希的虚拟节点怎么调参? → 节点数少时(<10),虚拟节点数应该调大(200+);节点数多时(50+),100 个虚拟节点就够。虚拟节点数过多会导致路由表变大,每次查找需要 O(log V) 的二分查找,V 是虚拟节点总数。Cassandra 官方推荐 256。
Q4:读写延迟怎么算? → 写入延迟 = max(写 W 个副本的延迟),读取延迟 = max(读 R 个副本的延迟)。副本间有网络往返,所以单副本延迟约 1ms(本地机房内),W=2 时延迟约 1-2ms,W=3 时约 2-5ms(取决于最慢的那个副本)。
总结
分布式 KV 存储的关键设计链路:
| 问题 | 解法 | 参考实现 | 注意 |
|---|---|---|---|
| 数据分片 | 一致性哈希 + 虚拟节点 | Cassandra、DynamoDB | 虚拟节点数 100~200 避免数据倾斜 |
| 副本与一致性 | Quorum(W+R > N) | Dynamo、Cassandra | 先确认业务是 AP 还是 CP |
| 冲突检测 | 向量时钟 / LWW | Dynamo、Riak | 向量时钟会无限增长,需截断 |
| 临时故障 | Hinted Handoff | Dynamo、Cassandra | 限制 hint 数量,避免堆积 |
| 数据同步 | Merkle 树 | Dynamo、Cassandra | 每 10 分钟触发一次,避免恢复期间流量洪峰 |
| 成员管理 | Gossip 协议 | Cassandra、Consul | 下线的节点要等 gossip 传播,不是立刻感知 |
面试时建议从需求出发:先问存储容量、QPS、一致性要求,再推导出分片策略和副本数,最后落到容错方案。不要一上来就堆术语——把一致性哈希的「为什么」讲清楚,比背出所有 Dynamo 论文细节更打动人。
参考
参考:Dynamo 论文(Amazon's Dynamo);Cassandra 官方文档(Partitioners / Virtual Nodes);《Designing Data-Intensive Applications》Ch.6(Partitioning)/ Ch.9(Consistency and Consensus);ScyllaDB 博客关于 Virtual Nodes 的实践