跳转到内容

Queue

Wippy 提供队列系统,用于异步消息处理,支持可配置的驱动和消费者。

flowchart LR
P[Publisher] --> D[Driver]
D --> Q[Queue]
Q --> C[Consumer]
C --> W[Worker Pool]
W --> F[Function]
  • Driver - 后端实现(memory、AMQP、SQS)
  • Queue - 绑定到驱动的逻辑队列
  • Consumer - 连接队列到处理器,带并发设置
  • Worker Pool - 并发消息处理器

多个队列可以共享一个驱动。多个消费者可以从同一队列处理消息。

Kind描述
queue.driver.memory内存队列驱动
queue.driver.amqpAMQP(RabbitMQ)驱动
queue.driver.sqsAWS SQS 驱动(也支持 LocalStack、ElasticMQ)
queue.queue带驱动引用的队列声明
queue.consumer处理消息的消费者

用于开发和单节点部署的进程内驱动。无外部依赖。

- name: memory_driver
kind: queue.driver.memory
lifecycle:
auto_start: true

用于 RabbitMQ 和 AMQP 0-9-1 兼容的 broker。

- name: amqp_driver
kind: queue.driver.amqp
url: "amqp://guest:guest@localhost:5672/"
vhost: "/"
connection_name: "wippy-service"
heartbeat: "10s"
connection_timeout: "30s"
reconnect_delay: "1s"
reconnect_max_delay: "30s"
default_message_ttl: "1h"
default_queue_expiry: "24h"
prefetch_count: 10
lifecycle:
auto_start: true
字段类型默认值描述
urlstringamqp://guest:guest@localhost:5672/Broker URL
vhoststring-虚拟主机覆盖
connection_namestring-在 broker UI 中显示的标识符
auth_mechanismstringPLAINPLAINEXTERNAL(mTLS)或 AMQPLAIN
heartbeatduration-Keep-alive 间隔
connection_timeoutduration-拨号超时
reconnect_delayduration1s初始重连退避
reconnect_max_delayduration30s最大重连退避
default_message_ttlduration-应用于已声明队列的默认消息 TTL
default_queue_ttlduration-应用于已声明队列的默认 TTL
default_queue_expiryduration-已声明队列的默认队列过期时间
prefetch_countint-Channel 级 prefetch 上限
frame_sizeint-AMQP frame 大小限制
channel_maxint-每连接最大 channel 数
tlsobject-TLS 设置(见下文)

TLS 块:

tls:
enabled: true
server_name: "rabbit.example.com"
cert_env: "AMQP_CLIENT_CERT"
key_env: "AMQP_CLIENT_KEY"
ca_env: "AMQP_CA_CERT"
insecure_skip_verify: false

内联 cert/key/ca 字段携带 PEM 内容;*_env 变体通过 env registry 解析。两种来源在每个字段上互斥。insecure_skip_verify 禁用证书验证(仅用于开发)。

用于 AWS SQS 和 SQS 兼容的 endpoint(LocalStack、ElasticMQ)。凭证、区域和其他 AWS SDK 设置来自共享的 config.aws 资源。

- name: aws_config
kind: config.aws
region: us-east-1
access_key_id_env: app:AWS_ACCESS_KEY_ID
secret_access_key_env: app:AWS_SECRET_ACCESS_KEY
- name: sqs_driver
kind: queue.driver.sqs
config: app:aws_config
endpoint: "http://localhost:9324"
message_retention_period: 345600
default_delay_seconds: 0
lifecycle:
auto_start: true
字段类型默认值描述
configRegistry ID必需提供区域和凭证的 config.aws 资源
endpointstring-自定义 endpoint URL(LocalStack、ElasticMQ);真实 AWS 时省略
message_retention_periodint345600(4天)队列级保留时间(秒)(60–1209600)
default_delay_secondsint0CreateQueue 时应用的默认投递延迟(0–900)
disable_message_checksum_validationboolfalse在发送/接收时禁用 SQS 消息校验和检查
use_fipsboolfalse使用 FIPS 兼容的 endpoint
use_dual_stackboolfalse使用 dual-stack(IPv4 + IPv6)endpoint

队列在首次使用时由驱动自动创建。在发布时使用 SQS 前缀的 header(sqs.*)来寻址 SQS 特定属性;像 correlation_idcontent_type 这样的中性键在可能的情况下会被翻译为 SQS 系统属性。

