如何用 iii queue worker 以 durable:subscriber 订阅 topic 处理发布/订阅消息? 📅 发布时间:2026/9/14 17:14:55 👁 浏览次数: 如何用 iii queue worker 以 durable:subscriber 订阅 topic 处理发布/订阅消息【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii当你需要多个独立消费者可靠地收到同一批事件例如订单变更、缓存失效通知并且不允许丢消息时iii 的queueworker 提供了持久化的发布/订阅通道生产者调用iii::durable::publish把消息发到命名 topic消费者 worker 通过注册durable:subscriber类型的 trigger 订阅该 topic。引擎为每条消息逐次运行消费函数消费函数正常返回即确认ack抛错则回退nack消息会重试并最终进入死信队列DLQ。前提已安装 iii CLI本地有一个运行中的 engine 和一个 Compose daemon。完整操作路径基于仓库文档 Queues 和 Workers。准备运行环境engine、Compose daemon 与 queue workercompose::add由运行中的 Compose daemon 提供所以在使用任何命令前需要在两个终端分别保持 engine 和 daemon 运行文档建议放在项目根目录# terminal 1 iii --config config.yaml # terminal 2, from the directory that contains worker-compose.yaml iii compose --namespace dev --engine ws://127.0.0.1:49134如果项目还没有 Compose 文件先创建worker-compose.yaml内容为containers: {}。然后在第三个终端与 daemon 相同的项目目录把queueworker 加入项目iii trigger -n dev compose::add workerqueue-n是iii trigger的--namespace短形式用于把函数解析到正在运行的 Compose daemon 所在的devnamespace。engine 与 SDK 包可以存在不同的 patch 版本但应保持同一个小版本线例如都在0.11.x除非发布说明另有说明。创建消费 worker 并注册 durable:subscriber 触发器消费者就是一个普通的 iii worker 进程安装对应语言的 iii SDK、通过III_URL连接 engine、注册函数和 trigger。按 创建新 worker 的方式建好项目后入口代码注册消费函数并把它绑定到 topic。以 Node / TypeScript 为例文档示例import { registerWorker } from iii-sdk; const url process.env.III_URL; if (!url) throw new Error(III_URL must be set); const worker registerWorker(url, { workerName: email-worker, namespace: orders, }); // receives the data from each published message worker.registerFunction(email::send, async (msg: { to: string; subject: string }) { // do the work here; throw to nack and let the message retry return { sent: true }; }); worker.registerTrigger({ type: durable:subscriber, function_id: email::send, config: { topic: emails }, });config.topic决定订阅哪个 topic函数收到的是每条已发布消息中的data字段。Python 版本的关键部分等价import os from iii import register_worker, InitOptions worker register_worker( os.environ[III_URL], InitOptions(worker_nameemail-worker, namespaceorders), ) def send(msg: dict) - dict: # do the work here; raise to nack and let the message retry return {sent: True} worker.register_function(email::send, send) worker.register_trigger({ type: durable:subscriber, function_id: email::send, config: {topic: emails}, })要让 Compose 启动这个本地 worker还需要在 worker 目录放一个iii.worker.yaml清单字段说明见 Workersname: email-worker description: Consumes emails topic and sends mail. scripts: start: pnpm start添加 worker 并发布消息在第三个终端把 worker 加入 Compose 使其启动iii trigger -n dev compose::add worker./email-worker确认消费端运行后向 topic 发布一条消息。引擎会把data投递给每个订阅者email::send每条消息运行一次# publish a message to the emails topic iii trigger iii::durable::publish --json {topic:emails,data:{to:ab.com,subject:hi}}也可以打开 console 的Traces标签页观察消息从 publish 到email::send执行的完整链路。验证订阅与投递状态pub/sub topic 只有在有函数订阅它之后才会出现在检查命令里向没有订阅者的 topic 发布会不会注册该 topic也就无从检查。列出所有 topic文档示例输出iii trigger engine::queue::list_topics[{ name: emails, broker_type: builtin, subscriber_count: 1 }]broker_type: builtin表示这是内置 broker 上的 pub/sub topic配置过的命名队列则显示broker_type: function_queue。查看 topic 统计文档示例输出depth是等待消费的积压消息数dlq_depth是已进入死信队列的消息数。消费者跟得上时depth应为 0iii trigger engine::queue::topic_stats topicemails{ depth: 0, consumer_count: 1, dlq_depth: 0, config: null }看到 topic 已注册且depth: 0即说明订阅和投递正常工作。处理失败重试、死信队列与重放投递失败时按指数退避重试1 秒然后 2 秒最多 3 次尝试之后消息进入 DLQ。因此正常情况下 DLQ 相关命令返回为空只有消息真的失败后才有内容。要验证 DLQ 路径可以让email::send改为抛错TypeScript 示例worker.registerFunction(email::send, async () { throw new Error(forced failure); }); worker.registerTrigger({ type: durable:subscriber, function_id: email::send, config: { topic: emails }, });再次发布消息重试耗尽后指数退避下大约几秒它会落到 DLQ。列出有死信消息的 topic文档示例输出id、时间戳、大小每次运行都不同iii trigger engine::queue::dlq_topics[{ topic: emails, broker_type: builtin, message_count: 1 }]浏览死信消息内容可以看到失败原因文档示例iii trigger engine::queue::dlq_messages topicemails[ { id: 0b9c…, payload: { to: ab.com, subject: hi }, error: ErrorBody { code: \invocation_failed\, message: \forced failure\..., failed_at: 1718900000, retries: 3, size_bytes: 64 } ]修复代码后把整个 topic 的死信消息移回主队列重新处理iii trigger iii::queue::redrive topicemails{ queue: emails, redriven: 1 }也可以只处理单条消息用engine::queue::dlq_messages返回的id重放一条iii trigger iii::queue::redrive_message topicemails message_id0b9c…或者彻底丢弃一条消息——该操作会从 DLQ 永久删除它文档示例输出iii trigger iii::queue::discard_message topicemails message_id0b9c…{ queue: emails, message_id: 0b9c…, redriven: 1 }可选用 queue_config 调整单个订阅者的投递方式topic选择消费什么queue_config则用于针对某一个订阅者调整投递。例如严格逐条处理而不是并发处理时注册 trigger 时指定fifo队列文档示例Node 与 Python 写法相同只是代码风格不同worker.register_trigger({ type: durable:subscriber, function_id: email::send, config: {topic: emails, queue_config: {type: fifo}}, })完整的 trigger 字段过滤、adapter 相关选项等以 queue worker 参考文档为准。限制与升级注意事项每个订阅者的持久队列以其 worker 的 namespace 为作用域同一 topic、同一 function id 的两个订阅者若处于不同 namespace是两条独立队列各自收到每条已发布事件而不是竞争消费同一条队列。使用 RabbitMQ adapter 的项目注意0.23.x 中队列命名已改为 namespace 限定格式旧队列不会自动迁移升级前需要排空并删除旧队列。详见 Upgrading from 0.22.x。0.23 起队列实现不再由 engine 提供而是独立的queueCompose 包从 0.22.x 升级时需要保留原有的file_path或 broker 配置避免旧任务丢失。下一步可以阅读 linkly 教程第 4 章其中演示了在真实项目里把数据库写入移到队列、并用 durable pub/sub 广播link.updated事件给缓存刷新器和 Python 分析 worker 的完整用法。【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考