ZooKeeper 数据同步机制深度解析:Leader 选举后的数据对齐

ZooKeeper 数据同步机制深度解析:Leader 选举后的数据对齐

    • 一、数据同步概述
      • 1.1 为什么需要数据同步?
      • 1.2 数据同步的核心目标
    • 二、同步前的准备工作
      • 2.1 Follower 连接 Leader
      • 2.2 关键数据结构
    • 三、四种同步策略详解
      • 3.1 DIFF 同步(差异化同步)
      • 3.2 TRUNC 同步(截断同步)
      • 3.3 SNAP 同步(全量快照同步)
      • 3.4 TRUNC+DIFF 同步(复合同步)
    • 四、同步完成确认
    • 五、完整数据同步示例
      • 5.1 场景一:3节点集群数据同步
      • 5.2 命令行验证
    • 六、数据同步的源码分析
      • 6.1 Leader 端决策逻辑
      • 6.2 Follower 端处理逻辑
    • 七、总结
      • 7.1 四种同步策略对比
      • 7.2 数据同步完整流程图
      • 7.3 一句话总结

🌺The Begin🌺点点关注,收藏不迷路🌺

摘要:Leader 选举完成后,新当选的 Leader 必须与集群中所有 Follower 进行数据同步,确保每个节点的数据视图一致。这是 ZooKeeper 保证数据一致性的关键步骤。本文将深入剖析 Leader 选举后的数据同步机制,详细讲解四种同步策略(DIFF、TRUNC、SNAP)的触发条件和执行流程,通过源码分析和实战示例,帮助读者全面理解这一核心机制。

一、数据同步概述

1.1 为什么需要数据同步?

Leader 选举完成后,集群中各个节点的数据状态可能不一致:

选举完成后的集群状态

新 Leader
maxZxid=120

Follower1
lastZxid=100

Follower2
lastZxid=120

Follower3
lastZxid=90

Follower4
lastZxid=130

节点 数据状态 需要同步
Leader maxZxid=120 基准节点
Follower1 lastZxid=100 缺失事务 101-120
Follower2 lastZxid=120 已同步,无需操作
Follower3 lastZxid=90 缺失事务 91-120,且落后较多
Follower4 lastZxid=130 有多余事务(可能是旧 Leader)

1.2 数据同步的核心目标

目标 说明
一致性 所有节点最终看到相同的数据视图
完整性 已提交的事务不能丢失
正确性 未提交的事务不能出现
高效性 根据落后程度选择最优同步策略

二、同步前的准备工作

2.1 Follower 连接 Leader

选举完成后,Follower 会主动连接新 Leader,并发送自己的状态信息:

新 Leader

Follower

新 Leader

Follower

6. 收集所有 Follower 的 lastZxid

7. 确定 minZxid 和 maxZxid

8. 为每个 Follower 决定同步策略

1. 建立 TCP 连接

2. 发送 FOLLOWERINFO(lastZxid)

3. 记录 Follower 的 lastZxid

4. 返回 LEADERINFO(newEpoch)

5. 发送 ACKEPOCH(lastZxid)

2.2 关键数据结构

Leader 端维护的同步信息:

// Leader.java - 维护 Follower 状态
public class Leader {
    // 记录每个 Follower 的 lastZxid
    private Map<Long, Long> followerLastZxid = new ConcurrentHashMap<>();
    // 已提交事务日志缓存
    private LinkedList<Proposal> committedLog = new LinkedList<>();
    // 最小 ZXID(committedLog 中最旧的事务)
    private long minCommittedLog;
    // 最大 ZXID(committedLog 中最新的一个)
    private long maxCommittedLog;
    public void processAckEpoch(QuorumPacket qp, LearnerHandler handler) {
        long lastZxid = qp.getZxid();
        long sid = handler.getSid();
        // 记录 Follower 的 lastZxid
        followerLastZxid.put(sid, lastZxid);
        // 更新全局最小和最大 ZXID
        updateZxidRange(lastZxid);
        // 如果收到过半 ACK,开始决定同步策略
        if (followerLastZxid.size() > half) {
            decideSyncStrategyForAll();
        }
    }
}

三、四种同步策略详解

Leader 根据 Follower 上报的 lastZxid 与自身 committedLog 的范围比较,决定采用哪种同步方式:

等于 Leader.maxZxid

小于 Leader.minZxid

在 Leader.minZxid 和 maxZxid 之间

大于 Leader.maxZxid

收到 Follower 的 lastZxid

比较 lastZxid

无需同步
直接进入服务

SNAP 全量同步
发送完整快照

DIFF 增量同步
发送缺失事务

TRUNC 截断同步
回滚多余事务

同步完成

3.1 DIFF 同步(差异化同步)

触发条件Leader.minZxid ≤ Follower.lastZxid ≤ Leader.maxZxid

