Skip to content

第 8 章 主从复制

学习目标:彻底搞清楚 Redis 主从复制的必要性、拓扑形态、全量同步五步流程、增量同步「断点续传」机制、命令传播的异步本质、主从延迟的来源、读写分离的常见翻车场景。看完之后能完整画出「从节点连上主节点到稳定运行」的全流程,并能回答大厂面试官追问「PSYNC 的三种返回值是什么?」「复制积压缓冲区设多大?」「写完立刻读读到旧值怎么办?」这类问题。


8.1 为什么需要主从复制

8.1.1 单机 Redis 的三大瓶颈

到目前为止我们一直把 Redis 当一个孤零零的进程在用——一台机器、一个 127.0.0.1:6379。这种部署能撑得住玩具项目,但只要业务一上量,三个问题立刻浮现

┌────────────────────────────────────────────────────────────┐
│   单机 Redis 的三大死穴                                      │
├────────────────────────────────────────────────────────────┤
│  ① 高可用差:进程挂了 / 机器宕机 → 业务直接不可用              │
│  ② 读扩展难:所有读请求都打一个实例,CPU/带宽都顶不住         │
│  ③ 数据安全弱:只有一份数据,磁盘坏了 RDB+AOF 一起没           │
└────────────────────────────────────────────────────────────┘

主从复制(Replication)就是用来解决这三个问题的最基础设施

问题主从复制如何解决
读扩展写命令只走主节点,读请求可以分摊到 N 个从节点
容灾备份数据在多台机器上各有一份 → 物理隔离的备份
故障转移基础主节点挂了,从节点还活着 → 哨兵 / Cluster 自动切换的前提

8.1.2 生活类比:图书馆总馆与分馆

📚 生活类比:把 Redis 主节点想象成一座总馆(图书的唯一编目权威),从节点是分布在各社区的分馆

  • 新书入库(写请求)必须先送到总馆登记 → 写只走主
  • 总馆每登记一本书,就把书目和书的副本派人送到各分馆归档 → 命令传播
  • 读者借书(读请求)就近去任何一个分馆即可,不必都涌向总馆 → 读扩展
  • 万一总馆失火(主节点宕机),分馆里的副本仍在,能临时承担一部分总馆职能 → 容灾 + 故障转移基础
  • 但派送图书需要时间。读者刚听说总馆有新书,跑到分馆借却被告知「还没到货」 → 主从延迟

这一节先记住四个字:就近借书。这正是 Redis 读写分离的核心收益,也是后面 8.7 那些坑的根源。

8.1.3 几个绕不开的术语

术语含义
master主节点,唯一可写
replica / slave从节点(Redis 5.0 起官方文档统一改用 replica,配置项 slaveof 改名 replicaof,本章混用)
全量同步(Full Resync)从节点拿到主节点的完整 RDB 快照 + 期间命令
增量同步(Partial Resync)从节点断线重连后,只补缺失的命令
replid(runid)主节点的 40 位十六进制身份证
offset主节点累计发出的字节数(从节点累计接收并执行的字节数)
backlog复制积压缓冲区(环形 byte buffer),增量同步的「弹药库」

8.2 主从架构图

主从复制不是只能「一主一从」。常见拓扑有三种。

8.2.1 一主一从(最简单,做热备)

   ┌──────────┐                ┌──────────┐
   │  master  │ ──────────────►│ replica  │
   │ 6379     │   命令流推送     │ 6380     │
   └──────────┘                └──────────┘
        ▲                            ▲
       写                          读 / 备份
  • ✅ 配置最简单。
  • ✅ 主挂了从顶上,是哨兵的最小可用单元。
  • ❌ 只有 1 个从,读扩展能力有限。

8.2.2 一主多从(最常见)

                ┌──────────┐
                │  master  │  ←── 写
                └────┬─────┘
        ┌───────────┼───────────┐
        ▼           ▼           ▼
   ┌────────┐ ┌────────┐ ┌────────┐
   │replica1│ │replica2│ │replica3│  ←── 读
   └────────┘ └────────┘ └────────┘

  ✅ 简单直接,读扩展能力随从节点数线性增长
  ❌ 主节点压力随从节点数线性增长(每多一个从,主就要多发一份命令流)
  ❌ 主节点带宽也会成为瓶颈

8.2.3 级联复制 / 链式(M → S1 → S2)

   ┌──────────┐
   │  master  │
   └────┬─────┘

   ┌──────────┐    ←── 既是 master 的从,
   │ replica1 │       也是 replica2/3 的「父节点」
   └────┬─────┘
   ┌────┴────┐
   ▼         ▼
┌────────┐ ┌────────┐
│replica2│ │replica3│
└────────┘ └────────┘

  ✅ 减轻 master 的复制压力(master 只需要发给 replica1)
  ✅ 适合「跨地域复制」:主只发一份到异地中转节点,再分发
  ❌ 复制延迟叠加:replica3 的数据要经过两跳
  ❌ 中间节点挂了,下游全断
  ❌ 维护复杂度升高

💡 生产建议:除非主节点带宽 / CPU 真的成了瓶颈,优先用一主多从;级联复制只在跨地域、跨可用区或从节点数量极大(> 10)时考虑。

8.2.4 拓扑选型速查

拓扑适用场景不适用场景
一主一从小流量业务、做热备大量读请求
一主多从大多数线上业务主节点带宽 / CPU 已经接近瓶颈
级联复制跨地域 / 极多从对延迟敏感的业务

8.3 全量同步(首次同步)

8.3.1 什么时候走全量?

凡是**主节点判断「我没法给你做增量」**的情况,都会走全量:

  1. 从节点第一次接入(手里没有任何 replid 和 offset,发的是 PSYNC ? -1)。
  2. 从节点的 replid 与主节点对不上(说明这个从原来跟的是另一个主)。
  3. 从节点的 offset 太旧,已经被主节点的复制积压缓冲区覆盖了(环形 buffer 转了一圈)。

8.3.2 五步流程总览

   阶段 ①  从节点发起 PSYNC
   阶段 ②  主节点回 +FULLRESYNC <replid> <offset>
   阶段 ③  主节点 fork 子进程 BGSAVE 生成 RDB
           (同时把后续写命令塞入 per-replica「复制缓冲区」)
   阶段 ④  主节点把 RDB 通过 socket 传给从节点
   阶段 ⑤  从节点清空旧数据 → 加载 RDB → 重放复制缓冲区命令 → 进入命令传播

8.3.3 ASCII 时序图

  replica                                          master
    │                                                 │
 ①  │── PSYNC ? -1 ──────────────────────────────────►│
    │                                                 │
    │                              ┌──────────────────┴──────────────┐
    │                              │ 主节点判断:从节点是新连接,必须全量 │
    │                              └──────────────────┬──────────────┘
    │                                                 │
 ②  │◄── +FULLRESYNC <replid> <offset> ───────────────│
    │                                                 │
    │                                ┌────────────────┴────────────────┐
    │                                │ ③ fork() 子进程 BGSAVE          │
    │                                │   生成 RDB;同时主线程收到的     │
    │                                │   每条写命令推到「复制缓冲区」    │
    │                                │   (client output buffer)         │
    │                                │                                  │
    │                                │   buffer 内容(举例):           │
    │                                │     SET user:99 Eve              │
    │                                │     INCR cnt                     │
    │                                │     LPUSH log "..."              │
    │                                └────────────────┬────────────────┘
    │                                                 │
    │                                ┌────────────────┴────────────────┐
    │                                │ BGSAVE 完成 → RDB 文件落盘 / 内存│
    │                                └────────────────┬────────────────┘
    │                                                 │
 ④  │◄── $<rdb_size>\r\n<rdb_bytes> ──────────────────│  传输 RDB
    │                                                 │
    │ ┌─ 收到 RDB ─────────────────────────────────┐  │
    │ │ ⑤a 清空自己原有的所有数据                  │  │
    │ │ ⑤b 加载 RDB,恢复成主节点 BGSAVE 那一刻的快照│  │
    │ └────────────────────────────────────────────┘  │
    │                                                 │
    │◄── 重放复制缓冲区里积攒的命令 ────────────────────│  追赶差距
    │     SET user:99 Eve                             │
    │     INCR cnt                                    │
    │     LPUSH log "..."                             │
    │                                                 │
 ✓  │── 数据追平 → master_link_status = up ─────────►│
    │                                                 │
    │── 进入命令传播阶段(持续接收写命令)──────────────►│

