3.2 消息并发:Actor、Channel 与背压
共享内存把问题表述为“谁可以访问这份状态”;消息模型换了一个角度:“谁拥有状态,其他执行体怎样提出请求”。这个转换能缩小共享范围,但不会自动消除排队、丢消息、重复处理和资源耗尽。
Actor:把状态和行为放在同一边界
一个 Actor 通常包含三部分:
- 只能由它自己直接修改的状态;
- 接收消息的邮箱;
- 根据消息决定下一步行为的处理逻辑。
调用方 ──Reserve(orderId)──> 库存 Actor
├─ 私有库存状态
├─ 去重记录
└─ 邮箱
调用方 <──Reserved / Rejected──┘如果某个 Actor 一次只处理一条消息,那么单个消息处理过程不需要和同一 Actor 的另一条消息竞争状态。它提供的是状态所有权和处理顺序,系统整体仍然存在并发。许多 Actor 可以并行运行,多个发送者的消息也可能以不同顺序到达。
不要把 Actor 误解成以下绝对规则:
- 本地消息不一定经过字节序列化;框架可能直接传递对象引用,因此消息最好保持不可变;
- 消息发送成功不一定等于目标已经处理;
- 顺序保证通常有明确范围。以 Erlang 为例,同一发送者发往同一接收者的信号保持发送顺序,但这不等于来自多个发送者的全局顺序;
- Actor 不等于无限邮箱。无界邮箱会把过载延迟转化为内存增长。
Actor 的工程重点:协议和故障
Actor 以一组消息协议作为接口:
sealed interface InventoryCommand {}
record Reserve(String orderId, int quantity) implements InventoryCommand {}
record Cancel(String orderId) implements InventoryCommand {}
record GetAvailability() implements InventoryCommand {}协议设计要明确:
- 消息是否幂等,重放会怎样;
- 请求与响应怎样关联;
- 处理失败后重试、跳过还是停止;
- 邮箱满时拒绝、阻塞、丢弃还是降级;
- Actor 重启后,状态从日志、快照还是外部存储恢复。
监督树是 Actor 生态常见的故障组织方式:父级观察子级失败,按策略重启或停止。它解决的是“谁负责处置失败”,不是让失败消失。若消息已经触发外部付款,再重启处理器,仍需幂等键或事务性发件箱避免重复副作用。
Channel:把通信本身变成一等对象
CSP 风格更强调独立执行过程通过 Channel 同步。发送者通常不必知道具体接收者,只需依赖 Channel 的元素类型和关闭协议。
type Job struct {
ID string
URL string
}
func worker(ctx context.Context, jobs <-chan Job, results chan<- error) {
for {
select {
case <-ctx.Done():
return
case job, ok := <-jobs:
if !ok {
return
}
results <- process(job)
}
}
}这个签名已经表达了三件事:
jobs <-chan Job只能接收;results chan<- error只能发送;context.Context承载取消和截止时间。
无缓冲 Channel 要等发送者和接收者同时就绪,通信同时形成同步点。有缓冲 Channel 允许生产与消费短暂错峰,但容量也定义了系统愿意吸收多少积压。
背压是稳定性协议
假设入口每秒产生 2,000 个任务,下游只能处理 1,200 个:
积压增长率 = 2,000 - 1,200 = 800 个/秒任何有限内存最终都会耗尽。扩大队列只是延后故障。系统必须在以下策略中做出可观测的选择:
- 阻塞生产者:把压力向上游传播;
- 拒绝新任务:快速失败,并给调用方明确重试信息;
- 丢弃或合并:只适合允许损失的遥测、刷新类任务;
- 扩容消费者:前提是瓶颈确实能横向扩展;
- 降级工作量:减少单个任务的成本。
队列容量应由可接受等待时间反推。例如下游稳定吞吐 1,200 个/秒,业务最多容忍 2 秒排队,初始容量估算不应远离 2,400,随后再结合突发流量和内存占用压测,而不是随手填一个十万。
Actor 与 Channel 不只是语法不同
| 维度 | Actor | CSP / Channel |
|---|---|---|
| 核心抽象 | 拥有状态和邮箱的实体 | 执行过程之间的通信通道 |
| 寻址 | 通常向某个 Actor 地址发送 | 通常向某个 Channel 发送 |
| 状态归属 | 状态封装在 Actor 内 | 由使用 Channel 的过程决定 |
| 协调方式 | 异步消息和行为切换 | 同步或缓冲通信、选择操作 |
| 常见强项 | 实体生命周期、监督、位置透明 | 流水线、扇入扇出、背压 |
| 常见风险 | 邮箱膨胀、协议演进、失败重放 | 泄漏 goroutine、关闭权不清、环路等待 |
现实系统可以组合二者:Actor 内部使用 Channel 调度工作,流水线阶段用 Actor 承载持久状态。选择模型时看状态所有权和故障边界,不要按语言阵营站队。
关闭协议决定系统能否收尾
生产者—消费者最常见的错误往往不在数据竞争,而在关闭权没有归属:
- 通常由发送方关闭 Channel,接收方不应猜测“不会再有消息”;
- 多发送方需要一个协调者,在全部发送方结束后关闭;
- 取消与正常完成是不同语义,不能只靠一个空值混在一起;
- 消费者退出后,仍在阻塞发送的生产者必须能收到取消,否则会泄漏。
这些规则应当出现在 API 和测试中,而不只存在于团队口头约定。
完成检查
为“图片处理流水线”画出:上传、病毒扫描、转码、入库四个阶段,然后标注:
- 每个阶段的并发上限;
- 队列容量及满载策略;
- 谁关闭队列;
- 某张图片失败时,是跳过、重试还是取消整条流水线;
- 如何防止同一任务重复入库。
如果图中只有箭头而没有容量、取消和失败路径,它还不是可运行的并发设计。