适用场景:Follower 只落后少量事务,是最常见的同步方式。

Leader (min=90, max=120)

Follower (lastZxid=100)

Leader (min=90, max=120)

Follower (lastZxid=100)

1. DIFF 指令

2. PROPOSAL(zxid=101)

2. PROPOSAL(zxid=102)

2. PROPOSAL(zxid=103)

…(直到 zxid=120)

3. 写入事务日志

4. 等待 COMMIT

5. COMMIT(zxid=101)

5. COMMIT(zxid=102)

5. COMMIT(zxid=103)

…(按顺序提交)

6. 应用到内存

7. 同步完成

源码实现

// LearnerHandler.java - DIFF 同步
private void syncFollower(long lastZxid) {
    // 获取从 lastZxid+1 开始的所有提案
    Iterator<Proposal> it = leader.getCommittedLog()
                                 .iterator(lastZxid + 1);
    // 发送 DIFF 指令
    queuePacket(new QuorumPacket(Leader.DIFF, 0, null, null));
    // 按顺序发送缺失的提案
    while (it.hasNext()) {
        Proposal p = it.next();
        queuePacket(p.packet);
    }
    // 发送 COMMIT 消息
    it = leader.getCommittedLog().iterator(lastZxid + 1);
    while (it.hasNext()) {
        Proposal p = it.next();
        queuePacket(new QuorumPacket(Leader.COMMIT, p.packet.getZxid(), null, null));
    }
}

3.2 TRUNC 同步(截断同步)

触发条件Follower.lastZxid > Leader.maxZxid

适用场景:Follower 是旧 Leader 恢复后重新加入,有多余的事务需要删除。

Leader (max=120)

Follower (lastZxid=130)

Leader (max=120)

Follower (lastZxid=130)

告诉 Follower 回滚到 zxid=120

1. TRUNC 指令

2. 删除 zxid 121-130 的事务日志

3. 回滚内存数据

4. 回滚完成确认

源码实现

// LearnerHandler.java - TRUNC 同步
private void handleTruncSync(long lastZxid) {
    long leaderZxid = leader.getMaxCommittedLog();
    // 发送 TRUNC 指令
    QuorumPacket trunc = new QuorumPacket(Leader.TRUNC, leaderZxid, null, null);
    queuePacket(trunc);
    LOG.info("Sending TRUNC to follower, follower lastZxid={}, leader lastZxid={}",
             Long.toHexString(lastZxid), Long.toHexString(leaderZxid));
}

3.3 SNAP 同步(全量快照同步)

触发条件Follower.lastZxid < Leader.minZxid

适用场景:Follower 落后太多(超过 committedLog 缓存范围),或新节点首次加入。

Leader (min=90, max=120)

Follower (lastZxid=50)

Leader (min=90, max=120)

Follower (lastZxid=50)

2. 生成内存数据快照

1. SNAP 指令

3. 序列化 DataTree

4. 发送快照数据

5. 清空当前数据

6. 加载快照

7. 快照加载完成

源码实现

// LearnerHandler.java - SNAP 同步
private void handleSnapSync() {
    // 1. 获取内存数据快照
    ByteBuffer snapshot = leader.getZKDatabase().getSnapshot();
    // 2. 发送 SNAP 指令
    queuePacket(new QuorumPacket(Leader.SNAP, 
                                  leader.getMaxCommittedLog(), 
                                  snapshot.array(), 
                                  null));
    // 3. 发送 UPTODATE 确认
    queuePacket(new QuorumPacket(Leader.UPTODATE, -1, null, null));
}

3.4 TRUNC+DIFF 同步(复合同步)

触发条件:Follower 既有缺失的事务,又有多余的事务

适用场景:复杂的网络分区故障后恢复。

Leader (min=95, max=120)

Follower (lastZxid=125)

Leader (min=95, max=120)

Follower (lastZxid=125)

2. 回滚到 zxid=120

1. TRUNC 指令

截断到 zxid=120

3. DIFF 指令

4. PROPOSAL(zxid=96-120)

5. COMMIT(zxid=96-120)

6. 应用事务

四、同步完成确认

无论采用哪种同步策略,最后都需要完成确认:

Leader

Follower

Leader

Follower

同步完成

收到过半 ACK

状态变更为 FOLLOWING

状态为 LEADING

集群恢复服务

NEWLEADER

如果是 SNAP,生成快照

ACK

UPTODATE

五、完整数据同步示例

5.1 场景一:3节点集群数据同步

Leader (max=120)

Follower2 (last=120)

Follower1 (last=100)

Leader (max=120)

Follower2 (last=120)

Follower1 (last=100)

Leader 选举完成

par

[并行同步]

收到过半 ACK

同步完成

FOLLOWERINFO(100)

决定 DIFF 同步

DIFF