🔑 三个易错点

  1. 从节点收到 RDB 后先清空自己,再加载——这一步从节点会进入 loading 状态,默认仍可读但读到旧值;可设 replica-serve-stale-data no 让它直接拒绝读。
  2. 复制缓冲区是 per-replica 的(每个从节点单独一块),全量同步专用,不要和复制积压缓冲区(backlog)搞混——后者是全局共享的环形 buffer,给增量同步用。
  3. RDB 是 BGSAVE 那一刻的快照,不包含 BGSAVE 期间产生的写命令,所以才需要复制缓冲区在 RDB 传完后补齐这段差距。

8.3.4 关键细节:fork 阻塞、复制缓冲区、无盘复制

① BGSAVE 与 fork 开销

主节点会 fork 一个子进程生成 RDB。fork 本身在 Linux 上是 O(1)(COW,写时复制),但如果主节点内存大(几十 GB),fork 期间的页表复制 + 写时复制带来的短暂阻塞会让主线程卡几百毫秒。生产配套:

  • 关闭透明大页:echo never > /sys/kernel/mm/transparent_hugepage/enabled
  • 避免在业务高峰期主动触发全量。
  • 监控 latest_fork_usec(INFO stats),fork 时长持续大于 1 秒就要警惕了。

② 复制缓冲区配置

配置项:client-output-buffer-limit replica 256mb 64mb 60

含义:
  - 硬限制 256mb:超过立即断开
  - 软限制 64mb 持续 60 秒:连续 60 秒超 64mb 也断开

如果 RDB 传输慢 + 主节点写入猛,buffer 撑爆 → 从节点被踢
→ 从节点重连 → 又一轮全量 → 还撑爆 → 复制风暴!

⚠️ 生产踩坑复制风暴的根源往往就是这块缓冲区配小了。大流量场景下建议放宽到 512mb 128mb 60 甚至更大,并配套调大 repl-backlog-size

③ 无盘复制(Diskless Replication)

Redis 2.8.18 起支持 repl-diskless-sync yes。主节点不再把 RDB 落盘再传,而是 fork 子进程后直接通过 socket 把 RDB 流式发给从节点。适合磁盘慢、网络快的场景(云上 SSD 一般,但万兆内网)。

配置项:
  repl-diskless-sync yes
  repl-diskless-sync-delay 5    # 等几秒看有没有更多 replica 一起来,攒一批

8.4 增量同步

8.4.1 PSYNC 命令的三种返回

Redis 2.8 引入 PSYNC(之前只有 SYNC,任何主从断连重连都要全量,代价巨大)。PSYNC 的核心思想:记下「我同步到哪儿了」,断线重连后只补差量

PSYNC 命令格式:
  PSYNC <replid> <offset>

主节点的三种返回:
  ① +FULLRESYNC <replid> <offset>   → 全量同步(建立新关系)
  ② +CONTINUE                        → 增量同步(接着之前的进度)
  ③ -ERR / 其他错误                   → 协议失败、重连

主节点的判断逻辑(伪代码):

python
def handle_psync(slave_replid, slave_offset):
    if slave_replid == "?" or slave_offset == -1:
        return "FULLRESYNC", master.replid, master.offset

    if slave_replid != master.replid and slave_replid != master.replid2:
        return "FULLRESYNC", master.replid, master.offset      # replid 完全对不上

    backlog_first = max(0, master.offset - master.backlog_size)
    if slave_offset < backlog_first:
        return "FULLRESYNC", master.replid, master.offset      # offset 已被覆盖

    return "CONTINUE"                                          # 走增量

8.4.2 关键数据结构 ① replication backlog

这是增量同步的心脏——一块环形(circular)字节缓冲区全局共享一份(不是 per-replica 的)。每条主节点执行的写命令在发给从节点的同时,也会写入这个环形 buffer。

   主节点的 master_repl_offset 在持续增长(从 0 开始累加,永不回退)

  ┌──────────────────────────────────────────────────────────┐
  │  环形缓冲区 (默认 1MB),写满后从头覆盖                     │
  │                                                            │
  │  [..命令D..][..命令E..][..命令F..][..命令A..][..命令B..]   │
  │   ↑                                  ↑                    │
  │   最旧的命令(即将被覆盖)            写指针                 │
  │   offset=12000                       offset=13024         │
  └──────────────────────────────────────────────────────────┘

  从节点重连时报 PSYNC <replid> <offset>:
    if backlog_first_byte_offset ≤ offset ≤ master_repl_offset:
        → +CONTINUE,把 offset 之后的命令一次性发过去
    else:
        → +FULLRESYNC,缺的命令已经被环形覆盖了,没法补

怎么算大小?

backlog_size 至少要 ≥ master_writes_per_second × max_disconnect_seconds

经验公式:
  写入速率 1MB/s,最大可容忍断网 60 秒 → backlog ≥ 60MB
  写入速率 100KB/s,断网 5 分钟        → backlog ≥ 30MB

配置项:repl-backlog-size 64mb默认 1mb 太小,生产几乎必调)。

8.4.3 关键数据结构 ② replid(runid)

每个 Redis 节点启动时都会生成一个 40 位十六进制的随机字符串作为身份证:

  • 2.8 ~ 4.0:叫 runid,每次重启都变 → 主节点重启后所有从必须全量。
  • 4.0 +:改叫 replid,并引入 PSYNC2 协议——主节点关闭前会把 replidmaster_repl_offset 持久化到 RDB 头部,重启后恢复,从节点不必全量;同时维护 replid2(上一次的 replid)支持哨兵切换不全量。
INFO server 里看:
  run_id:7f0d9c5b8a2e1f4d6c8a3b5e7f9d1c3a5e7f9d1c

INFO replication 里看:
  master_replid:7f0d9c5b8a2e1f4d6c8a3b5e7f9d1c3a5e7f9d1c
  master_replid2:0000000000000000000000000000000000000000   ← 历史 replid
  second_repl_offset:-1                                       ← 历史 offset 断点

8.4.4 关键数据结构 ③ 复制偏移量 offset

主从都会维护一个 64 位整数 offset:

  • master_repl_offset:主节点累计发出的字节数。
  • slave_repl_offset:从节点累计接收并执行的字节数。
正常情况:master_repl_offset == slave_repl_offset  → 完全同步
延迟情况:master_repl_offset >  slave_repl_offset  → 差距即「滞后字节数」

INFO replication 里的关键字段:
  master_repl_offset:13024
  slave0:ip=192.168.1.10,port=6379,state=online,offset=12950,lag=1
                                                  ↑      ↑
                                            从同步到的位置  滞后秒数

🎯 易错点:offset 是字节数,不是命令条数——RESP 协议下不同命令大小不同,按字节计才能精确定位 backlog。