- name: tasks
kind: queue.queue
driver: app.queue:memory_driver
codec: json/plain
queue_name: "app_tasks"
driver_options:
memory:
max_length: 500
dead_letter:
queue: app.queue:tasks_dlq
max_attempts: 5
字段类型必需描述
driverRegistry ID队列驱动
codecstring消息正文的线格编码。默认为 json/plain(参见编解码器
queue_namestring外部队列名(默认为 entry 名)
driver_optionsobject按驱动 kind 索引的子配置
dead_letter.queueRegistry ID失败消息的队列 ID
dead_letter.max_attemptsint路由到 DLQ 之前的尝试次数

driver_options 下的键按驱动名称分组。驱动只读取自己的子配置——其他键处于休眠状态,这允许单个队列条目在需要时为多个驱动声明设置。

memory:

描述
max_length有界缓冲区大小(0 = 无界)

amqp:

描述
durable在 broker 重启后保留
auto_delete当最后一个消费者断开时删除
message_ttl每队列消息 TTL 覆盖
queue_expiry未使用队列的过期时间
max_length保留的最大消息数

codec 选择消息正文在交给 broker 之前如何序列化。它是一个 payload 格式字符串,默认为 json/plain

编解码器格式
json/plainJSON(默认)
application/msgpackMessagePack

AMQP 驱动在发布的消息上设置相应的 content-typeapplication/jsonapplication/msgpack)。未知的编解码器会在声明队列时失败,而不是在发布时。

- name: task_consumer
kind: queue.consumer
queue: app.queue:tasks
func: app.queue:task_handler
concurrency: 4
prefetch: 20
auto_ack: false
driver_options:
amqp:
consumer_tag: "worker-1"
exclusive: false
lifecycle:
auto_start: true
depends_on:
- app.queue:tasks
字段默认值描述
queue必需队列 registry ID
func必需处理函数 registry ID
concurrency1并行 worker 数量
prefetch10每个 worker 的缓冲区大小
auto_ackfalse为 true 时,runtime 不调用 broker ack;处理器成功/失败是唯一的 settle 信号
driver_options-按驱动的子配置(与队列结构相同)

amqp 消费者选项:

描述
exclusive单消费者队列访问
no_local拒绝在同一连接上发布的消息
no_wait订阅时不等待 broker 确认
consumer_tag此订阅的标识符
消费者遵循调用上下文,可以受安全策略约束。在 lifecycle 级别配置 actor 和策略。参见 Security

Worker 作为并发 goroutine 运行:

concurrency: 3, prefetch: 10
1. Driver delivers up to 10 messages to buffer
2. 3 workers pull from buffer concurrently
3. As workers finish, buffer refills
4. Backpressure when all workers busy and buffer full

消费者处理器以解码后的消息体作为第一个参数。使用 queue.message() 访问投递元数据(id、headers)。

local queue = require("queue")
local logger = require("logger")
local function main(body)
local msg = queue.message()
logger:info("processing", {
id = msg:id(),
correlation_id = msg:header("correlation_id")
})
local ok, err = process_task(body)
if err then
return false -- nack: redelivery or DLQ
end
return true -- ack: remove from queue
end
return { main = main }
- name: task_handler
kind: function.lua
source: file://task_handler.lua
method: main
modules:
- queue
- logger

Runtime 根据处理器返回值自动 settle:

处理结果动作
true 或非 false 返回Ack
falseNack(根据驱动重新投递或 dead-letter)
抛出错误Nack

仅在需要提前 settle 时显式调用 msg:ack()msg:nack()。Settlement 是单次的:先到达的调用获胜。

当队列上配置了 dead_letter 时,nack 超过 max_attempts 的消息会被路由到 DLQ,驱动会设置 x_dead_letter_reasonx_original_queue header。发布者不得设置任何 x_* header——这些保留给 DLQ 簿记使用。

从 Lua 代码发布:

local queue = require("queue")
queue.publish("app.queue:tasks", {
id = "task-123",
action = "process",
data = payload
})

参见 Queue 模块 了解完整 API。

消费者停止时:

  1. 停止接收新消息
  2. 取消 worker 上下文
  3. 等待正在处理的消息(有超时)
  4. 如果 worker 未及时完成则返回错误