PROPOSAL(101-120)

COMMIT(101-120)

ACK

FOLLOWERINFO(120)

已同步,无需操作

UPTODATE

ACK

NEWLEADER

NEWLEADER

ACK

ACK

UPTODATE

UPTODATE

5.2 命令行验证

# 1. 启动集群(3节点)
./bin/zkServer.sh start
# 2. 创建一些节点
./bin/zkCli.sh -server localhost:2181 create /sync-test "data"
./bin/zkCli.sh -server localhost:2181 create /sync-test/node1 "node1"
./bin/zkCli.sh -server localhost:2181 create /sync-test/node2 "node2"
# 3. 停止 Leader
./bin/zkServer.sh stop
# 4. 重新启动原 Leader
./bin/zkServer.sh start
# 5. 查看日志中的数据同步过程
tail -f logs/zookeeper.out | grep -E "SYNC|DIFF|TRUNC|SNAP"
# 预期输出
2024-01-01 10:00:01 - Synchronizing with follower 2
2024-01-01 10:00:01 - Using DIFF sync for follower 2
2024-01-01 10:00:02 - Sending proposals from 101 to 120
2024-01-01 10:00:03 - UPTODATE sent to follower 2

六、数据同步的源码分析

6.1 Leader 端决策逻辑

// Leader.java - 决定同步策略
public class Leader {
    public void decideSyncStrategy(LearnerHandler handler, long lastZxid) {
        long minZxid = getMinCommittedLog();
        long maxZxid = getMaxCommittedLog();
        if (lastZxid == maxZxid) {
            // 已同步,直接发送 UPTODATE
            handler.queuePacket(new QuorumPacket(Leader.UPTODATE, -1, null, null));
        } else if (lastZxid > maxZxid) {
            // Follower 有多余事务,需要截断
            LOG.info("Follower has newer data, sending TRUNC to {}", handler.getSid());
            handler.queuePacket(new QuorumPacket(Leader.TRUNC, maxZxid, null, null));
        } else if (lastZxid < minZxid) {
            // Follower 落后太多,需要全量同步
            LOG.info("Follower is too old, sending SNAP to {}", handler.getSid());
            handleSnapSync(handler);
        } else {
            // Follower 落后少量事务,增量同步
            LOG.info("Follower is slightly behind, sending DIFF to {}", handler.getSid());
            handleDiffSync(handler, lastZxid);
        }
    }
}

6.2 Follower 端处理逻辑

// Learner.java - Follower 处理同步
public class Learner {
    public void syncWithLeader(long newLeaderZxid) throws Exception {
        // 接收 Leader 的同步指令
        QuorumPacket qp = leaderServer.readPacket();
        if (qp.getType() == Leader.DIFF) {
            // 差异化同步
            LOG.info("Received DIFF from leader");
            snapshotNeeded = false;
        } else if (qp.getType() == Leader.TRUNC) {
            // 截断同步
            long zxid = qp.getZxid();
            zk.getZKDatabase().truncateLog(zxid);
            LOG.info("Truncated log to {}", Long.toHexString(zxid));
        } else if (qp.getType() == Leader.SNAP) {
            // 全量同步
            LOG.info("Received SNAP from leader");
            snapshotNeeded = true;
            // 应用快照
            zk.getZKDatabase().clear();
            zk.getZKDatabase().deserializeSnapshot(leaderServer.getInputStream());
        }
        // 同步完成后确认
        leaderServer.writePacket(new QuorumPacket(Leader.ACK, 0, null, null));
    }
}

七、总结

7.1 四种同步策略对比

策略 触发条件 操作 适用场景 性能
DIFF minZxid ≤ lastZxid ≤ maxZxid 增量发送缺失事务 短暂故障后恢复
TRUNC lastZxid > maxZxid 回滚多余事务 旧 Leader 恢复
SNAP lastZxid < minZxid 发送完整快照 新节点或长时间故障
TRUNC+DIFF 复杂情况 先回滚后增量 网络分区恢复 中等

7.2 数据同步完整流程图

Leader 选举完成

等待 Follower 连接

收到 FOLLOWERINFO

记录 lastZxid

收到过半 ACK?

为每个 Follower 决定策略

DIFF 同步

TRUNC 同步

SNAP 同步

无需同步

发送 NEWLEADER

等待 Follower ACK

收到过半 ACK?

发送 UPTODATE

集群就绪

7.3 一句话总结

ZooKeeper 在 Leader 选举后,通过 DIFF、TRUNC、SNAP 三种核心同步策略,根据每个 Follower 的 lastZxid 与 Leader 的日志范围进行比对,以最小开销实现所有节点的数据对齐,确保分布式数据的一致性和完整性。

在这里插入图片描述

🌺The End🌺点点关注,收藏不迷路🌺
© 版权声明

相关文章