Skip to content

设计一个分布式 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)。

java
// 一致性哈希 + 虚拟节点
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 是读写重叠的条件,保证强一致。面试常问的配置组合:

WRN特点适用场景延迟典型值
313写入慢(3次确认),读取快(1次)电商商品详情,写少读多写 5ms, 读 1ms
133写入快,读取慢(3次聚合)日志收集,写多读少写 1ms, 读 5ms
223读写折中,Dynamo 默认用户会话,读写均衡各约 2ms
113弱一致性,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 有两个不同版本

此时 V2V3并发写入,没有先后关系,无法自动合并。

向量时钟的版本比较

每个节点维护 {nodeId: version} 映射。写入时递增本节点版本号。

比较规则:假设 v1 = {A:2, B:1},v2 = {A:2, B:2}:

  • v1 的每个分量 ≤ v2 的对应分量,且至少一个严格小于 → v1 是 v2 的祖先(可自动合并)
  • 如果 v1[A]=2, v1[B]=1v2[A]=1, v2[B]=2,两者互有大小 → 并发冲突,返回客户端解决

生产注意:向量时钟会无限增长。Dynamo 的做法是每 N 个版本合并一次(截断)。Cassandra 的解决方式是:如果节点数固定且不大,用时间戳+节点 ID 的简单方式替换。实践中,大多数场景不需要向量时钟,用时间戳选最大值(Last Write Wins, LWW)就够用,但代价是可能丢数据。

java
// 向量时钟的完整实现
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 并删除 hint

Hinted 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
冲突检测向量时钟 / LWWDynamo、Riak向量时钟会无限增长,需截断
临时故障Hinted HandoffDynamo、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 的实践

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。