生产者消费者模式
《生产者-消费者模式》终极完全版笔记(含背压、超时与生产级防护)
一、 一句话穿透本质
生产者只负责“受理业务”(丢进筐),不负责“办理业务”(不等结果);消费者只负责“办理业务”(从筐里拿),不关心是谁丢的。
两者通过“队列(缓冲区)”和“任务ID(漂流瓶)”进行彻底解耦,实现异步削峰填谷。
二、 底层运行原理(阻塞与唤醒机制)
这是整个模式能跑起来的“物理定律”。绝对不能靠“死循环”去轮询(那会把CPU烧干),必须依靠操作系统的**“等待-通知”机制**。
| 场景 | Java(多线程)的底层做法 | Node.js(单线程异步)的底层做法 | 本质共性 |
|---|---|---|---|
| ① 队列为空时 (消费者没活干) |
线程调用 wait(),操作系统挂起线程(移出CPU队列),CPU占用为0。 |
执行 await 一个 pending 的 Promise,主线程交还给事件循环,去处理其他HTTP请求。 |
都让出了CPU,绝不空转消耗资源。 |
| ② 有任务到达时 (生产者派活) |
调用 notify(),操作系统唤醒消费者线程。 |
执行存好的 resolve() 函数,Promise 状态变更,消费者被唤醒。 |
精准唤醒,没有浪费CPU去轮询。 |
| ③ 队列满了 (生产者太快) |
生产者线程调用 wait(),被挂起休眠。 |
执行 await new Promise,生产者函数被挂起,直到消费者腾出空位并唤醒它。 |
主动限流,让生产者慢下来。 |
核心结论:Java 是“让线程睡觉”,Node 是“让 Promise 挂起、让主线程去干别的活”。两者殊途同归,目的都是在不干活的时候,把宝贵的 CPU 资源让出来。
三、 背压机制(Backpressure)与极简源码
定义:当下游(消费者)处理速度跟不上上游(生产者)生产速度时,为防止内存溢出而采取的强制限流手段。
1. 带背压的异步队列源码(Node.js 实现)
1 | class AsyncQueue { |
2. 背压的行为表现(业务视角)
- 队列没满:任务“嗖”的一下入队,
await几乎瞬间返回(仅一个微任务周期),用户毫秒级收到“提交成功”。 - 队列满了:代码确实“停”在
await queue.push()这一行(函数挂起),但主线程并未卡死(因为挂起的是函数,主线程去处理其他请求了)。等消费者腾出空位,这个请求被唤醒,才继续返回“提交成功”。
四、 避坑指南:为什么 push 里必须 new Promise?
错误写法:
await this.spaceResolvers.push(resolve)Array.push方法返回的是数组的新长度(数字),await发现不是 Promise,会瞬间跳过,根本起不到挂起等待的作用,生产者会直接往下执行,逻辑直接乱套。正确写法:
await new Promise(resolve => this.spaceResolvers.push(resolve))new Promise创建了一个pending(悬空)状态的 Promise,await必须等它被resolve才能继续。我们把resolve存进数组,相当于把“唤醒按钮”交给了消费者。每个挂起的请求都拥有一个独立的 Promise 实例,互不干扰。
五、 生产级风险:大量挂起 Promise 会炸内存吗?(必加防护)
会! 纯内存队列如果只靠 maxSize 限制 items,但不限制 spaceResolvers(等待空位的生产者数量),瞬间涌入 10 万个请求,10 万个 Promise 对象和它们携带的 req/res 上下文依然会撑爆内存。
生产环境必须加上以下三道防线:
| 防护策略 | 实现方式(Node.js) | 效果 |
|---|---|---|
| 1. 最大等待数(拒绝策略) | push 方法开头判断:若 this.spaceResolvers.length > 1000,直接 throw Error('系统繁忙')。 |
快速失败。第 1001 个请求直接返回 503,绝不挂起等待,保住服务器命脉。 |
| 2. 超时控制(Timeout) | 利用 Promise.race 包裹等待逻辑,设置超时(如 5 秒)。超时则 reject,并清理掉数组里对应的 resolve(防止内存泄漏)。 |
防止死等。若消费者因故障无法腾出空位,等待的请求不会永远挂起,而是超时报错让用户重试。 |
| 3. 网关层限流(终极护盾) | 在 Nginx / 负载均衡层配置最大连接数(如 1000)。 | 物理隔离。超过连接数的请求根本进不到 Node 进程,直接在最外层被挡掉。 |
超时控制的代码片段参考:
1 | // 在 push 方法内部,替换单纯的 await new Promise |
六、 核心疑问破解:“我怎么知道任务做完了?”
生产者不直接返回结果,但任务状态被持久化在数据库里。获取结果有三种标准姿势:
| 方案 | 实现方式 | 适用场景 |
|---|---|---|
| ① 前端轮询 | 前端拿 taskId 每隔 2 秒调接口查数据库 status 字段。 |
CMS后台最常用,极其简单可靠。 |
| ② WebSocket 推送 | 消费者处理完,主动 emit 事件给前端。 |
实时聊天、金融行情,0延迟。 |
| ③ 第三方回调 | 消费者处理完,POST 请求第三方 callbackUrl。 |
跨系统/微服务间调用。 |
七、 三大实现模型对比(多线程 / Node / 分布式MQ)
无论底层是哪种技术栈,这个模式的思想完全一致,只是落地的“材料”不同:
| 对比维度 | Java 多线程阻塞 | Node.js 单线程异步 | 分布式 MQ(RabbitMQ/Kafka) |
|---|---|---|---|
| 缓冲区位置 | JVM 堆内存(如 ArrayBlockingQueue)。 |
Node 进程内存(JS数组)。 | 独立中间件(Broker)的磁盘/内存。 |
| 等待机制 | wait/notify 挂起线程(内核态切换)。 |
await 挂起函数(用户态切换)。 |
长轮询(Long Polling)或 TCP 连接挂起。 |
| 数据安全性 | 不安全。进程重启,数据全丢。 | 不安全。进程重启,数据全丢。 | 安全。消息持久化到磁盘,重启不丢。 |
| 防护机制 | offer 拒绝 + poll(timeout) 超时。 |
内存 maxSize + spaceResolvers 限制 + Promise.race。 |
prefetch 预取限制 + 消息 TTL + 死信队列。 |
八、 最终收益总结(实战价值)
- 削峰填谷:瞬间来 1000 个请求,生产者 1 毫秒塞进队列(或触发背压挂起)。消费者按自己的速度慢慢消化,服务器稳如泰山。
- 故障隔离:消费者崩溃了,生产者依然毫秒级返回“提交成功”。等消费者重启,从队列(或磁盘)捞起任务继续干,用户无感知。
- 逻辑解耦:想加新功能(如摘要+配图)?完全不用动 API 代码,只需新增一个消费者监听同一队列或处理中间结果,扩展极其灵活。
九、 终极心法(刻进脑子里)
“同步直连是‘人等货’(用户死等结果),异步队列是‘货等人’(结果存好了等用户来取)。
我们把‘等待’从用户的浏览器转移到了数据库的status字段里;把‘计算压力’从主线程转移到了后台独立的消费者进程中。
背压不是简单的‘满了就等’,而是‘满了先等一会儿,等太久或等太多就直接拒绝’。用‘快速失败’和‘超时释放’换取系统的‘整体存活’。 ”

