分布式系统设计
#系统 · #分布式 · #CAP · #BASE · #一致性 · #共识 · #微服务
当单机无法承载业务规模,分布式是唯一的出路。但分布式会引入网络分区、时钟不同步、部分故障等问题。本节从 CAP 定理出发,梳理分布式系统的核心概念、一致性模型、共识算法和架构模式。
1. CAP 定理
1.1 定义
2000 年 Eric Brewer 提出,后来由 Lynch 等人证明:
在分布式系统中,当发生网络分区(Partition)时,
只能在一致性(Consistency)和可用性(Availability)之间二选一。mermaid
graph TD
CAP["CAP 定理"]
CAP --> C["一致性 (Consistency)<br/>所有节点同时看到相同数据"]
CAP --> A["可用性 (Availability)<br/>每个请求都能获得响应"]
CAP --> P["分区容错 (Partition Tolerance)<br/>系统在网络分区时仍能工作"]
P -.->|"网络分区发生时"| CHOICE["必须选 C 或 A"]
CHOICE --> CP["CP 系统<br/>保证一致性,牺牲可用性"]
CHOICE --> AP["AP 系统<br/>保证可用性,牺牲一致性"]1.2 CP vs AP
| CP(一致性优先) | AP(可用性优先) | |
|---|---|---|
| 代表 | ZooKeeper, Etcd, HBase | Cassandra, DynamoDB, Eureka |
| 行为 | 分区时少数节点不可用 | 分区时允许读写不一致 |
| 适合 | 配置中心、分布式锁、元数据 | 用户数据、社交动态、日志 |
| 风险 | 可用性降低,影响面大 | 数据冲突,甚至丢数据 |
1.3 常见误解
误解:CAP 只能三选二。
正解:分区是不可避免的。真正的选择是:
- 正常运行时:C + A + P(全部满足)
- 分区时:选择 C 或 A2. BASE 理论
对 CAP 中 AP 系统的进一步扩展:
| 缩写 | 含义 | 说明 |
|---|---|---|
| BA | Basically Available | 基本可用:允许降级,但不宕机 |
| S | Soft State | 软状态:允许中间状态,不要求时刻一致 |
| E | Eventually Consistent | 最终一致性:经过一段时间后数据收敛 |
ACID (传统数据库): BASE (分布式):
- Atomicity 原子性 - Basically Available
- Consistency 强一致性 - Soft State
- Isolation 隔离性 - Eventually Consistent
- Durability 持久性3. 一致性模型
从强到弱:
强一致性 ──────────────────────────────→ 弱一致性
│ │ │
线性一致性 顺序一致性 最终一致性
(Linearizable) (Sequential) (Eventual)
│
│ 读总是返回最近写入的值
│ 所有操作看起来像在一个全局时钟下执行
│
│ 代表:Etcd (Raft), Zookeeper (ZAB)3.1 线性一致性(Linearizability)
go
// 例:非一致性系统的异常行为
// 时刻 T1: Client A 写入 x = 2
// 时刻 T2: Client B 读取 x → 返回 1(旧值!)
// 时刻 T3: Client B 读取 x → 返回 2
// 线性一致性系统:T2 时刻一定会返回 23.2 最终一致性(Eventual Consistency)
场景:DNS 记录更新
t=0: A 记录改为 1.2.3.4(权威 DNS)
t=1: 用户 A 查询 → 返回旧 IP(命中旧缓存)
t=2: 用户 B 查询 → 返回旧 IP
t=3: TTL 过期 → 全部缓存刷新
t=4: 所有用户查询 → 返回新 IP
→ 最终所有用户都能看到 1.2.3.43.3 因果一致性(Causal Consistency)
有因果关系的操作必须按顺序看到,无因果关系的可以任意顺序。
因果关系: 发帖 → 回复 → 必须按序
非因果关系: 张三发帖 → 李四发帖 → 可以乱序4. 共识算法
4.1 为什么需要共识
mermaid
graph LR
A["3个节点"] --> B{"谁是 Leader?"}
B -->|"节点1: 我是"| C1["冲突!"]
B -->|"节点2: 我是"| C2["冲突!"]
B -->|"节点3: 我是"| C3["冲突!"]
C1 --> D["共识算法<br/>保证只有一个 Leader"]4.2 Paxos
Leslie Lamport 1990 年提出。理解门槛极高,分 Basic Paxos、Multi-Paxos 等变体。
角色:
Proposer: 提议值
Acceptor: 接受/拒绝提议
Learner: 学习最终值
Basic Paxos 两阶段:
Phase 1 (Prepare): Proposer → Acceptor: "我要提议,编号 N"
Acceptor → Proposer: "OK,我已接受的最大编号是 M"
Phase 2 (Accept): Proposer → Acceptor: "接受编号 N, 值 V"
Acceptor → Proposer: "Accepted"4.3 Raft
为可理解性设计的共识算法,将问题分解为三个子问题:
mermaid
flowchart TB
subgraph LeaderElection["1. Leader 选举"]
L1["随机超时 → 发起投票"]
L2["获得多数选票 → 成为 Leader"]
L3["Leader 定期发心跳"]
end
subgraph LogReplication["2. 日志复制"]
R1["客户端 → Leader"]
R2["Leader → Followers (AppendEntries)"]
R3["多数确认 → 提交"]
end
subgraph Safety["3. 安全性"]
S1["选举限制:只有最新日志的节点能成为 Leader"]
S2["Leader 不覆盖已提交的日志"]
end
LeaderElection --> LogReplication
LogReplication --> SafetyRaft 关键机制
| 机制 | 作用 |
|---|---|
| Term(任期) | 逻辑时钟,每次选举递增 |
| 随机超时 | 避免多个节点同时发起选举(分裂投票) |
| AppendEntries | 心跳 + 日志复制,Leader 用它维持权威 |
| committed vs applied | 已提交(多数确认)≠ 已应用(写入状态机) |
成员变更(Joint Consensus)
mermaid
flowchart LR
C1["C_old = {A, B, C}"] --> T["过渡期: C_old+new = {A,B,C,D,E}"]
T --> C2["C_new = {A, B, C, D, E}"]两阶段过渡,确保任意时刻不会有"两个 Leader 分别被新老配置中的多数选中"。
5. 分布式事务
详见 [分布式事务],这里简述核心:
5.1 两阶段提交(2PC)
mermaid
sequenceDiagram
participant C as Coordinator
participant P1 as Participant 1
participant P2 as Participant 2
C->>P1: Prepare
C->>P2: Prepare
P1-->>C: YES
P2-->>C: YES
C->>P1: Commit
C->>P2: Commit问题:协调者单点故障 → 参与者不确定是否提交 → 资源阻塞。
5.2 TCC(Try-Confirm-Cancel)
Try: 预留资源(冻结库存)
Confirm: 确认使用(扣减库存)
Cancel: 释放资源(解冻库存)
每阶段需要业务幂等,适合金融场景。5.3 Saga
长事务拆分为一组本地事务:
T1 → T2 → T3
如果 T3 失败:
T3 补偿 → T2 补偿 → T1 补偿
每步有对应的补偿操作(不是回滚,而是"对冲")。6. 分布式 ID
6.1 方案对比
| 方案 | 优点 | 缺点 |
|---|---|---|
| UUID | 无中心,简单 | 非单调,字符串存索引性能差 |
| DB 自增 | 单调递增,简单 | 单点瓶颈 |
| 号段模式 | 批量取 ID,接近单调 | 需 DB 支持 |
| Snowflake | 高性能,趋势递增 | 依赖机器时钟 |
| Redis INCR | 简单,自增 | 持久化问题 |
6.2 Snowflake 算法
64 位 ID 结构:
┌──────────┬──────────┬──────────┬──────────┐
│ 1bit │ 41bit │ 10bit │ 12bit │
│ 符号位 │ 时间戳 │ 机器ID │ 序列号 │
│ (未使用) │ (毫秒) │ (1024) │ (4096) │
└──────────┴──────────┴──────────┴──────────┘
每毫秒可生成 4096 个 ID,理论 QPS = 4096000go
type Snowflake struct {
mu sync.Mutex
epoch int64 // 起始时间戳(2020-01-01)
machineID int64
sequence int64
lastMs int64
}
func (s *Snowflake) NextID() int64 {
s.mu.Lock()
defer s.mu.Unlock()
now := time.Now().UnixMilli()
if now == s.lastMs {
s.sequence = (s.sequence + 1) & 0xFFF
if s.sequence == 0 {
// 同一毫秒用完,等下一毫秒
for now <= s.lastMs {
now = time.Now().UnixMilli()
}
}
} else {
s.sequence = 0
}
s.lastMs = now
return ((now - s.epoch) << 22) | (s.machineID << 12) | s.sequence
}6.3 时钟回拨问题
场景:NTP 同步导致时钟回拨 2 秒
应对:
1. 拒绝生成,等待时钟追上(简单但影响可用性)
2. 使用拓展序列号(号段模式下不受影响)
3. 美团 Leaf:在 ZK 中记录当前时间戳,如果回拨则报警并拒绝7. 分布式锁
7.1 Redis 分布式锁(Redlock 的争议)
go
// 单节点 Redis 锁(非生产级)
func acquireLock(key string, ttl time.Duration) bool {
return redis.SetNX(ctx, key, "locked", ttl).Val()
}
// Redlock:向 N 个独立 Redis 实例依次加锁
// 如果 > N/2 成功 → 获得锁
// 问题:Martin Kleppmann 论证其在 GC pause/时钟跳变时不是安全的7.2 Etcd 分布式锁
go
// 基于 Lease + 事务的分布式锁
func (c *Client) Lock(ctx context.Context, key string) (*Lease, error) {
// 1. 创建 Lease(带 TTL)
lease, err := c.Grant(ctx, int64(ttl.Seconds()))
if err != nil {
return nil, err
}
// 2. 用事务 CAS 写入 key(if key not exists → put key with lease)
txn := c.Txn(ctx).If(clientv3.Compare(clientv3.CreateRevision(key), "=", 0))
txn = txn.Then(clientv3.OpPut(key, "locked", clientv3.WithLease(lease.ID)))
resp, err := txn.Commit()
if err != nil || !resp.Succeeded {
return nil, ErrLocked
}
// 3. 持续续约
keepAlive, _ := c.KeepAlive(ctx, lease.ID)
go func() { for range keepAlive {} }()
return lease, nil
}7.3 选型
| Redis | Etcd | Zookeeper | |
|---|---|---|---|
| 一致性 | 弱(异步复制) | 强(Raft) | 强(ZAB) |
| 性能 | 极高 | 高 | 中 |
| 实现 | 简单 | 中等 | 复杂 |
| 适合 | 高性能要求、容忍短时不一致 | 强一致性要求 | 大量历史遗留系统 |
8. 服务发现与注册
mermaid
flowchart TB
S["服务启动"] --> R["注册到注册中心<br/>(Etcd/Consul/Nacos)"]
R --> H["定期心跳续约"]
C["客户端"] --> D["从注册中心获取服务列表"]
D --> LB["负载均衡选择实例"]
LB --> CALL["发起 RPC 调用"]
H -.->|"心跳超时"| R8.1 注册中心对比
| Etcd | Consul | Nacos | Zookeeper | |
|---|---|---|---|---|
| CAP | CP | CP | CP+AP 可切换 | CP |
| 健康检查 | Lease | Agent 心跳 | Agent 心跳 | 临时节点 |
| 多数据中心 | ❌ | ✅ | ✅ | ❌ |
| 配置管理 | ✅(需自己做) | ✅(KV) | ✅✅ | ✅(需自己做) |
9. 限流算法
9.1 算法对比
| 算法 | 原理 | 优点 | 缺点 |
|---|---|---|---|
| 固定窗口 | N 秒内最多 M 次 | 简单 | 边界突发 |
| 滑动窗口 | 过去 N 秒内最多 M 次 | 平滑 | 实现稍复杂 |
| 漏桶 | 恒定速率流出 | 严格平滑 | 无法应对突发 |
| 令牌桶 | 恒定速率放入令牌 | 允许一定突发 | 实现稍复杂 |
9.2 令牌桶实现
go
type TokenBucket struct {
rate float64 // 每秒放入 tokens
capacity float64 // 桶容量
tokens float64 // 当前令牌数
lastTime time.Time
mu sync.Mutex
}
func (tb *TokenBucket) Allow() bool {
tb.mu.Lock()
defer tb.mu.Unlock()
now := time.Now()
elapsed := now.Sub(tb.lastTime).Seconds()
tb.tokens += elapsed * tb.rate
if tb.tokens > tb.capacity {
tb.tokens = tb.capacity
}
tb.lastTime = now
if tb.tokens >= 1 {
tb.tokens--
return true
}
return false
}10. 微服务架构模式
10.1 服务拆分原则
| 原则 | 说明 |
|---|---|
| 单一职责 | 一个服务只做一件事 |
| 高内聚低耦合 | 内部紧密,外部松散 |
| 按业务域 | DDD 限界上下文划分 |
| 数据独立 | 每个服务自己的数据库(Database per Service) |
10.2 通信模式
同步:HTTP/gRPC — 请求-响应
异步:消息队列 — 事件驱动
混合:同步查询 + 异步通知(CQRS)10.3 常见模式
| 模式 | 解决什么问题 |
|---|---|
| API Gateway | 统一入口、认证、限流、路由 |
| BFF(Backend for Frontend) | 为每种客户端定制 API |
| CQRS | 读写分离,不同数据模型 |
| Event Sourcing | 存储事件而非状态,可重放 |
| SAGA | 分布式事务,补偿机制 |
| Circuit Breaker | 防止级联故障 |
| Bulkhead | 资源隔离(线程池、连接池) |
| Sidecar | 将基础设施能力与业务解耦 |
10.4 熔断器(Circuit Breaker)
mermaid
stateDiagram-v2
[*] --> Closed: 正常状态
Closed --> Open: 错误率 > 阈值
Open --> HalfOpen: 超时后尝试
HalfOpen --> Closed: 试探请求成功
HalfOpen --> Open: 试探请求失败go
// 简易熔断器
type CircuitBreaker struct {
failures int
threshold int // 失败阈值
timeout time.Duration // 熔断时间
lastFailure time.Time
state int // 0=Closed, 1=Open, 2=HalfOpen
mu sync.Mutex
}
func (cb *CircuitBreaker) Call(fn func() error) error {
cb.mu.Lock()
if cb.state == 1 { // Open
if time.Since(cb.lastFailure) > cb.timeout {
cb.state = 2 // → HalfOpen
} else {
cb.mu.Unlock()
return ErrCircuitOpen
}
}
cb.mu.Unlock()
err := fn()
cb.mu.Lock()
defer cb.mu.Unlock()
if err != nil {
cb.failures++
cb.lastFailure = time.Now()
if cb.failures >= cb.threshold {
cb.state = 1 // → Open
}
} else {
cb.failures = 0
cb.state = 0 // → Closed
}
return err
}10.5 幂等性设计(Idempotency)
幂等性是分布式系统的第一道防线——它的重要性怎么强调都不为过。在分布式环境中,网络超时、进程崩溃、消息重试是常态而非异常。如果没有幂等性保护,一次"支付请求超时导致的重试"可能让用户被扣两次款。这不是 bug,这是分布式系统的物理现实——你需要从设计层面接受"重试必然发生",并在业务逻辑中做好准备。
下面的时序图展示了经典场景:客户端以为请求失败于是重试,但服务端其实已经处理成功——幂等性保证了第二次请求不会重复扣款:
mermaid
flowchart TB
Client["客户端"] -->|"支付请求 (txnId=abc123)"| Server["服务端"]
Server -->|"处理成功, 返回 OK"| X[❌ 网络丢包]
Client -->|"重试: 支付请求 (txnId=abc123)"| Server
Server -->|"检查 txnId=abc123<br/>→ 已处理, 返回缓存结果"| Client幂等实现策略:
| 策略 | 原理 | 适用场景 | 优点 | 缺点 |
|---|---|---|---|---|
| 唯一键/ID | DB 唯一约束 + INSERT IGNORE/ON DUPLICATE KEY | 支付、下单 | 最简单可靠 | 依赖 DB |
| 版本号/乐观锁 | UPDATE SET v=v+1 WHERE v=old_v | 并发更新 | 无额外存储 | 冲突需重试 |
| Token 机制 | 先获取 token → 请求带上 token → 服务端校验并删除 token | API 调用 | 通用性强 | 需要 token 管理 |
| 状态机 | 已支付的订单不再支付 | 订单、工单 | 业务语义清晰 | 需定义状态图 |
| 去重表 | 独立表记录 (request_id, result),查表判断 | 通用 | 可与业务表分离 | 额外存储开销 |
go
// 幂等支付示例:利用 DB 唯一约束
func ProcessPayment(txnId string, amount int) error {
// 1. 先插入幂等记录(唯一约束保证不会重复)
_, err := db.Exec(
"INSERT INTO idempotency_keys (txn_id, status, created_at) VALUES (?, 'pending', NOW())",
txnId,
)
if err != nil {
if isDuplicateKey(err) {
// 已处理过,读取已有结果返回
return getCachedResult(txnId)
}
return err
}
// 2. 执行实际业务
if err := deductBalance(amount); err != nil {
db.Exec("UPDATE idempotency_keys SET status='failed' WHERE txn_id=?", txnId)
return err
}
// 3. 标记成功
db.Exec("UPDATE idempotency_keys SET status='success' WHERE txn_id=?", txnId)
return nil
}10.6 降级(Degradation)
降级的核心理念是"宁可少给,不可不给"——当某个非核心依赖不可用时,系统应有能力返回部分结果而非整个请求失败。这和熔断器不同:熔断器是"发现故障 → 停止调用",降级是"发现故障 → 用备选方案继续服务"。好的降级策略让用户几乎感受不到后端故障:
mermaid
flowchart TB
Request["用户请求商品详情"] --> Core["核心: 商品基本信息<br/>必须返回"]
Request --> Reviews["依赖: 评价服务"]
Request --> Recommend["依赖: 推荐服务"]
Reviews -->|"正常"| ReviewData["显示评价"]
Reviews -->|"超时/故障"| ReviewDefault["降级: '评价加载中...'"]
Recommend -->|"正常"| RecData["显示推荐"]
Recommend -->|"超时/故障"| RecDefault["降级: 显示默认推荐/缓存"]| 降级策略 | 做法 | 适用 |
|---|---|---|
| 返回默认值 | 超时时返回预设数据 | 非核心展示 |
| 返回缓存 | 用上次成功的数据 | 相对稳定的数据 |
| 屏蔽功能 | 直接跳过该功能模块 | 非关键路径 |
| 静态化 | 降级到静态页面 | 大促/秒杀场景 |
| 简化处理 | 跳过复杂计算用简单逻辑 | 推荐算法等 |
10.7 限流(Rate Limiting)
限流是保护系统不被突发流量打垮的最后一道防线,算法原理与实现详见本章 #_9-限流算法。从微服务架构角度看,限流有两个部署位置:
| 部署位置 | 代表 | 适用场景 |
|---|---|---|
| 网关层集中限流 | Nginx limit_req, Kong rate-limiting | 全局入口流量控制 |
| 服务层独立限流 | Go x/time/rate, Java Guava RateLimiter | 每个服务独立保护自己 |
| SDK/边车限流 | Sentinel, Hystrix | 多语言微服务统一治理 |
生产实践:Go 用
golang.org/x/time/rate,Java 用 GuavaRateLimiter,网关层用 Nginxlimit_req_zone+limit_req。
三大支柱
| 支柱 | 作用 | 工具 |
|---|---|---|
| Metrics | 聚合数据:计数器、仪表、直方图 | Prometheus + Grafana |
| Logging | 离散事件记录 | ELK, Loki |
| Tracing | 请求链路追踪 | Jaeger, Zipkin |
黄金信号(回顾)
来自 [性能分析与调优]:Latency, Traffic, Errors, Saturation。
11. 秒杀系统设计
秒杀是高并发场景中最具挑战性的系统设计题之一。核心矛盾:瞬时海量流量 vs 有限库存。
11.1 核心挑战
| 挑战 | 描述 | 影响 |
|---|---|---|
| 高并发 | 瞬时 QPS 可达日常 100-1000 倍 | 服务器过载、连锁雪崩 |
| 超卖 | 热点库存扣减串行化 | 数据库行锁竞争、死锁 |
| 刷子/黄牛 | 脚本自动抢购 | 正常用户无法抢到 |
| 热点隔离 | 商品详情页成为热点 | 拖垮整个集群 |
| 削峰 | 瞬间流量冲垮系统 | DB 连接池耗尽 |
11.2 架构分层
mermaid
flowchart TB
User["用户"] --> CDN["CDN / 静态化"]
CDN --> GW["API 网关<br/>限流、鉴权、黑名单"]
GW --> Queue["消息队列<br/>削峰填谷"]
Queue --> App["秒杀服务<br/>库存扣减"]
App --> Cache["Redis<br/>库存预热 + Lua 原子扣减"]
Cache --> DB["MySQL<br/>异步落库"]11.3 逐层优化策略
第一层:前端限流
javascript
// 按钮防抖 + 灰置,防止用户重复点击
const btn = document.getElementById('seckill-btn');
let canClick = true;
btn.onclick = async () => {
if (!canClick) return;
canClick = false;
btn.disabled = true;
btn.innerText = '排队中...';
try {
await seckill();
} finally {
// 接口返回后不解禁 — 秒杀只能点一次
}
};text
前端策略:
- CDN 静态化商品详情页(不请求后端)
- 按钮置灰 + 防抖
- 倒计时结束后才露出抢购按钮
- 请求中添加验证码/滑块(防脚本)第二层:网关层
nginx
# Nginx 连接数限制 — 防 DDoS
limit_conn_zone $binary_remote_addr zone=addr:10m;
limit_conn addr 10; # 单 IP 最多 10 个并发连接
# Nginx 请求频率限制 — 令牌桶
limit_req_zone $binary_remote_addr zone=req:10m rate=100r/s;
limit_req zone=req burst=200 nodelay;text
网关层策略:
- IP 级别限流 + 黑名单
- 反爬虫 (UA 检查、验证码)
- 鉴权失败直接拒绝
- 秒杀链接动态生成(防提前曝光)第三层:削峰(消息队列)
go
// 秒杀请求 → MQ → 异步处理
func handleSeckill(userID, productID string) error {
// 1. Redis 预检:库存是否 > 0?用户是否已抢过?
// 2. 快速失败返回(告知用户排队中)
// 3. 投递到 MQ
return mq.Publish(SeckillMessage{
UserID: userID,
ProductID: productID,
Timestamp: time.Now(),
})
}第四层:库存扣减 — Redis + Lua 原子操作
这是秒杀系统最核心的部分。库存必须预热到 Redis,扣减必须是原子操作。
lua
-- seckill.lua: Redis Lua 原子扣减
-- KEYS[1]: 库存 key (seckill:stock:{productID})
-- KEYS[2]: 已购用户 set (seckill:purchased:{productID})
-- ARGV[1]: 用户 ID
local stock = tonumber(redis.call('GET', KEYS[1]))
if not stock or stock <= 0 then
return -1 -- 库存不足
end
-- 检查是否已购买过(防重复)
local purchased = redis.call('SISMEMBER', KEYS[2], ARGV[1])
if purchased == 1 then
return -2 -- 已购买
end
-- 扣减库存 + 记录用户
redis.call('DECR', KEYS[1])
redis.call('SADD', KEYS[2], ARGV[1])
return 1 -- 成功go
// Go 调用 Redis Lua
func seckill(ctx context.Context, rdb *redis.Client, productID, userID string) (int, error) {
script := `
local stock = tonumber(redis.call('GET', KEYS[1]))
if not stock or stock <= 0 then return -1 end
local purchased = redis.call('SISMEMBER', KEYS[2], ARGV[1])
if purchased == 1 then return -2 end
redis.call('DECR', KEYS[1])
redis.call('SADD', KEYS[2], ARGV[1])
return 1
`
return rdb.Eval(ctx, script,
[]string{
fmt.Sprintf("seckill:stock:%s", productID),
fmt.Sprintf("seckill:purchased:%s", productID),
},
userID,
).Int()
}第五层:异步落库
go
// MQ 消费者:将扣减成功的订单写入 MySQL
func processSeckillOrder(msg SeckillMessage) {
// 幂等性:用 orderID 防重
tx, _ := db.Begin()
// INSERT INTO orders (order_id, user_id, product_id) VALUES (?, ?, ?)
// UPDATE product SET stock = stock - 1 WHERE id = ? AND stock > 0
tx.Commit()
// 通知用户抢购结果(推送/短信)
notifyUser(msg.UserID, "恭喜抢到!")
}11.4 高并发陷阱与避坑指南
| 陷阱 | 错误做法 | 正确做法 |
|---|---|---|
| 超卖 | stock - 1 先读后写 | Redis Lua 原子操作 + DB WHERE stock > 0 |
| MySQL 行锁瓶颈 | 直接UPDATE商品表扣库存 | Redis 预热库存,异步写 DB |
| 缓存击穿 | 秒杀开始瞬间库存 key 过期 | 预热时不设 TTL 或设置足够长 |
| 消息积压 | 无限制接收请求 | 前端+网关限流,库存为 0 直接拒绝 |
| 链接暴露 | 固定秒杀 URL | 活动开始前动态生成加密链接 |
| 重复抢购 | 不加幂等校验 | Redis Set 已购用户 + 订单幂等 |
| 热点商品详情 | 每次请求都查 DB | CDN 静态化 + 本地缓存 |
11.5 秒杀系统时序
mermaid
sequenceDiagram
participant U as 用户
participant GW as 网关
participant R as Redis(库存)
participant MQ as 消息队列
participant DB as MySQL
Note over U: 10:00:00 秒杀开始
U->>GW: 点击抢购
GW->>GW: IP限流/鉴权/黑名单
GW->>R: Lua 脚本原子扣减
alt 库存不足
R-->>GW: -1 → 返回"已抢光"
GW-->>U: 很遗憾
else 已购买
R-->>GW: -2 → 返回"已参与"
GW-->>U: 您已抢过
else 成功
R-->>GW: 1 → 成功
GW->>MQ: 投递订单消息
GW-->>U: 排队中,请稍候
MQ->>DB: 异步落库
DB-->>MQ: 写入成功
Note over U: 推送通知:抢购成功!
end10. Gossip 协议
10.1 协议原理
Gossip(流言)协议模仿人类社会中流言的传播方式:每个节点周期性地随机选择几个邻居,交换所知道的信息。没有中心节点,具有极高的鲁棒性。
text
Gossip 协议的核心循环(每个节点独立执行):
1. 每秒随机选择 k 个邻居(fanout=3)
2. 交换"摘要"(自己知道的最新信息)
3. 更新本地状态,合并对方提供的新信息
4. 重复 → 最终所有节点达成一致
传染模型:
- 初始 1 个节点持有消息 M
- 第一轮: 随机发给 fanout 个节点 → fanout 个节点知道
- 第二轮: 多个节点继续转发 → 指数级扩散
- O(log N) 轮后: 所有节点都知道(概率上保证)10.2 Go 模拟实现
go
package main
import (
"fmt"
"math/rand"
"sync"
"time"
)
type Node struct {
id int
data map[string]int64 // key → version (版本号越大越新)
peers []*Node // 邻居列表
fanout int // 每轮选择几个节点交换
mu sync.RWMutex
}
func NewNode(id, fanout int) *Node {
return &Node{
id: id,
data: make(map[string]int64),
fanout: fanout,
}
}
// 与其他节点交换数据 (Push-Pull 模式)
func (n *Node) GossipRound() {
n.mu.RLock()
// 随机选择 fanout 个邻居
selected := randomSample(n.peers, n.fanout)
n.mu.RUnlock()
for _, peer := range selected {
n.exchange(peer)
}
}
func (n *Node) exchange(peer *Node) {
// Push: 发送自己的数据
n.mu.RLock()
mySnapshot := make(map[string]int64)
for k, v := range n.data {
mySnapshot[k] = v
}
n.mu.RUnlock()
// Pull: 获取对方的数据
peer.mu.RLock()
peerSnapshot := make(map[string]int64)
for k, v := range peer.data {
peerSnapshot[k] = v
}
peer.mu.RUnlock()
// 合并(取版本号更大的)
n.mu.Lock()
for k, v := range peerSnapshot {
if myV, ok := mySnapshot[k]; !ok || v > myV {
n.data[k] = v
}
}
n.mu.Unlock()
peer.mu.Lock()
for k, v := range mySnapshot {
if peerV, ok := peerSnapshot[k]; !ok || v > peerV {
peer.data[k] = v
}
}
peer.mu.Unlock()
}
// 引入新数据
func (n *Node) Introduce(key string) {
n.mu.Lock()
n.data[key] = time.Now().UnixNano()
n.mu.Unlock()
}
// 统计持有某 key 的节点数
func convergence(nodes []*Node, key string) int {
count := 0
for _, n := range nodes {
n.mu.RLock()
if _, ok := n.data[key]; ok {
count++
}
n.mu.RUnlock()
}
return count
}go
// 模拟 100 个节点的 Gossip 收敛
func main() {
const N = 100
nodes := make([]*Node, N)
for i := 0; i < N; i++ {
nodes[i] = NewNode(i, 3)
}
// 构建随机拓扑(每个节点连接所有其他节点)
for _, n := range nodes {
n.peers = nodes
}
// 节点 0 引入新数据
nodes[0].Introduce("key-001")
// 每秒一轮,观察收敛速度
for round := 0; round < 10; round++ {
var wg sync.WaitGroup
for _, n := range nodes {
wg.Add(1)
go func(n *Node) {
defer wg.Done()
n.GossipRound()
}(n)
}
wg.Wait()
c := convergence(nodes, "key-001")
fmt.Printf("Round %d: %d/%d nodes aware (%.0f%%)\n",
round+1, c, N, float64(c)/N*100)
}
}text
模拟输出(fanout=3):
Round 1: 3/100 nodes aware (3%) ← 指数扩展初期慢
Round 2: 8/100 nodes aware (8%)
Round 3: 22/100 nodes aware (22%)
Round 4: 53/100 nodes aware (53%)
Round 5: 85/100 nodes aware (85%)
Round 6: 98/100 nodes aware (98%) ← 6 轮接近全收敛
Round 7: 100/100 nodes aware (100%)
结论: fanout=3, 100 节点 → ~6 轮 (6 秒) 收敛10.3 Gossip 协议的收敛速度分析
mermaid
flowchart TB
subgraph Spread["流言传播模型 (fanout=3)"]
R1["Round 1: 1 → 3"]
R2["Round 2: 3 → 9"]
R3["Round 3: 9 → 27"]
R4["Round 4: 27 → 81"]
end
R1 --> R2 --> R3 --> R4
Note["收敛轮数 ≈ log_{fanout}(N)<br/>fanout=3, N=100 → log₃(100) ≈ 4.2 轮<br/>但前几轮可能有重复 → 实际 ~6-8 轮"]| fanout | 100 节点收敛轮数 | 1000 节点收敛轮数 | 网络开销/轮 |
|---|---|---|---|
| 2 | ~8-10 | ~12-14 | 低 |
| 3 | ~6-8 | ~9-11 | 中 |
| 5 | ~4-5 | ~6-8 | 高 |
| 10 | ~3-4 | ~4-5 | 很高 |
10.4 Gossip 协议的容错性
text
节点故障: 不影响传播
- 某个节点宕机 → 其他节点继续传播
- 消息通过不同路径最终到达所有存活节点
网络分区: 分区内独立传播
- 两个分区各自内部收敛
- 分区愈合后 → 首次 Gossip 交换 → 数据合并(版本号仲裁)
- 自动愈合,无需人工干预
消息丢失: 多轮重试补偿
- 单次 Gossip 可能丢包 → 但下轮会有其他节点传播相同信息
- 丢失概率 P → 连续 k 轮全部丢失的概率 = P^(k*fanout)
- fanout=3, P=0.1 → 连续 3 轮全丢 = 0.1^9 = 10^-9(基本不可能)10.5 Gossip 在工程中的应用
| 系统 | 用途 | 实现 |
|---|---|---|
| Cassandra | 集群成员管理、节点发现 | 每 1 秒 gossip 一次 |
| Consul | 服务发现、健康检查 | Serf (Gossip 变体) |
| Redis Cluster | 集群拓扑传播 | 类 Gossip 协议 |
| DynamoDB | 数据同步、故障检测 | 原始 Gossip 架构 |
| Blockchain | P2P 网络传播交易/区块 | Gossip 广播 |
参考
- CAP Theorem: Brewer's Conjecture and the Feasibility of Consistent, Available, Partition-Tolerant Web Services
- Raft: In Search of an Understandable Consensus Algorithm
- Gossip Protocols (Demers et al., 1987)
- Martin Kleppmann: Designing Data-Intensive Applications
- How to do distributed locking
- Google SRE Book
登录后即可发表评论 👇