Kafka / K8s / etcd / Redis / ZooKeeper / MySQL 扒开看都是同几招: 心跳→选举→租约 · 分片→复制→Quorum · 重试→幂等 · 限流→熔断 — 八大核心问题一张图
把 Kafka、Kubernetes、etcd、Redis 想成不同菜系, 后厨却只有二十口锅: 选举、心跳、租约、分片、复制、共识、限流、熔断、幂等……单机时代一个问题一个函数, 分布式时代要先接受三条公理——机器随时会没、网络会丢会乱序、没有统一的"现在"。
八大核心问题就是这套世界观的目录: 从"谁还活着"出发, 一路推出选主、数据分布、复制、一致性、容错、流量治理和多服务协作。看懂这 20 个基础机制, 再遇到任何新中间件都是"换皮重认"——哦, 这货本质就是选举 + 日志复制 + 分片 + 租约——而不是从零再学一套名词。
// etcd 选主: 所有副本竞争同一个前缀 key, 谁被多数派接受谁是主 e := concurrency.NewElection(sess, "batch/settle/leader") e.Campaign(ctx) // → 阻塞直到当选; 会话结束自动让位
for range ticker.C { if err := ping(ctx); err != nil { misses++ } // 关键: 记次数, 不一次定死 } // → misses >= 3 才标记 suspect, 进选举流程
// etcd: 30s TTL 的租约 + 心跳续期 lease, _ := cli.Grant(ctx, 30) cli.KeepAlive(ctx, lease.Id) // 关键: 进程暴毙, 租约到点自动释放, 不会死锁
NX(互斥) + PX(过期) + owner 值(身份), 三件缺一: 没过期会死锁, 没 owner 会删别人的锁。 SET lock:job w1 NX PX 30000 // → OK; 拿到锁且 30s 自动过期 // 关键: value 存 owner 身份, 释放时校验后再删
if req.Token < store.LastFencing() { return ErrFenced // → 旧主的写被拒: "token 6 < 7" }
w + r > n 读写集合必有交集, 读就能见到最新值。代价: w、r 越大延迟越高、可用节点数要求越高。 n, w, r := 3, 2, 2 // 关键: w + r > n → 读写交集非空 // → 3 副本挂 1 台仍满足多数派, 服务继续
1/N 数据(mod-N 要迁大半)。虚拟节点把每台机器拆成 100~200 个环上点, 抹平机器异构导致的数据倾斜。 h := crc32(key) i := sort.Search(vnodes, func(i int) bool { return vnodes[i].hash >= h }) // 关键: 顺时针找第一个 >= h 的虚拟节点; 扩容只动 1/N
-- MySQL 半同步: 至少 1 从库收到 binlog 才向客户端确认 SET GLOBAL rpl_semi_sync_master_enabled = 1; -- 关键: 异步复制切主丢数据的解药, 不是备份
// P 发生: 取 C(拒绝服务) 或 A(容忍旧数据), 二选一 // Else 正常: 取 L(低延迟) 或 C(强一致), 仍要选 // 关键: CAP 不是"三选二", 平时 CA 是伪命题
CREATE TABLE pay_idem ( idem_key VARCHAR(64) NOT NULL, PRIMARY KEY (idem_key) -- 关键: 唯一索引是幂等最后的闸门 ); -- → 重复 INSERT 报 1062 Duplicate entry, 捕获后直接返回上次结果
tokens := math.Min(cap, tokens + rate*dt) // 关键: 匀速补充, 上限 cap if tokens >= 1 { tokens--; allow() } else { reject() }
Closed(正常放行) → 错误率超阈值 → Open(直接快速失败, 不打下游) → 冷却后 → Half-Open(放少量探测请求, 成功则闭合)。本质: 快速失败保护自己也保护下游。 Closed --错误率>50%--> Open --冷却30s--> Half-Open --探测过--> Closed
// 关键: Open 期间请求要走兜底, 不是报错给用户peers[rand.Intn(len(peers))].Exchange(state)
// 关键: 无中心也能同步; Redis Cluster / Consul 都用它lc = max(lc, msg.clock) + 1 // 收到消息时推进 // 关键: 物理时钟会说谎(回拨), 因果序不会
定时任务多副本部署却没有选主, 谁的 cron 先到谁执行 → 数据重复处理。用 etcd 租约 + Election 选出唯一执行者。
// 结算任务部署 10 副本, 只允许 1 个执行: etcd lease + election sess, err := concurrency.NewSession(cli, concurrency.WithTTL(15)) if err != nil { log.Fatal(err) } e := concurrency.NewElection(sess, "batch/settle/leader") go func() { for { if err := e.Campaign(ctx); err != nil { continue } // 竞选, 阻塞到当选 runDailySettle(ctx) // 当选: 唯一执行批任务 e.Resign(ctx) // 主动让位, 发布时平滑换主 } }() <-ctx.Done() // 关键: 进程死→会话死→lease 过期, 其他副本 15s 内自动补位
上线后 10 副本每天只有 leader 跑批, 主宕机切换窗口 = 会话 TTL 15 秒。
心跳 5s 超时即切主 → 旧主 GC 结束醒来继续写, 新旧两主并存。修法: client-go 的 lease 选举参数放宽 + 存储层校验 fencing。
lock := &resourcelock.LeaseLock{
LeaseClient: coordClient,
LeaseMeta: metav1.ObjectMeta{Namespace: "prod", Name: "batch-runner"},
Identity: id,
}
leaderelection.RunOrDie(ctx, leaderelection.LeaderElectionConfig{
Lock: lock,
LeaseDuration: 15 * time.Second, // 15s 没续约才判死 — 别设 3s, 一次 GC 就误切
RenewDeadline: 10 * time.Second,
RetryPeriod: 2 * time.Second,
Callbacks: leaderelection.LeaderCallbacks{
OnStartedLeading: runSettle,
OnStoppedLeading: func() { log.Fatal("lost leadership, exit") }, // 失主即自杀防双主
},
})
// 关键: 下游存储仍要校验 fencing token — GC 20s 的旧主写进来直接 409 拒绝
hash(orderID) % 4 改成 % 8, 3/4 的行路由全变, 停机迁移要 3 天。换成一致性哈希环, 只有新旧路由不一致的行才搬。
type Ring struct{ vnodes []vnode } // 每个物理库挂 200 个虚拟节点 func (r *Ring) Get(key string) *DB { h := crc32.ChecksumIEEE([]byte(key)) i := sort.Search(len(r.vnodes), func(i int) bool { return r.vnodes[i].hash >= h }) return r.vnodes[i%len(r.vnodes)].db // 关键: 顺时针第一个 >= h 的 vnode } // 迁移期双写: 新旧库都写, 读走旧库, 后台只搬 route 变化的行, 校验后切读 for _, order := range backlog { if ring8.Get(order.ID) == ring4.Get(order.ID) { continue } migrate(ctx, order) } // 收益: 4 亿行只搬约 1/8, 迁移窗口从 3 天缩到 4 小时
写主库、读从库, 复制延迟 3s 内用户看到旧数据, 工单刷屏。做"读己之写": 记住写操作的 GTID, 从库追上才读, 追不上读主。
// 写: 记录本次写入的全局事务序 db.Exec("UPDATE users SET avatar=? WHERE id=?", avatar, uid) gtid, _ := queryScalar(db, "SELECT @@GLOBAL.gtid_executed") session.Set("w_gtid", gtid) // 读: 从库先追上这个 GTID 才服务本次读, 1s 追不上降级读主 err := replica.QueryRow( "SELECT WAIT_FOR_EXECUTED_GTID_SET(?, 1)", gtid).Err() if err != nil { row = master.Query("SELECT avatar FROM users WHERE id=?", uid) } // 关键: 只对"写后立刻读"的请求付出这个代价, 其余照走从库
单条消息处理慢, 两次 poll 间隔超过 max.poll.interval.ms → 被踢出组 → 全组停止消费重新分配 → 雪崩式堆积。调参 + 换增量重平衡。
# consumer.properties — 重平衡风暴的三板斧 session.timeout.ms=30000 # 30s 没心跳判死 (broker 侧判定) heartbeat.interval.ms=3000 # 心跳间隔 ≈ session/3 max.poll.interval.ms=600000 # 两次 poll 上限: 给慢处理留足时间 max.poll.records=500 # 关键: 少拉快提, 别一次 5000 条噎死自己 partition.assignment.strategy=CooperativeStickyAssignor # 增量重平衡 (2.4+): 只挪动受影响的分区, 其余消费者不停
重平衡从每天 288 次降到 0 次(只在新消费者加入时增量发生), 消费延迟稳定在 3 万条以内。
无超时无熔断: 下游 P99 3s, 上游线程池 200 被慢慢占满, 最后全站 503。三件套: 超时递减 + gobreaker 熔断 + 兜底返回。
cb := gobreaker.NewCircuitBreaker(gobreaker.Settings{
Name: "inventory",
Interval: 10 * time.Second, // Closed 态错误统计窗口
Timeout: 30 * time.Second, // Open 冷却 30s 后进 Half-Open
ReadyToTrip: func(c gobreaker.Counts) bool {
return c.Requests >= 50 && float64(c.TotalFailures)/float64(c.Requests) > 0.5
},
})
resp, err := cb.Execute(func() (any, error) {
ctx, cancel := context.WithTimeout(ctx, 800*time.Millisecond) // 上游剩 1s, 我只花 0.8s
defer cancel()
return invClient.Reserve(ctx, req)
})
if errors.Is(err, gobreaker.ErrOpenState) {
return fallbackEstimate(req) // 关键: 熔断走兜底, 不是把错误抛给用户
}
限流做在应用内存里, 20 台机器 = 20 份配额互不知晓; 改成 Redis + Lua 把"判定+扣减"做成原子操作, 集群级精确限流。
-- KEYS[1]=令牌桶 ARGV: rate 令/秒, capacity 桶容量, now(ms), 本次请求 n local rate, cap = tonumber(ARGV[1]), tonumber(ARGV[2]) local now = tonumber(ARGV[3]) local b = redis.call('HMGET', KEYS[1], 'tokens', 'ts') local tokens = math.min(cap, (tonumber(b[1]) or cap) + (now - (tonumber(b[2]) or now)) / 1000 * rate) local allowed = 0 if tokens >= tonumber(ARGV[4]) then tokens = tokens - tonumber(ARGV[4]); allowed = 1 end redis.call('HMSET', KEYS[1], 'tokens', tokens, 'ts', now) redis.call('PEXPIRE', KEYS[1], 60000) return allowed -- 关键: 判定+扣减一个原子脚本, 多机也不会放超
第三方支付 at-least-once 重试回调是常态, 处理函数不幂等就会重复加钱。模板: 唯一索引挡重, 重复回调直接回成功止住对方重试。
func HandlePayCallback(req PayNotify) error { // 幂等闸门: 回调带 order_no, 唯一索引挡重 res, err := db.Exec( `INSERT IGNORE INTO pay_callback(order_no, amount, status) VALUES (?, ?, 'paid')`, req.OrderNo, req.Amount) if err != nil { return err } if n, _ := res.RowsAffected(); n == 0 { return nil // 关键: 重复回调当成功返回, 上游就不再重试 } return creditUser(req) // 只有首个回调走到真加钱 }
先写订单库再发 MQ, 中间服务崩溃就永久丢事件。Outbox 模式: 业务数据与事件同事务落库, relay 异步投递, 用至少一次 + 消费幂等换最终一致。
func CreateOrder(o Order) error { tx, _ := db.Beginx() defer tx.Rollback() tx.Exec("INSERT INTO orders(order_no, sku, qty) VALUES(?,?,?)", o.No, o.Sku, o.Qty) tx.Exec("INSERT INTO outbox(topic, payload) VALUES('order.created', ?)", toJSON(o)) // 关键: 业务表和事件同生共死, 同一事务提交 return tx.Commit() } // relay 轮询投递 MQ (至少一次), 消费端用去重表兜住重复: // SELECT * FROM outbox WHERE sent=0 ORDER BY id LIMIT 100 FOR UPDATE SKIP LOCKED
大堆 JVM 启动要 4 分钟, 默认探针 30s 判死 → Pod 被反复杀 → 越杀越慢。三层探针各司其职: startup 给启动预算, readiness 摘流量, liveness 只查进程自身。
# deployment.yaml — 启动慢的应用必须配 startupProbe startupProbe: httpGet: { path: /healthz, port: 8080 } periodSeconds: 5 failureThreshold: 60 # 5s × 60 = 300s 启动预算, 期间不杀不摘 readinessProbe: httpGet: { path: /ready, port: 8080 } periodSeconds: 5 failureThreshold: 3 # 未就绪只摘流量不重启 — 扛住预热窗口 livenessProbe: httpGet: { path: /live, port: 8080 } periodSeconds: 10 # 关键: 只查"自己还能动吗", 不碰 DB/下游
SET key val 无 TTL, 锁成了"永生锁". 正解: SET ... NX PX 30000 三件套。 // 错: SET lock:job worker-1 → 持有者崩溃, 全集群永久阻塞 // 对: SET lock:job worker-1 NX PX 30000 → OK, 30s 后自动可抢
KeepAlive)。 // 错: PX 30000 但任务要跑 90s → 第 30s 第二实例抢到, 双跑 // 对: lease,_ := Grant(ctx,30); KeepAlive(ctx, lease.Id) → 自动续期
// 错: DEL lock:job → 删掉的是 B 的锁 // 对: if redis.call('GET',KEYS[1])==ARGV[1] then redis.call('DEL',KEYS[1]) end
// 错: if time.Since(last) > 5*time.Second { elect() } → 一次 STW 就误切 // 对: suspicion 达阈值才选举, 且存储层拒 token < last 的旧主写
// 错: for { if err != nil { retry() } } → 重试风暴 // 对: backoff = min(base * 2^n, 10s) + rand(0, base) → 错峰衰减
// 错: client.Do(req) 超时 → 原样重发 → 扣款 ×2 // 对: req.Header.Set("Idempotency-Key", orderNo) → 服务端唯一索引挡重
// 错: 网关×3 + 服务×3 + SDK×3 = 27 倍写放大 // 对: retryBudget: 10% # 超过预算直接失败, 链路只留一层重试
// 错: 每层 WithTimeout(3s) → 5 层链路用户等 15s // 对: ctx 从入口传入, 子预算 = 剩余 deadline × 0.8
readiness。 # 错: livenessProbe httpGet /healthz → /healthz 内部连 DB # 对: liveness 查 /live(纯进程); DB 检查放 readiness /ready
// 错: db := hash(uid) % 4 → % 8 后 3/4 的行要搬家 // 对: 一致性哈希环: 只迁移新旧路由不一致的 ~1/8 行
tenant_id + ORM 拦截器/RLS 兜底。 // 错: SELECT * FROM orders WHERE id=1001 → 可能查出别家租户的 // 对: SELECT ... WHERE id=1001 AND tenant_id='t_a' → 串号被挡
// 错: 所有读都打 product:9527 → 单分片 50w QPS // 对: product:9527:{0..15} 随机后缀拆 16 份 + 本地缓存 100ms
-- 错: 默认异步复制 → 主挂时最新 2s 的写可能没到从库 -- 对: SET GLOBAL rpl_semi_sync_master_enabled=1; 切主先验 GTID
// 错: 旧主分区期间继续写 → 双份数据 // 对: if !haveQuorum() { stepDown() }; 写带 token, 旧 token 被拒
// 错: UPDATE 后 SELECT 走从库 → 延迟窗口内读到旧昵称 // 对: session 写后 5s 内强制读主 / WAIT_FOR_EXECUTED_GTID_SET
id=-1 请求, 缓存永远 miss, DB 每秒 5 万次空查询. 原因: 不存在的数据没有缓存价值但照样穿库. 正解: Bloom Filter 拦截 + 空值短 TTL 缓存。 # 错: GET user:-1 miss → SELECT ... WHERE id=-1 → 打爆 DB # 对: 布隆过滤器拦截 + SET user:-1 "" EX 60 # 空值缓存 60s
// 错: 1w 个请求同时 miss → 同时 SELECT → 连接池爆 // 对: g.singleflight(key, loadFromDB) → 1 个回源, 其余等结果
// 错: SET k v EX 3600 → 1 小时后同一秒集体失效 // 对: EX (3600 + rand(0,600)) → 过期点错开, DB 压力摊平
# 错: DEL cache; UPDATE db → 间隙读把旧值 SET 回缓存 # 对: UPDATE db; DEL cache → 后续读 miss 回源拿到新值
// 错: 处理成功但 offset 提交失败 → 重投 → 又发一件货 // 对: INSERT INTO inbox(msg_id) 唯一键 → 冲突即跳过, 只处理一次