8.4.5 三种数据结构如何配合判定增量 / 全量

       从节点 PSYNC <replid_from_slave> <offset_from_slave>


   ┌──────────────────────────────────────────────────────┐
   │ 主节点判断:                                          │
   │                                                       │
   │ 1. replid_from_slave == master.replid                │
   │    → 检查 offset 是否在积压缓冲区范围                  │
   │      → 在:CONTINUE(增量)                           │
   │      → 不在:FULLRESYNC(全量)                       │
   │                                                       │
   │ 2. replid_from_slave == master.replid2 (PSYNC2)      │
   │    → 老主重启 / 主从切换的场景                         │
   │      → 同上判断 offset                                │
   │                                                       │
   │ 3. 都对不上                                           │
   │    → FULLRESYNC(全量)                              │
   └──────────────────────────────────────────────────────┘
场景replid 比较offset 比较结果
首次接入?-1FULLRESYNC
同主断线短时重连一致在 backlogCONTINUE
同主断线时间太久一致已被覆盖FULLRESYNC
哨兵切换后等于 replid2在 backlogCONTINUE(PSYNC2)
跟过别的主都对不上FULLRESYNC

8.5 命令传播

8.5.1 异步推送,不等 ACK

进入命令传播阶段后,主节点会做两件事:

主节点每执行一条写命令:
  ① 在自己内存里执行(修改数据 + 更新 master_repl_offset)
  ② 把命令字节追加到:
     - replication backlog(环形 buffer)
     - 每个 replica 的 client output buffer
  ③ 通过 socket 异步把命令推给所有 replica

⚠️ 注意:主节点不会等 replica 回 ACK,写完自己内存就返回客户端。

这就是「异步复制」的本质——天然弱一致

8.5.2 心跳:REPLCONF ACK

从节点每秒主动给主节点发:

REPLCONF ACK <slave_repl_offset>

作用:
  - 上报「我已经同步到 offset 多少了」
  - 同时充当心跳:30 秒内主节点没收到 ACK,会断开这个 replica
  - 主节点据此更新 INFO 里的 lag 字段

8.5.3 异步复制的代价

T0: client → master:  SET user:1 "Alice"     ← master 写完立即返回 OK
T0+1ms: master 把命令塞进 replica 的 output buffer
T0+10ms: replica 收到命令并执行
T0+10ms+ε: client 跑去 replica 读 user:1 → 此时才能读到 "Alice"
一致性模型是否满足
强一致(写完所有副本立即可见)
顺序一致(同一客户端读到的版本单调)❌(如果客户端连不同 replica)
最终一致(足够时间后所有副本一致)

💡 想要强一致?Redis 提供了 WAIT N timeout 命令——客户端写完后调 WAIT 1 1000 表示最多等 1 秒,至少 1 个从同步完才返回。但 WAIT 有个隐藏坑:它不保证后续不丢——主节点挂了就算 WAIT 已返回也可能丢,因此 Redis 的复制只是「尽力」一致,不是 quorum 写

8.5.4 主从异步复制 vs 强一致方案对比

方案一致性性能适用场景
默认异步复制最终一致最高99% 业务
WAIT N timeout最多等 N 个从确认写延迟变高关键写入场景
min-replicas-to-write不满足条件就拒写牺牲可用性金融级写入安全
关闭读写分离读也走主,避免延迟读容量受限于主强一致核心数据

8.6 主从延迟问题

8.6.1 延迟的三大来源

┌──────────────────────────────────────────────────────────┐
│  主从延迟 = 网络延迟 + 命令处理延迟 + 持久化阻塞延迟         │
├──────────────────────────────────────────────────────────┤
│                                                            │
│  ① 网络延迟                                                │
│     - 跨机房:可能十几 ms 起步                              │
│     - 同机房:通常 < 1ms                                   │
│     - TCP 拥塞、丢包、网卡打满                              │
│                                                            │
│  ② 命令处理延迟                                            │
│     - 从节点也是单线程,慢命令会让回放卡顿                   │
│     - 大 Key 的删除、复杂的 Lua 脚本                        │
│     - 大 Hash / List 的同步本身就慢                         │
│                                                            │
│  ③ 持久化阻塞                                               │
│     - 从节点开了 AOF always 刷盘 → 每条命令都要 fsync      │
│     - 从节点正在做 BGSAVE / BGREWRITEAOF → fork 阻塞        │
│                                                            │
└──────────────────────────────────────────────────────────┘

8.6.2 监控指标:master_repl_offset - slave_repl_offset

排查主从延迟的第一指标

delta_bytes = master_repl_offset - slave_repl_offset

  delta = 0       → 完全同步
  delta < 1KB     → 微小波动,正常
  delta > 1MB     → 明显积压,要排查
  delta 持续上涨   → 复制即将断(backlog 撑爆 → FULLRESYNC)

辅助指标:

指标来源健康阈值
slave0.lag主节点 INFO replication≤ 1 秒
master_last_io_seconds_ago从节点 INFO replication≤ 1 秒
master_link_status从节点 INFO replicationup
repl_backlog_histlen / repl_backlog_size主节点 INFO replication利用率不应长期 100%
latest_fork_usecINFO stats< 1 秒

8.6.3 优化思路速查表

问题症状优化方向
同机房网络抖动排查丢包、换交换机、用 repl-disable-tcp-nodelay no 让数据立即发
跨机房延迟高拓扑改成级联(M → 异地 S1 → 本地 S2,S3)
大 Key 同步慢拆 Key、用 SCAN 渐进、避免 KEYS 命令
从节点 fsync 太频从节点关 appendfsync always,改 everysec
backlog 撑爆调大 repl-backlog-size 至 64mb 起步
fork 阻塞长关闭透明大页、限制实例内存不超过 10GB

8.6.4 min-replicas-to-write 兜底

如果对「写入了多少份副本」有要求(比如金融场景,至少要 1 个从同步成功才认写入有效):

min-replicas-to-write 1     # 至少要有 1 个从在线
min-replicas-max-lag  10    # 且其延迟 ≤ 10 秒

效果:当在线 + 满足延迟条件的 replica 数 < 1,
     master 直接拒绝所有写命令(返回 NOREPLICAS 错误)
     用「不可写」换「写入安全性」

⚠️ 副作用:所有 replica 集体延迟过高时,主节点会拒绝写入,业务侧要做好处理。这是用一致性换可用性的典型选择。


8.7 读写分离的「坑」

读写分离听起来美好,实际上有几个经典翻车现场。一句话总结:任何主从延迟敏感的业务,都要慎用读写分离

8.7.1 坑 1:刚写完立刻读,读到旧值

T0:    client → master:  SET user:1 "Alice"     ✅ 返回 OK
T0+ε:  client → replica: GET user:1
                          → 返回 (nil) 或上一次的值  ❌

业务表现

  • 用户改了昵称,刷新页面还是旧的;
  • 下单后跳转到列表页发现订单不在;
  • 新注册立即登录失败(注册写主,登录校验读从,从还没收到这条用户记录)。

应对方案

方案说明适用场景
写后短时间路由到主同一个会话/用户在 N 秒内强制读主简单粗暴,最常用
强一致读走主关键查询不走 replica配合业务标记
WAIT N timeoutRedis 命令,等待 N 个从同步完才返回强一致要求
业务幂等接受短暂不一致,下次刷新自动修正可接受最终一致

代码示例(写后强制读主):

python
def register(username):
    master.set(f"user:{username}", json.dumps(user))
    # 标记:这个用户接下来 3 秒内的读都打主节点
    session_router.pin_to_master(username, ttl=3)

def login(username):
    if session_router.should_read_master(username):
        return master.get(f"user:{username}")
    return replica.get(f"user:{username}")

8.7.2 坑 2:从节点宕机/重启,所有压力压回主节点

读写分离一旦从节点出问题:

  • 客户端没及时剔除故障从 → 部分请求报错;
  • 全部 fail-over 回主 → 主可能瞬间扛不住读流量。

