一群机器怎么不打架: 心跳超时触发竞选 → 多数派承认才算主 → 租约到期自动让位 → fencing token 把脑裂挡在门外
把十台机器想成一个班的值日生: 活儿只有一个人能干(跑结算/写元数据/发配置), 于是只能选一个值班班长。规矩是过半数同意才算当选(quorum), 值班有轮班时限(lease), 到点自动换班; 换班时新班长领一个递增的班次号(fencing token), 仓库只认最新班次 — 老班长从打盹里醒来, 拿旧班号连门都进不去。这套流程解决「单活组件谁来当」的冲突, 自己制造的麻烦是「误判会造成双主」, 所以选举 + 租约 + 防脑裂, 三件永远成套出现。
K8s 的 Lease 对象、etcd 的 concurrency 包、ZooKeeper 的临时节点, 都是把这三件事打包好的成品。看懂状态机(Follower→Candidate→Leader)、看懂时间轴(TTL+续期)、看懂拒绝语义(token 单调递增), 你就掌握了 K8s 控制面、Kafka Controller、所有分布式锁的共同骨架。
sess, _ := concurrency.NewSession(cli, concurrency.WithTTL(15)) e := concurrency.NewElection(sess, "jobs/settle/leader") e.Campaign(ctx) // 阻塞, 直到 etcd 多数派把你的 key 抬到第一位 // 关键: 合法性来自多数派落盘, 不是本机自封 → 天然防双主
// Raft 投票规则(伪代码): term 更高才投, 同 term 先到先投, 一生一票 if req.Term > rf.currentTerm { rf.currentTerm, rf.votedFor = req.Term, req.CandidateId } // 关键: 「term 高才投 + 一生一票」两条合起来, 数学上防住双主
n, quorum := 5, 5/2+1 // → quorum=3: 5 节点最多容忍 2 台宕机 // 关键: 两个多数派必有交集 → 双主在数学上不可能同时合法
member add / Raft 联合共识), 不能每台机器手工改配置 — 两半各改各的名册, 集群就真的裂开了。Gossip 型成员表最终一致, 共识型(etcd)强一致。 etcdctl member add infra-3 --peer-urls=http://10.0.3.11:2380 # → Member 8e9e27c5... added to cluster ef37ad9d..., 并打印新节点启动变量 # 关键: 成员变更经多数派落盘, 新节点要用 member list 的结果启动
NX(不存在才写入, 保证互斥) + PX(TTL, 持有者暴毙也能回收) + 唯一 owner 值(释放前校验, 防止删掉别人的锁)。局限: 锁服务无法在持有者失联时立刻收回授权 → 必须配 fencing token。 SET lock:cfg-push a1f2 NX PX 30000 # → OK (拿锁失败返回 nil, 不报错) # 关键: NX+PX+owner 三件缺一: 没 PX 会死锁, 没 owner 会删别人的锁
lease, _ := cli.Grant(ctx, 15) // 15s TTL 的租约 kaCh, _ := cli.KeepAlive(ctx, lease.Id) // 持续续期, 约每 TTL/3 一次回执 // 关键: 进程 kill -9 也不怕 — 到点服务端自动删除, 绝不死锁
if req.Token <= store.LastFencing() { return ErrFenced // → 旧主 token=6 被拒: "token 6 is stale (current 7)" } // 关键: 校验+拒绝必须原子地发生在存储层, 客户端自检无效
// 伪代码: 多数派不可达时的正确姿势是降级, 不是继续当家 if !majorityAlive() { enterReadonlyMode() } // 关键: 少数派一侧只读, 「两个主」从源头就不成立
jute.maxbuffer=1MB, 塞大对象会拖垮 raft。 etcdctl put /meta/route.json "$(cat 8MB-route.json)" # → etcdserver: request is too large # 关键: 协调服务只放「指针」(URL+版本号), 大对象进对象存储
WithRev(否则窗口内变更全漏); ZK 的 Watch 是一次性触发, 收到事件要重新注册。 wch := cli.Watch(ctx, "cfg/gateway/", clientv3.WithPrefix(), clientv3.WithRev(lastRev+1)) for wresp := range wch { apply(wresp.Events) } // 关键: 断线重连带 WithRev(上次版本号), 断连窗口内的事件一个不漏
// 控制面: Watch etcd → 原子替换路由表快照; 数据面: 每请求只读快照 atomic.StorePointer(&routes, newTable) // 关键: 数据面请求路径上没有 etcd — 协调服务挂了转发不受影响
for range ticker.C { // 5s 一轮, 不依赖事件, 漏了也会补齐 applyDiff(loadSpec(), loadActual()) // 期望 vs 实际 → 补差 } // 关键: diff+apply 幂等, 「对账式」设计消灭了对事件可靠性的依赖
结算 job 按 Deployment 部了 10 个副本, 全都执行一遍 = 重复扣款。用 etcd concurrency 选主: 只有 leader 进任务, 其余 9 个在 Campaign 上排队, 会话断开自动让位。
cli, _ := clientv3.New(clientv3.Config{
Endpoints: []string{"etcd-1:2379", "etcd-2:2379", "etcd-3:2379"},
DialTimeout: 5 * time.Second,
})
sess, _ := concurrency.NewSession(cli, concurrency.WithTTL(15)) // 15s 会话租约
e := concurrency.NewElection(sess, "jobs/settle/leader")
if err := e.Campaign(ctx); err != nil {
log.Fatal(err) // 阻塞, 直到多数派承认当选
}
defer e.Resign(context.Background()) // 退出主动让位, 不等 TTL
runSettlement(ctx) // 只有 leader 走到这里, 其余 9 个还堵在 Campaign
<-sess.Done() // etcd 端会话失效(暴毙/网络断)才返回, 进程退出
收益: 从「分布式锁 + 重试轮询」的 40 行轮询代码收敛成 10 行, 且进程被 kill -9 后接管耗时 ≤ TTL(15s)。
自研元数据服务要求秒级接管, 手写心跳最容易写歪(忘记 renewDeadline、忘记失主回调)。k8s.io/client-go 的 leaderelection 包把三个超时和回调都定好了。
lock := &resourcelock.LeaseLock{
LeaseMeta: metav1.ObjectMeta{Namespace: "infra", Name: "meta-master"},
Client: kube.Client(),
LockConfig: resourcelock.ResourceLockConfig{Identity: hostname},
}
leaderelection.RunOrDie(ctx, leaderelection.LeaderElectionConfig{
Lock: lock,
LeaseDuration: 15 * time.Second, // 主失联后 15s 内其他副本可接管
RenewDeadline: 10 * time.Second, // 每 10s 内必须完成一次续约
RetryPeriod: 2 * time.Second, // 续约失败的重试间隔
Callbacks: leaderelection.LeaderCallbacks{
OnStartedLeading: startServing, // 当选: 先 warmup 再对外服务
OnStoppedLeading: drainAndExit, // 失主: 停写 + 退出进程, 防赖着不走
},
})
经验值: LeaseDuration > RenewDeadline > RetryPeriod × 2, 这个配比在 K8s 自家组件(kube-scheduler/controller-manager)里已验证多年。
配置平台 8 个实例都可能触发「下发到网关」, 并发下发造成版本交错。SET NX PX + Lua 校验 owner 释放, 三件套一个不能少。
owner := uuid.NewString() // 每个持锁者唯一身份 ok, err := rdb.SetNX(ctx, "lock:cfg-push", owner, 30*time.Second).Result() if !ok { return ErrLocked } // 别人持有中, 直接拒绝而不是排队 defer func() { // 释放必须原子: 校验 owner 后再删, 防止 A 超时后误删 B 的锁 rdb.Eval(ctx, `if redis.call("GET", KEYS[1]) == ARGV[1] then return redis.call("DEL", KEYS[1]) end return 0`, []string{"lock:cfg-push"}, owner) }() pushConfig(ctx) // 锁内只做下发动作; 超过 30s 的活要靠 watchdog 续期
300 个网关实例每秒全量 GET etcd, 读 QPS 300+ 且变更生效慢 1s。改成 Watch 事件流, etcd CPU 从 60% 降到 8%, 生效从秒级变毫秒级。
var lastRev int64 for { // 外层循环只处理「Watch 断了」, 内层消费事件 wch := cli.Watch(ctx, "cfg/gateway/", clientv3.WithPrefix(), clientv3.WithRev(lastRev+1)) for wresp := range wch { for _, ev := range wresp.Events { apply(ev.Kv.Key, ev.Kv.Value) // PUT/DELETE 都会推过来 lastRev = ev.Kv.ModRevision // 记进度, 断线从这继续 } } // 走到这里 = Watch 断了; 带 lastRev 重连, 断连窗口内的事件不丢 }
曾把「每请求查一次 etcd 拿路由」写进数据面, etcd 一次 200ms 抖动 = 全站请求 +200ms。改成控制面维护本地快照, 数据面零依赖。
// 控制面(低频): Watch 到变更 → 解析 → 原子替换快照, 不加锁 func onConfigChange(ev *clientv3.Event) { tbl := parseRules(ev.Kv.Value) atomic.StorePointer(&routes, unsafe.Pointer(&tbl)) } // 数据面(高频): 每请求只读快照 — etcd 故障时继续用最后已知配置 func route(req *Request) *Target { tbl := (*RoutingTable)(atomic.LoadPointer(&routes)) return tbl.Match(req) }
量化: 数据面 p99 从 42ms 降到 3ms; etcd 停机演练 10 分钟, 转发成功率 100%(用旧快照)。
调度器必须单活(两个调度器抢同一个 pod 会绑定两次)。K8s 官方做法是 Lease 资源锁, 双实例部署也只有一个在工作, 直接抄组件配置即可。
# kube-scheduler.yaml — 单活由 Lease 实现, 主失联后备位自动接管 apiVersion: kubescheduler.config.k8s.io/v1 kind: KubeSchedulerConfiguration leaderElection: leaderElect: true resourceLock: "leases" # Lease 对象; 1.14 前的 endpoints/configmap 已废弃 resourceName: "kube-scheduler" resourceNamespace: "kube-system" leaseDuration: 15s # 失联 15s 后其他副本可接管 renewDeadline: 10s # 10s 内续不上就主动退出, 不赖位 retryPeriod: 2s # 续约重试间隔
事件会丢、回调会乱序, 一旦依赖「变更触发」逻辑就会漏。改成对账式: 每 5s 全量 diff 期望与实际, 任何单轮失败都被下一轮自愈。
func (c *Controller) reconcile(ctx context.Context) error { want, err := c.listSpec(ctx) // 期望状态: 用户声明的 replicas=5 if err != nil { return err } got, err := c.listActual(ctx) // 实际状态: 现在活着的实例 if err != nil { return err } for _, inst := range missing(want, got) { c.create(ctx, inst) // 少了就补 } for _, inst := range extra(want, got) { c.delete(ctx, inst) // 多了就删 } return nil // 不重试不补偿: 下一轮 ticker 会再 diff 一次 }
凌晨对账 job 双跑告警。排查发现 KeepAlive 协程在 etcd 重连期间 panic, 被顶层 recover 吞掉 — 没有报错、没有让位, 租约 30s 后过期, 备位无感接管, 两个实例同时跑批。
// 现场日志: lease keep alive failed: etcdserver: requested lease not found // 根因: etcd 重连期间续期报错 → 防御性 recover() 吞掉 → 无人续期 → 静默失约 go func() { defer func() { if r := recover(); r != nil { log.Fatalf("keepalive died: %v", r) // 修法: 续期死了必须退出进程 } }() for range kaCh { } // 正常时每约 TTL/3 收到一次续期回执 log.Fatal("session closed, resign") // 通道关闭 = 会话失效 → 让位退出 }()
修复后补了一条黄金指标: lease 秒数剩余量, 剩余 < 5s 告警, 再也没发生过静默失约。
选型不是背参数, 是对齐「互斥强度」和「一致性承诺」。给团队的选型注释, 贴在架构文档里。
# 选型速查(以各官方文档为准): # etcd : Raft 强一致 + Watch + Lease; K8s 原生, 客户端最成熟, 首选默认 # ZooKeeper : ZAB 强一致 + 临时节点 + 一次性 Watch; 老牌稳定, 客户端 API 偏啰嗦 # Consul : LAN/WAN Gossip, 服务发现/多数据中心强; KV 一致性语义弱于 etcd # Redis : SET NX PX 只能做「能容忍失效」的弱互斥, 无一致性保证, 必须配 fencing # 规则: 选主/元数据 → etcd/ZK; 服务发现 → Consul/K8s SVC; 弱锁 → Redis + fencing
机房网络分区 4 分钟, 降级失败的旧主与新主都收过写。恢复后主从复制报冲突, 第一步永远是: 对比两侧 gtid_executed 的差集, 找出只在旧主上的孤儿事务。
-- 现场定位(新主上执行), MySQL 8.0, gtid_mode=ON SELECT @@GLOBAL.gtid_executed; -- 新主: 3E11FA47-...:1-82001 -- 在旧主上执行: 差集 = 只在旧主落过盘的孤儿事务 SELECT GTID_SUBTRACT('3E11FA47-...:1-81990', '3E11FA47-...:1-82001'); -- → '3E11FA47-...:81991-82001' 这 11 个事务就是脑裂窗口的孤儿写 -- 处置: 旧主先钉死只读, 孤儿事务按业务语义人工核对后补写/丢弃 SET GLOBAL read_only = ON; SET GLOBAL super_read_only = ON;
复盘结论: 事故根因是旧主降级失败(read_only 未生效), 终极解法是存储层 fencing, 而不是更快的切换。
PX 过期 + 主动释放。 # 错: SET lock:job w1 NX # 崩溃后无人释放 → 永久 ErrLocked # 对: SET lock:job w1 NX PX 30000 # → OK, 30s 后即使崩溃也自动释放
# 错: DEL lock:job # 把 B 的锁删了 → 互斥被击穿 # 对: EVAL "if redis.call('GET',KEYS[1]) == ARGV[1] then return redis.call('DEL',KEYS[1]) end return 0" 1 lock:job <owner>
// 错: acquire(30s); runLongBatch() // 任务 3min ≫ 锁 30s, 中途已易主 // 对: go watchdog(lock, 10*time.Second) // 每 10s 续期, 失败立刻停任务
// 错: if !ping() { elect() } // 单次超时 → Full GC 就误切 // 对: if miss >= 3 { elect() } // 连续 3 次失败才进入竞选
// 错: if !ping(master) { becomeMaster() } // 两半分区都成立 → 双主 // 对: e.Campaign(ctx) // 只有 etcd 多数派承认的一侧能当选
too many open files, fd 数每天涨. 原因: 断线重连逻辑里新建 Watch 流, 旧的从不 Cancel. 正解: Watch 流用 ctx 管生命周期, 重连前先 cancel 上一轮。 // 错: for { cli.Watch(ctx, key, ...) } // 每次重连新建, 旧流不释放 // 对: wctx, cancel := context.WithCancel(ctx) // 每轮循环先 cancel() 上一轮的 Watch 流, 再重建
// 错: defer func() { recover() }() // 吞 panic → 静默失约 // 对: log.Fatalf("keepalive died: %v", r) // 续期线程死了就喊出来
# 错: ttl=3s, 失败立即重试 # 一次抖动 → 全体一起重选 # 对: ttl=15s + rand(0,2s) 退避 # 抖动期间大家错峰, 不共振
sess.Done()。 // 错: if time.Now().After(expireAt) { resign() } // 本地时钟会说谎 // 对: <-sess.Done() // etcd 服务端裁决, 到点关闭通道
# 错: cluster: n1,n2 # 挂 1 台 → 无多数派 → 全停 # 对: cluster: n1,n2,n3 # 挂 1 台仍 2/3, 服务继续
// 错: store.Put(req) // 谁的请求都收 → 双写错乱 // 对: if req.Token <= store.LastFencing() { return ErrFenced }
etcdserver: request is too large, raft 日志膨胀, 快照越传越慢. 原因: etcd/ZK 是元数据存储不是对象存储(etcd 默认单 value ≤ 1.5MiB). 正解: 协调服务只放指针, 大对象进 OSS/S3。 # 错: etcdctl put /meta/route.json "$(cat 8MB.json)" # → etcdserver: request is too large # 对: etcdctl put /meta/route "s3://cfg/route@v17" # 放指针
// 错: for { time.Sleep(time.Second); cli.Get(ctx, key) } // QPS 随实例数涨 // 对: wch := cli.Watch(ctx, key, clientv3.WithRev(rev+1)) // 推送, 0 轮询
// 错: l1 := acquire("lock:job"); l2 := acquire("lock:job") // 第二次失败 // 对: ctx 带上锁凭证向下传, 内层函数见到凭证直接复用
// 错: e.Campaign(ctx); serve() // 当选即服务, 缓存是冷的 // 对: warmup(); atomic.StoreInt32(&ready, 1); serve() // 探针先过
// 错: A: lock(r1); lock(r2) B: lock(r2); lock(r1) // 互等 → 卡死 // 对: keys := sortKeys(need); for k := range keys { lock(k) } // (全局按字典序加锁)
# 错: TTL 3s # 抖动就重选, 稳定性崩 # 对: TTL 15s + Watch DELETE 事件触发竞选 # 稳且接管毫秒级
// 错: etcd 写 leader=nodeA, 同时 redis.set("leader", nodeA) // 会漂移 // 对: 只信 etcd; Redis 里的是只读缓存, 冲突时以 etcd 为准
# 错: registry.list("svc")[0] 当主用 # AP 列表无互斥 → 双主 # 对: 选主走 etcd Campaign; 注册中心只回答「有哪些实例在线」
# 错: etcd 无 auth, 谁都能 put /jobs/settle/leader/ # 主被顶掉 # 对: etcdctl role grant job-svc readwrite /jobs/settle/ # 限定命名空间