跳转到内容

队列消费者

队列消费者使用工作池处理队列中的消息。

flowchart LR
subgraph Consumer[消费者]
QD[队列驱动] --> DC[投递通道<br/>prefetch=10]
DC --> WP[工作池<br/>concurrency]
WP --> FH[函数处理器]
FH --> AN[Ack/Nack]
end
选项默认值最大值说明
queue必填-队列注册表 ID
func必填-处理函数注册表 ID
concurrency11000工作线程数
prefetch1010000消息缓冲区大小
auto_ackfalse-在运行处理器之前自动确认
driver_options{}-驱动特定的消费者选项
- name: order_consumer
kind: queue.consumer
queue: app:orders
func: app:process_order
concurrency: 5
prefetch: 20
lifecycle:
auto_start: true
depends_on:
- app:orders

处理函数接收消息体:

-- process_order.lua
local json = require("json")
local function handler(body)
local order = json.decode(body)
-- 处理订单
local result, err = process_order(order)
if err then
-- 返回错误触发 Nack(重新入队)
return nil, err
end
-- 成功触发 Ack
return result
end
return handler
- name: process_order
kind: function.lua
source: file://process_order.lua
modules:
- json
结果动作效果
成功Ack消息从队列移除
错误Nack消息重新入队(取决于驱动)
  • 工作线程作为并发 goroutine 运行
  • 每个工作线程一次处理一条消息
  • 消息从投递通道轮询分发
  • 预取缓冲区允许驱动提前投递
  • 当所有工作线程繁忙且缓冲区满时产生背压
concurrency: 3
prefetch: 10
流程:
1. 驱动向缓冲区投递最多 10 条消息
2. 3 个工作线程并发从缓冲区拉取
3. 工作线程完成后,缓冲区重新填充
4. 所有工作线程繁忙且缓冲区满时产生背压

停止时:

  1. 停止接收新投递
  2. 取消工作线程上下文
  3. 等待处理中的消息(带超时)
  4. 如果工作线程未完成则返回超时错误
# 队列驱动(开发/测试用内存驱动)
- name: queue_driver
kind: queue.driver.memory
lifecycle:
auto_start: true
# 队列定义
- name: orders
kind: queue.queue
driver: app:queue_driver
queue_name: orders # 覆盖名称(默认:入口名称)
codec: json # 负载编解码器(可选)
dead_letter: # 死信处理(可选)
queue: app:dlq
max_attempts: 5
driver_options:
memory:
max_length: 10000 # 内存驱动:有界队列大小
字段说明
queue_name覆盖队列名称(默认:入口 ID 名称)
codec负载编解码器名称
dead_letter.queue死信队列的注册表 ID
dead_letter.max_attempts路由到 DLQ 之前的最大投递尝试次数
driver_options按驱动名称键控的驱动特定设置
死信路由取决于驱动。AMQP 遵循代理级别的 DLX 配置;内存驱动不会将消息路由到 DLQ。

用于开发/测试的内置内存队列:

  • 类型:queue.driver.memory
  • 消息存储在内存中
  • Nack 将消息重新入队到队列尾部
  • 重启后不持久化