做读写分离前必须有兜底

  • 客户端做健康检查 + 自动剔除故障从。
  • 用 LB(如 Twemproxy / Codis / 云厂商 Proxy)/ Sentinel 做透明切换。
  • 主节点容量按「主自身流量 + 1 ~ 2 个从挂掉的流量」预留。

8.7.3 坑 3:从节点 loading 状态返回 stale 数据

从节点正在加载 RDB(master_sync_in_progress=1)时,默认会返回旧数据replica-serve-stale-data yes)。如果业务无法容忍,可设:

replica-serve-stale-data no   # loading 期间所有命令返回 LOADING 错误

但这意味着每次全量同步期间该 replica 不可读——要在客户端做健康检查(探测 master_sync_in_progress)。

8.7.4 坑 4:过期 Key 的同步歧义

  • 从节点不主动过期 Key,要等主节点发 DEL 命令同步过来。
  • 3.2 之前从节点甚至会返回已过期的 Key。
  • 4.0+ 已修复,但仍有极小窗口期。

应对:升级新版本;关键场景显式 TTL 检查。

8.7.5 坑 5:大事务 / Lua 脚本期间从节点延迟暴增

  • 主节点单条命令耗时长(比如 EVAL 跑 500ms 的 Lua),命令传到从后回放也耗时长。
  • 应对:拆小事务、限制 Lua 执行时长、监控 slowlog

8.7.6 用 WAIT 命令实现「至少 N 个从确认」

bash
# 写入后等待至少 1 个从同步完,最多等 100ms
SET user:1 Alice
WAIT 1 100
# 返回值:实际同步到的 replica 数(如果是 0,说明 100ms 内没等到)

注意 3 个限制

  1. WAIT 不保证后续不丢——主节点挂了仍可能丢。
  2. WAIT 会拖慢写延迟——慎用在高 QPS 路径。
  3. WAIT 不保证「具体哪个 replica」同步完了——只保证「至少 N 个」。

🎯 统一原则

  • 能接受最终一致 → 做读写分离 + 兜底(重试 / 强制读主 / 业务幂等)。
  • 强一致核心业务 → 直接所有读写都走主(牺牲读扩展换简单)。
  • 金融级写入安全min-replicas-to-write + WAIT。

8.8 实战:用 INFO replication 看复制状态

8.8.1 主节点视角

127.0.0.1:6379> INFO replication
# Replication
role:master
connected_slaves:2
slave0:ip=192.168.1.11,port=6379,state=online,offset=15234,lag=0
slave1:ip=192.168.1.12,port=6379,state=online,offset=15180,lag=1
master_failover_state:no-failover
master_replid:7f0d9c5b8a2e1f4d6c8a3b5e7f9d1c3a5e7f9d1c
master_replid2:0000000000000000000000000000000000000000
master_repl_offset:15234
second_repl_offset:-1
repl_backlog_active:1
repl_backlog_size:1048576
repl_backlog_first_byte_offset:14186
repl_backlog_histlen:1048

逐字段解读:

字段含义
role:master当前节点角色
connected_slaves:2在线从节点数
slave0/slave1各从节点的 ip/port/状态/offset/延迟秒数
master_replid当前 replication ID
master_replid2上一次的 replication ID(PSYNC2 用)
master_repl_offset主节点累计发出字节数
repl_backlog_size积压缓冲区大小(默认 1MB,生产必调
repl_backlog_first_byte_offsetbacklog 中最旧字节的 offset
repl_backlog_histlenbacklog 中实际有多少字节有效

判断从节点能不能增量同步:从节点的 offset 必须 ≥ repl_backlog_first_byte_offset,否则只能全量。

8.8.2 从节点视角

127.0.0.1:6380> INFO replication
# Replication
role:slave
master_host:192.168.1.10
master_port:6379
master_link_status:up
master_last_io_seconds_ago:1
master_sync_in_progress:0
slave_read_repl_offset:15234
slave_repl_offset:15234
master_link_down_since_seconds:-1
slave_priority:100
slave_read_only:1
字段含义
role:slave当前是从节点
master_link_status:up与主连接正常(down 即断连)
master_last_io_seconds_ago距离上次收到主的数据多久了
master_sync_in_progress是否正在做全量同步(0/1)
slave_repl_offset已同步到的 offset
slave_priority哨兵选主时的优先级(值越小越优先,0 表示永不当选)
slave_read_only从节点是否只读(建议保持 1)

8.8.3 跑一遍配套代码

实战代码见 08_replication/code/

  • 01_check_replication.py:用 INFO replication 看主从信息(role、connected_slaves、master_repl_offset 等)。单机环境会自动 mock 一份典型主从输出,方便理解每个字段。
  • 02_replication_lag.py:测量「主写入 → 从读到」的延迟(写一个时间戳 key,从节点轮询读,计算时间差)。单机环境会用「镜像 key + 后台同步线程」模拟主从异步复制
  • 03_wait_command.py:演示 WAIT N timeout 命令——客户端写入后调 WAIT 等待至少 N 个从确认。

浏览器演示见 08_replication/demo.html

  • 全量同步流程动画:可视化 5 步流程,按「下一步」单步看,或「自动播放」一气呵成。
  • 环形 backlog 可视化:圆环 16 个槽,主节点写入往里塞,模拟「从节点离线 N 秒」看是否还能 CONTINUE。
  • 主从拓扑模拟器:拖动添加从节点,演示「一主多从」「级联复制」两种拓扑;点击主节点写入看命令往各从节点扩散。

8.9 本章小结

┌──────────────────────────────────────────────────────────┐
│                      本章核心要点                           │
├──────────────────────────────────────────────────────────┤
│                                                            │
│  ① 主从复制解决三件事:读扩展 + 容灾备份 + 故障转移基础      │
│                                                            │
│  ② 三种拓扑:一主一从、一主多从、级联复制                   │
│     - 默认优先一主多从                                     │
│                                                            │
│  ③ 全量同步五步:PSYNC → FULLRESYNC → BGSAVE 出 RDB →      │
│     传 RDB → 加载 + 重放复制缓冲区                         │
│     - 复制缓冲区(per-replica)撑不住 → 复制风暴            │
│                                                            │
│  ④ 增量同步靠三件套:                                       │
│     - replication backlog(环形 buffer,全局共享)         │
│     - replid(节点身份证)                                 │
│     - offset(累计字节数)                                 │
│     PSYNC 三种返回:FULLRESYNC / CONTINUE / -ERR           │
│                                                            │
│  ⑤ 命令传播是异步的 → 天然弱一致 → 主从延迟无可避免         │
│     - 监控核心:master_repl_offset - slave_repl_offset    │
│                                                            │
│  ⑥ 读写分离 5 个坑:写后立即读、从挂打主、loading stale、    │
│     过期 key、大事务延迟。强一致就别做读写分离。            │
│                                                            │
│  ⑦ 强一致工具:WAIT N timeout / min-replicas-to-write      │
│     都是「以可用性换一致性」                                │
│                                                            │
│  ⑧ 排障必看:INFO replication 的 master_repl_offset、       │
│     repl_backlog_*、slave_repl_offset、master_link_status  │
│                                                            │
└──────────────────────────────────────────────────────────┘

8.10 面试高频题

Q1:Redis 主从复制的全量同步流程是怎样的?

考察点:复制流程整体把控、五步流程的细节。

标准答案

完整的全量同步分为五步:

  1. 从节点发起 PSYNC:从节点连上主节点后发送 PSYNC ? -1(首次接入时 replid 未知、offset 为 -1)。
  2. 主节点回 +FULLRESYNC:主节点判断这是新连接,回复 +FULLRESYNC <master_replid> <master_offset>,告诉从节点「我们要走全量同步」以及当前主节点的身份和位置。
  3. 主节点 BGSAVE 生成 RDB:fork 子进程执行 BGSAVE,把当前内存数据快照写入 RDB 文件。此期间主节点新接收的写命令会被塞入这个 replica 专属的「复制缓冲区」(client output buffer),等 RDB 传完再补发。
  4. 传输 RDB:BGSAVE 完成后,主节点把 RDB 文件通过 socket 流式发给从节点($<size>\r\n<bytes> 协议)。
  5. 从节点加载 + 增量补齐:从节点收到 RDB 后先清空自己原有数据,再加载 RDB,最后重放复制缓冲区里积攒的命令追上主节点。完成后 master_link_status 变为 up,进入命令传播阶段。

加分项

  • 主节点 fork 子进程会有短暂阻塞,建议关闭透明大页。
  • 2.8.18+ 支持 repl-diskless-sync,RDB 不落盘直接通过 socket 流式传输。
  • 复制缓冲区(per-replica)和复制积压缓冲区(全局环形)是两个不同的东西——前者全量同步用,后者增量同步用,面试别答混。

易错点

  • 把「从节点收到 RDB 后立即接命令」说错——其实从节点要先清空自己原有数据再加载 RDB,这一步会让从节点短暂不可用。
  • 把复制缓冲区说成「环形 buffer」——它是普通的 client output buffer,撑爆了会断连。

Q2:增量同步是怎么做到「断点续传」的?

考察点:PSYNC 协议、replid + offset 的协作。

标准答案

增量同步的核心是 PSYNC <replid> <offset> 命令 + 主节点维护的复制积压缓冲区(replication backlog)

工作机制

  1. 主节点每执行一个写命令,两个动作并发
    • 把命令字节追加到 backlog(环形 buffer);
    • 异步推给所有在线 replica。
  2. 从节点维护自己已同步到的 slave_repl_offset
  3. 当从节点掉线重连时,发送 PSYNC <自己保存的 replid> <slave_repl_offset>
  4. 主节点判断:
    • 如果 replid 和 offset 都对得上,并且 offset 仍在 backlog 范围(backlog_first_byte_offset ≤ slave_offset ≤ master_offset),就回 +CONTINUE,把 backlog 里 offset 之后的命令一次性发过去;
    • 否则回 +FULLRESYNC,走全量。

类比

backlog 就像主节点身边的一本「最近写过的命令日志本」,从节点掉线回来时拿着「我上次看到第 12950 行」的书签来对——只要日志本里还有这一行就能接着读,否则只能重抄一份。

PSYNC 的三种返回

返回含义触发条件
+FULLRESYNC <replid> <offset>全量同步首次接入 / replid 对不上 / offset 已被覆盖
+CONTINUE增量同步replid 一致 + offset 在 backlog 范围
-ERR ...协议失败不支持 PSYNC2、版本不兼容等

加分项:4.0+ 引入 PSYNC2,主节点维护 replid2 + second_repl_offset,让主节点重启 / 哨兵切换后从节点也能走增量。

易错点:把 offset 当成命令条数——其实它是字节数


Q3:复制积压缓冲区(replication backlog)的作用?设置多大合适?

考察点:增量同步底层机制、生产参数调优经验。

标准答案

作用

  • 是主节点维护的一块固定大小、环形、全局共享的字节缓冲区。
  • 每条主节点执行的写命令在发给从节点的同时,也写入这个 buffer。
  • 当从节点短暂掉线(网络抖动、重启),重连后报上自己的 offset只要这个 offset 还在 backlog 范围内,主节点就能从 backlog 里捞出缺失的命令做增量补发,避免昂贵的全量同步。

默认大小 1MB,生产环境几乎必调。经验公式:

backlog_size ≥ 平均写入速率 (bytes/s) × 最大可容忍断连秒数

例如:
  写入 1MB/s + 容忍断连 60s → 至少 60MB
  写入 100KB/s + 容忍断连 5 分钟 → 至少 30MB

配置项:repl-backlog-size 64mb(建议起步)。

设置太小的后果:从节点稍微断线长一点(甚至只是网络抖动几秒),重连时 offset 已被覆盖 → 全量同步 → BGSAVE + RDB 传输 → 复制缓冲区可能撑爆 → 又触发全量 → 复制风暴

设置太大的后果

  • 占用主节点内存(10GB 的 backlog 就吃 10GB 内存)。
  • 实际能用上的部分不多——绝大多数业务断连不会超过几分钟。

加分项

  • 当主节点没有任何在线从节点超过 repl-backlog-ttl(默认 3600 秒),backlog 会被释放。
  • 区分两个 buffer:复制缓冲区(client output buffer)是 per-replica 的,全量同步期间用复制积压缓冲区是全局共享的环形 buffer,给增量同步用。这两个经常被搞混。

易错点:以为越大越好,结果占用大量内存又用不上。


Q4:主从异步复制有什么问题?怎么实现强一致?

考察点:对一致性模型的理解、WAITmin-replicas-to-write 的使用。

标准答案

异步复制的问题

主节点执行完写命令后不等任何从节点 ACK 就返回客户端,因此:

  1. 数据可能丢失:主节点写入后立即宕机,命令还没推到从节点 → 切换主后这条数据丢了。
  2. 读到旧值:客户端写主后立刻读从,从节点还没收到 → 读到旧值(即「写后读不一致」)。
  3. 顺序歧义:客户端连不同 replica,可能看到「时光倒流」般的数据视图。

Redis 提供的 3 种「向强一致靠拢」方案

方案命令 / 配置一致性强度代价
WAIT N timeout客户端命令写完等 N 个从确认才返回写延迟变高,仍可能因主宕机丢数据
min-replicas-to-writeserver 配置在线 replica 数不够时拒绝写可用性下降(NOREPLICAS 错误)
关闭读写分离应用层读也走主,避免主从延迟读容量受限于主节点

WAIT 用法示例

SET user:1 Alice
WAIT 1 100        # 最多等 100ms,至少 1 个从确认才返回
                  # 返回值:实际同步到的 replica 数

重要提醒

  • WAIT 不是 quorum 写——它只保证「至少 N 个从确认了当前 offset」,不保证「即使主挂了这个数据也存活」。
  • Redis 从根本上不是强一致系统——真要强一致请用 etcd / ZooKeeper / Spanner。

加分项:能提到 Redis 7.4 的 client-side caching invalidation、Redis Cluster 的故障转移模型、Raft 协议如何用 quorum 写实现强一致。

易错点

  • 以为「WAIT 1 = 至少 1 个从一定收到 → 数据就不丢了」——其实主节点宕机仍可能丢,因为 Redis 没有 fsync-then-ack 协议。
  • 把 WAIT 用在高 QPS 写路径,导致整体性能崩盘。

Q5:读写分离会有什么坑?

考察点:实际工程经验、对最终一致性的理解。

标准答案

五大坑

  1. 写后立即读,读到旧值

    • 原因:主从异步复制有延迟,从节点还没收到这条命令。
    • 应对:① 同会话 N 秒内强制读主;② 用 WAIT N timeout 等待复制;③ 业务侧幂等设计。
  2. 从节点宕机/重启,流量打回主节点

    • 原因:客户端没及时剔除故障从,或全部回退到主时主撑不住。
    • 应对:客户端做健康检查 + 自动剔除;用 LB / Sentinel 做透明切换;主节点容量按「自身流量 + 几个从挂掉的流量」预留。
  3. 从节点 loading 状态返回 stale 数据

    • 原因:从节点正在加载 RDB(全量同步中),默认 replica-serve-stale-data yes,会返回旧数据。
    • 应对:业务无法容忍时设 no 拒绝读;或健康检查探测 master_sync_in_progress
  4. 过期 Key 在从节点上没及时删除

    • 原因:从节点不主动过期,要等主节点发 DEL 命令同步过来;3.2 之前从节点甚至会返回已过期的 Key。
    • 应对:升级到新版本(已修复);关键场景显式 TTL 检查。
  5. 大事务 / Lua 脚本期间从节点延迟暴增

    • 原因:主节点单条命令耗时长,命令传到从后回放也耗时长。
    • 应对:拆小事务、限制 Lua 执行时长、监控 slowlog

统一应对原则

  • 强一致业务直接所有读写都走主(牺牲读扩展换简单)。
  • 能接受最终一致的业务再做读写分离,并设计兜底(重试 / 强制读主 / 业务幂等)。
  • 金融级写入安全min-replicas-to-write + WAIT

加分项:能提到 Redis 6.0+ 的 CLIENT NO-EVICT / CLIENT REPLY 等更细粒度的客户端控制;提到「写后读」一致性方案在 MySQL 主从场景也有相同的解法(GTID 校验等),思想是相通的。

易错点:把读写分离当万能方案;忽视客户端层的健康检查;忘记从节点也是单线程。


Q6:主从延迟怎么排查和优化?

考察点:实际生产问题排障思路。

标准答案

第一步:定位延迟程度——用 INFO replication 看:

master 端:
  master_repl_offset:15234
  slave0:...,offset=12950,lag=1

            最直观的滞后秒数(理论 0~1 秒为正常)

  delta_bytes = master_repl_offset - slave_offset
              = 15234 - 12950 = 2284  bytes 还在传

slave 端:
  master_link_status:up
  master_last_io_seconds_ago:1

第二步:定位延迟来源——对照三大原因排查:

来源排查命令解决方向
网络延迟redis-cli --latency -h master_ip同机房部署、万兆内网、repl-disable-tcp-nodelay no
慢命令SLOWLOG GET 10拆 KEYS / 大 Key、限制 Lua 时长
持久化阻塞INFO persistenceaof_pending_bio_fsyncrdb_bgsave_in_progress从节点改 appendfsync everysec、错峰 BGSAVE
backlog 撑爆INFO replicationrepl_backlog_histlen 利用率调大 repl-backlog-size
fork 阻塞INFO statslatest_fork_usec关闭透明大页、限制实例内存

第三步:优化思路速查

  • 网络层:主从同机房同 VPC 部署;万兆内网;避免跨地域复制(要跨地域用专门的 CRDT/Cluster Bus)。
  • 慢命令:禁用 KEYS、SMEMBERS(大集合),用 SCAN 系列;拆大 Key;Lua 脚本控制时长。
  • 持久化:从节点不开 AOF always;错峰执行 BGSAVE;冷备从单独部署不挂业务流量。
  • 监控:盯 INFO replicationlag 字段和 master_last_io_seconds_ago,超阈值告警。
  • 业务层:写后强制读主、用 WAIT 命令、接受最终一致。
  • 配置层min-replicas-to-write + min-replicas-max-lag 强制保证至少 N 个从延迟在 M 秒内才允许写。

加分项

  • 能提到 repl-disable-tcp-nodelay no(默认)让数据立即发送,开 yes 则 40ms 合并发送以省带宽——跨机房时反而要保持 no
  • 能提到 Redis 的 latency monitorCONFIG SET latency-monitor-threshold 100 + LATENCY HISTORY event 是 Redis 自带的延迟事件审计工具。

易错点

  • 把「从节点延迟」和「主从断连」混为一谈。延迟是连着的但慢;断连是 master_link_status:down
  • 一上来就调大 backlog——backlog 大只是减少全量概率,对实时延迟无效

📌 下一章预告:第 9 章我们看 哨兵(Sentinel)——主从复制只是「数据多副本」的基础,但主节点宕机后谁来发现、谁来选新主、客户端怎么知道? 哨兵就是来解决这套自动化故障转移流程的。它的「主观下线 / 客观下线」「Raft-like 选举」「客户端订阅切换通知」每一个都是经典面试考点。

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

python
"""
Ch8 配套代码 1 / 3 —— 用 INFO replication 看主从信息

⚠️ 本脚本需要主从环境才能看到完整的 connected_slaves / slave0 等字段。
   若教程环境只有单机 Redis,脚本会自动 fallback:先打真实输出,
   再额外渲染一份「典型主从环境」的 mock 输出,对照学习每个字段。

主要演示:
  1. 主节点视角:role / connected_slaves / master_repl_offset
                 / master_replid / repl_backlog_*
  2. 从节点视角:role / master_link_status / slave_repl_offset
                 / master_last_io_seconds_ago / master_sync_in_progress
  3. 派生指标:积压缓冲区利用率、各从节点滞后字节数等
"""

import redis

r = redis.Redis(host="127.0.0.1", port=6379, decode_responses=True)


FIELD_DOC = {
    "role": "节点角色:master / slave",
    "connected_slaves": "在线从节点数量",
    "master_replid": "复制 ID(40 位 hex),节点身份证。从节点用它判断「跟的是不是同一个主」",
    "master_replid2": "上一次的复制 ID(PSYNC2),用于支持主重启 / 哨兵切换不全量",
    "master_repl_offset": "主节点累计发出字节数,单调递增,永不回退",
    "second_repl_offset": "切换 replid 时的 offset 断点,配合 replid2",
    "repl_backlog_active": "积压缓冲区是否启用(1 / 0)",
    "repl_backlog_size": "积压缓冲区大小(字节)。**生产必调**,默认 1MB 太小",
    "repl_backlog_first_byte_offset": "backlog 中最旧字节的 offset,从节点 offset < 它就只能全量",
    "repl_backlog_histlen": "backlog 中实际有效字节数(未写满时 < repl_backlog_size)",
    "master_host": "主节点 IP(slave 视角)",
    "master_port": "主节点端口(slave 视角)",
    "master_link_status": "与主连接状态:up / down",
    "master_last_io_seconds_ago": "距上次收到主的数据多少秒(长时间不动可能要排查心跳)",
    "master_sync_in_progress": "是否正在做全量同步(1 = loading 中)",
    "slave_repl_offset": "从节点已同步到的 offset,与 master_repl_offset 之差 = 滞后字节数",
    "slave_priority": "哨兵选主优先级(值越小越优先;0 = 永不当选)",
    "slave_read_only": "从节点是否只读(强烈建议保持 1)",
    "master_link_down_since_seconds": "与主断连了多少秒(-1 = 当前正常)",
}


def section(title: str) -> None:
    print("\n" + "=" * 70)
    print(title)
    print("=" * 70)


def show_real_info() -> dict:
    section("Demo 1: 当前真实环境的 INFO replication 输出")
    info = r.info("replication")
    for k, v in info.items():
        print(f"  {k:36} = {v}")
    print(f"\n  → 当前节点 role = {info.get('role')}, "
          f"在线 replica = {info.get('connected_slaves', 0)}")
    return info


def explain_each_field(info: dict) -> None:
    section("Demo 2: 字段逐项解读")
    for k, v in info.items():
        doc = FIELD_DOC.get(k, "(暂未收录解释)")
        print(f"  📌 {k} = {v}")
        print(f"     → {doc}\n")


def derive_metrics(info: dict) -> None:
    section("Demo 3: 派生指标 + 健康度评估")
    if info.get("role") == "master":
        size = info.get("repl_backlog_size", 0)
        used = info.get("repl_backlog_histlen", 0)
        print(f"  积压缓冲区: {size:,} 字节 ({size/1024/1024:.2f} MB)")
        print(f"  已使用    : {used:,} 字节 (利用率 "
              f"{(used/size*100 if size else 0):.1f}%)")
        if size <= 1024 * 1024:
            print("  ⚠️  backlog 仅 1MB(默认值),生产建议 ≥ 64MB")
            print("     配置示例: repl-backlog-size 64mb")
        if info.get("connected_slaves", 0) == 0:
            print("  ℹ️  当前没有在线从节点(单机 demo 环境正常)")


def show_mock_master() -> None:
    section("Demo 4: 典型【主节点】INFO replication 的样子(mock)")
    mock = """\
# Replication
role:master
connected_slaves:2
slave0:ip=192.168.1.11,port=6379,state=online,offset=15234,lag=0
slave1:ip=192.168.1.12,port=6379,state=online,offset=15180,lag=1
master_failover_state:no-failover
master_replid:7f0d9c5b8a2e1f4d6c8a3b5e7f9d1c3a5e7f9d1c
master_replid2:0000000000000000000000000000000000000000
master_repl_offset:15234
second_repl_offset:-1
repl_backlog_active:1
repl_backlog_size:1048576
repl_backlog_first_byte_offset:14186
repl_backlog_histlen:1048
"""
    print(mock)
    master_offset = 15234
    print("  📊 由 mock 输出推导:")
    print(f"     - slave0 滞后 = {master_offset - 15234} bytes (完全同步)")
    print(f"     - slave1 滞后 = {master_offset - 15180} bytes (54 字节落后)")
    print(f"     - backlog 有效区间 = [14186, {master_offset}]")
    print(f"     - 任何 slave_offset < 14186 → FULLRESYNC;否则 CONTINUE")


def show_mock_slave() -> None:
    section("Demo 5: 典型【从节点】INFO replication 的样子(mock)")
    mock = """\
# Replication
role:slave
master_host:192.168.1.10
master_port:6379
master_link_status:up
master_last_io_seconds_ago:1
master_sync_in_progress:0
slave_read_repl_offset:15234
slave_repl_offset:15234
master_link_down_since_seconds:-1
slave_priority:100
slave_read_only:1
connected_slaves:0
"""
    print(mock)
    print("  📊 健康判断要点:")
    print("     - master_link_status=up + master_sync_in_progress=0 → 正常服务")
    print("     - master_last_io_seconds_ago ≤ 1s → 心跳正常")
    print("     - slave_read_only=1 → 写命令会被拒绝(只读保护)")


if __name__ == "__main__":
    try:
        info = show_real_info()
        explain_each_field(info)
        derive_metrics(info)
        if info.get("role") == "master" and info.get("connected_slaves", 0) == 0:
            print("\n  ↓ 单机环境检测到无 replica,下面给出主/从典型 mock 输出 ↓")
            show_mock_master()
            show_mock_slave()
        print("\n✅ 完成。把这个脚本拿到真实主从集群上跑,能看到更丰富的字段。")
    except redis.ConnectionError as e:
        print(f"❌ Redis 连接失败: {e}")
python
"""
Ch8 配套代码 2 / 3 —— 测量「主写入 → 从读到」的延迟

⚠️ 真实主从环境用法:
    把 MASTER_HOST/PORT 指向主,REPLICA_HOST/PORT 指向从,
    脚本会真实测量主写入到从读到的延迟(毫秒级)。

⚠️ 单机环境用法(教程默认):
    脚本检测到 MASTER == REPLICA 时,会自动用「镜像 key + 后台同步线程」
    模拟主从异步复制:
        主键   :  demo:lag:master   ← 模拟「写主」
        镜像键 :  demo:lag:replica  ← 模拟「读从」
        线程    :  每 SYNC_LAG_MS 毫秒把主键复制到镜像键
    这样能在单机上重现「写完立即读读到旧值 / 延迟分布」的现象。

测量方式:
    主:SET demo:lag:master "<纳秒时间戳>"
    从:循环 GET demo:lag:replica,直到值等于刚才写入的时间戳
    延迟 = 当前时间 - 时间戳
"""

import threading
import time
import statistics
import redis

MASTER_HOST,  MASTER_PORT  = "127.0.0.1", 6379
REPLICA_HOST, REPLICA_PORT = "127.0.0.1", 6379

KEY_MASTER  = "demo:lag:master"
KEY_REPLICA = "demo:lag:replica"

ROUNDS         = 30
WRITE_INTERVAL = 0.1            # 每 100 ms 写一次
POLL_INTERVAL  = 0.001          # 从节点轮询读,1 ms 间隔
POLL_TIMEOUT   = 2.0            # 2 秒还读不到就判超时
SYNC_LAG_MS    = 30             # 单机 mock 模式下模拟的同步延迟


master = redis.Redis(host=MASTER_HOST, port=MASTER_PORT, decode_responses=True)
replica = redis.Redis(host=REPLICA_HOST, port=REPLICA_PORT, decode_responses=True)

stop_event = threading.Event()


def section(title: str) -> None:
    print("\n" + "=" * 70)
    print(title)
    print("=" * 70)


def is_single_instance() -> bool:
    return (MASTER_HOST, MASTER_PORT) == (REPLICA_HOST, REPLICA_PORT)


def mock_replica_worker():
    """单机环境:模拟一个有 SYNC_LAG_MS 延迟的「主→从复制管道」。"""
    while not stop_event.is_set():
        v = master.get(KEY_MASTER)
        if v is not None:
            replica.set(KEY_REPLICA, v)
        time.sleep(SYNC_LAG_MS / 1000.0)


def measure_one_round() -> float:
    """一轮:写主 → 轮询读从 → 返回延迟(秒),超时返回 -1。"""
    ts = str(time.perf_counter_ns())
    write_at = time.perf_counter()
    master.set(KEY_MASTER, ts)

    deadline = write_at + POLL_TIMEOUT
    while time.perf_counter() < deadline:
        if replica.get(KEY_REPLICA) == ts:
            return time.perf_counter() - write_at
        time.sleep(POLL_INTERVAL)
    return -1.0


def run_demo():
    section("Demo: 测量主从复制延迟")

    if is_single_instance():
        print("  ⚠️  检测到单机环境(MASTER == REPLICA),")
        print(f"      启动后台线程模拟一个 ~{SYNC_LAG_MS}ms 同步延迟的「主→从管道」")
        t = threading.Thread(target=mock_replica_worker, daemon=True)
        t.start()
        # 等同步线程对齐一次
        master.set(KEY_MASTER, "init")
        time.sleep(SYNC_LAG_MS / 1000.0 * 2)
    else:
        print(f"  ✅ 真实主从模式:master={MASTER_HOST}:{MASTER_PORT}, "
              f"replica={REPLICA_HOST}:{REPLICA_PORT}")

    print(f"  共 {ROUNDS} 轮,每轮间隔 {int(WRITE_INTERVAL*1000)}ms")
    print(f"\n  {'轮':>4} | {'延迟 (ms)':>10} | 状态")
    print(f"  {'-'*4}-+-{'-'*10}-+-{'-'*20}")

    samples = []
    timeouts = 0
    for i in range(1, ROUNDS + 1):
        lag = measure_one_round()
        if lag < 0:
            print(f"  {i:>4} | {'TIMEOUT':>10} | ❌ 超过 {POLL_TIMEOUT}s 还没同步过来")
            timeouts += 1
        else:
            samples.append(lag * 1000)
            tag = "✅" if lag * 1000 < 100 else ("⚠️" if lag * 1000 < 500 else "🔥")
            print(f"  {i:>4} | {lag*1000:>10.2f} | {tag}")
        time.sleep(WRITE_INTERVAL)

    stop_event.set()

    section("结果统计")
    print(f"  总轮数: {ROUNDS}, 超时: {timeouts}, 有效样本: {len(samples)}")
    if samples:
        print(f"  min    = {min(samples):.2f} ms")
        print(f"  max    = {max(samples):.2f} ms")
        print(f"  mean   = {statistics.mean(samples):.2f} ms")
        print(f"  median = {statistics.median(samples):.2f} ms")
        if len(samples) >= 2:
            print(f"  stdev  = {statistics.stdev(samples):.2f} ms")
        # 简易 p95
        sorted_s = sorted(samples)
        p95 = sorted_s[max(0, int(len(sorted_s) * 0.95) - 1)]
        print(f"  p95    = {p95:.2f} ms")

    print("\n  💡 关键结论:")
    print("     - 真实同机房主从延迟通常 < 1ms;跨机房可能十几 ms")
    print("     - 写主后立即读从 → 落在「延迟窗口」内就读到旧值")
    print("     - 业务对一致性敏感 → 写后短时间路由到主 / 用 WAIT / 关闭读写分离")

    master.delete(KEY_MASTER, KEY_REPLICA)


if __name__ == "__main__":
    try:
        run_demo()
    except redis.ConnectionError as e:
        print(f"❌ Redis 连接失败: {e}")
python
"""
Ch8 配套代码 3 / 3 —— WAIT N timeout 演示

⚠️ WAIT 命令需要真实主从环境才能看到「>0」的返回值。
   单机环境(无 replica)会一直返回 0,脚本会自动检测并打印 mock 演示。

WAIT N timeout 含义:
   - N        : 期望至少有多少个 replica 接收到「当前已发出的所有写命令」
   - timeout  : 最多等待多少毫秒(0 = 无限等)
   - 返回值   : 实际确认的 replica 数(可能小于 N,如果超时)

典型用法:
   r.set("key", "value")
   acked = r.wait(numreplicas=1, timeout=100)   # 最多等 100ms
   if acked < 1:
       # 没等到至少 1 个 replica 确认 → 业务可重试 / 降级 / 报警

⚠️ WAIT 不是 quorum 写:
   - 主节点宕机仍可能丢数据(Redis 没有 fsync-then-ack 协议)
   - 它只保证「至少 N 个 replica 收到 + 应用了当前 offset」
"""

import time
import redis

r = redis.Redis(host="127.0.0.1", port=6379, decode_responses=True)


def section(title: str) -> None:
    print("\n" + "=" * 70)
    print(title)
    print("=" * 70)


def detect_env() -> tuple[bool, int]:
    info = r.info("replication")
    role = info.get("role")
    n = info.get("connected_slaves", 0)
    print(f"  当前角色 : {role}")
    print(f"  在线 replica: {n}")
    is_real = (role == "master") and (n > 0)
    return is_real, n


def demo_real_wait():
    section("Demo: 真实主从环境下的 WAIT 演示")
    print("  场景:写一批 key 后调 WAIT,观察等待行为")
    for i in range(5):
        key = f"demo:wait:{i}"
        r.set(key, f"v{i}")

    print("\n  调用 r.wait(numreplicas=1, timeout=200)")
    t0 = time.perf_counter()
    acked = r.wait(num_replicas=1, timeout=200)
    elapsed = (time.perf_counter() - t0) * 1000
    print(f"    返回 = {acked}(确认数)")
    print(f"    耗时 = {elapsed:.1f} ms")
    if acked >= 1:
        print("    ✅ 至少 1 个 replica 已经同步到了当前 offset")
    else:
        print("    ⚠️ 200ms 内没有 replica 完成同步(要么没 replica,要么延迟太高)")

    print("\n  对比:r.wait(numreplicas=99, timeout=300)(要求过高)")
    t0 = time.perf_counter()
    acked = r.wait(num_replicas=99, timeout=300)
    elapsed = (time.perf_counter() - t0) * 1000
    print(f"    返回 = {acked}, 耗时 = {elapsed:.1f} ms")
    print("    ✅ WAIT 的语义:等不到也不会报错,最多等 timeout 毫秒就返回实际数")

    for i in range(5):
        r.delete(f"demo:wait:{i}")


def demo_mock_wait():
    section("Demo: 单机环境(无 replica)的 WAIT 演示")
    print("  调用真实的 WAIT 命令观察返回值(单机必为 0):\n")
    r.set("demo:wait:single", "v1")
    t0 = time.perf_counter()
    acked = r.wait(num_replicas=1, timeout=100)
    elapsed = (time.perf_counter() - t0) * 1000
    print(f"    r.wait(1, 100) → 返回 {acked}, 耗时 ≈ {elapsed:.0f} ms")
    print(f"    (等不到任何 replica,等到 timeout 后返回 0)\n")
    r.delete("demo:wait:single")

    section("MOCK:典型主从环境下 WAIT 的行为")
    rows = [
        ("SET k v + WAIT 0 0",        "0",    "0",     "立即返回(不要求确认)"),
        ("SET k v + WAIT 1 100",      "1",    "5",     "1 个 replica 5ms 内确认"),
        ("SET k v + WAIT 2 100",      "2",    "12",    "2 个 replica 都确认"),
        ("SET k v + WAIT 3 100",      "2",    "100",   "只有 2 个 replica,等 100ms 超时返回 2"),
        ("大批量写 + WAIT 1 50",      "0",    "50",    "写得快 + 50ms 太短 → 超时返回 0"),
        ("大批量写 + WAIT 1 1000",    "1",    "320",   "等 1s 充足,1 个 replica 320ms 内追上"),
    ]
    print(f"  {'命令':<32} | {'返回':>4} | {'耗时(ms)':>8} | 说明")
    print(f"  {'-'*32}-+-{'-'*4}-+-{'-'*8}-+-{'-'*30}")
    for cmd, ret, ms, desc in rows:
        print(f"  {cmd:<32} | {ret:>4} | {ms:>8} | {desc}")


def show_wait_caveats():
    section("WAIT 的 3 个隐藏坑(必看)")
    print("""\
  ① WAIT 不保证「数据不丢」
     - 即便 WAIT 1 返回成功,如果主节点立刻宕机:
         · 那 1 个收到命令的 replica 可能恰好被切成新主 → 数据保留 ✅
         · 但 fail-over 由哨兵 / Cluster 决定,不一定选中那台 → 数据可能丢 ⚠️
     - 真要强一致,得用 etcd / ZooKeeper / Spanner 这种带 quorum 写的系统

  ② WAIT 会拖慢写延迟
     - 高 QPS 写路径慎用,否则整个集群吞吐被拉低
     - 关键写 + 偶发使用比较合理

  ③ WAIT 不指定具体 replica
     - 返回 N 表示「至少 N 个 replica 同步到当前 offset」
     - 但是哪 N 个?不知道。所以也不能用它做「特定从节点优先选主」

  💡 实践建议:
     - 大多数业务用默认异步复制 + 业务幂等就够了
     - 关键写入(订单创建、支付完成)用 WAIT 1 短超时
     - 强一致核心数据 → 直接所有读写都走主,放弃读写分离
     - 金融级写入安全 → WAIT + min-replicas-to-write 双保险\
""")


if __name__ == "__main__":
    try:
        section("环境检测")
        is_real, n = detect_env()
        if is_real:
            demo_real_wait()
        else:
            demo_mock_wait()
        show_wait_caveats()
        print("\n✅ 完成。把脚本拿到真实主从(1 主 N 从)跑能看到 WAIT 返回 > 0。")
    except redis.ConnectionError as e:
        print(f"❌ Redis 连接失败: {e}")

01_check_replication.py ↗ · 02_replication_lag.py ↗ · 03_wait_command.py ↗