
如何用 iii queue worker 通过 TriggerAction.Enqueue 把函数入队并处理重试与死信【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii当一次函数调用耗时长、或需要重试保证时才执行比如发邮件、跑批处理同步等待会拖住调用方。iii 的queueworker 提供的TriggerAction.Enqueue动作可以把这次调用放进一个命名的队列调用在消息入队时立即返回一个messageReceiptId由队列在后台按配置的重试次数投递目标函数连续失败的消息最终进入死信队列DLQ可以用内置的引擎函数检查、重投redrive或丢弃。适用前提你已经在用 iii engine 和 Compose daemon 管理一组 worker本文用 Node / TypeScript SDK 举例命令对所有 SDK 通用。启动 engine 与 Compose daemoncompose::add由运行中的 Compose daemon 处理需要 engine 和一个 daemon 各占一个终端。如果项目目录里还没有 Compose 文件先创建worker-compose.yaml内容只写containers: {}# 终端 1 iii --config config.yaml # 终端 2在包含 worker-compose.yaml 的目录下 iii compose --namespace dev --engine ws://127.0.0.1:49134之后所有iii trigger命令都在终端 3、同一个项目目录里执行。添加 queue worker在终端 3 添加独立queueworkeriii trigger -n dev compose::add workerqueue队列能力现在由这个独立 worker 承载不再内建在 engine 里engine 会把TriggerAction.Enqueue和durable:subscriber路由到queueworker。如果入队时报enqueue_error: … engine::queue::enqueue not found就是漏了这一步补上上面的compose::add即可来源queue worker 说明。定义命名队列email-jobsTriggerAction.Enqueue要求队列名已配置。编辑worker-compose.yaml中queue容器的条目通过config_override声明队列queue_configs下每个 key 就是一个队列名带独立的重试、并发和顺序设置containers: queue: worker: package://api.workers.iii.dev/queue config_name: queue config_override: queue_configs: email-jobs: max_retries: 3 # 投递失败的最大次数之后进 DLQ默认 3 concurrency: 5 # 同时处理的最大任务数默认 10 type: standard # standard并发或 fifo组内有序 adapter: name: builtin config: store_method: file_based # in_memory | file_based file_path: ./data/queue_store字段含义与默认值以 queue worker README 的字段表为准max_retries默认3concurrency默认10fifo队列会把它覆盖为prefetch1backoff_ms默认1000重试间隔按backoff_ms × 2^(attempt−1)指数退避即默认 1 秒、2 秒type: fifo时message_group_field必填且该 JSON 字段必须出现在每条 payload 里。改完后重启 queue 容器让配置生效——compose::restart会在重启前重新读取 Compose 文件只影响queue这一个容器iii trigger -n dev compose::restart workerqueue注册目标函数并入队入队模式不需要注册任何 trigger被入队的函数本身就是投递目标。先在某个 worker 里注册要异步执行的函数这里是email::send例如import { registerWorker } from iii-sdk; const url process.env.III_URL; if (!url) throw new Error(III_URL must be be set); const worker registerWorker(url, { workerName: email-worker, namespace: orders, }); worker.registerFunction(email::send, async (msg: { to: string; subject: string }) { // 在这里做实际工作抛异常即视为本次投递失败消息会被重试 return { sent: true }; });然后把这个 worker 加入 Compose 并启动worker./email-worker指向你的 worker 源码目录替换为你的实际路径iii trigger -n dev compose::add worker./email-worker在任意 worker 里用worker.trigger入队。和普通的worker.trigger调用唯一的区别是action字段传TriggerAction.Enqueuequeue指定队列名import { TriggerAction, type EnqueueResult } from iii-sdk; const { messageReceiptId } await worker.triggerunknown, EnqueueResult({ function_id: email::send, payload: { to: ab.com, subject: hi }, action: TriggerAction.Enqueue({ queue: email-jobs }), }); // messageReceiptId identifies the enqueued jobPython 写法等价TriggerAction.Enqueue(queueemail-jobs)返回值里取receipt[messageReceiptId]见 Functions 文档。注意调用语义传了Enqueue后调用在消息入队时就返回而不是等函数执行完不带action的默认调用则会同步等函数返回或超时。验证入队与投递已配置的命名队列会出现在引擎的队列列表里broker_type为function_queue对比pub/sub 主题是builtiniii trigger engine::queue::list_topics文档示例输出具体条目随你的队列名不同[ { name: email-jobs, broker_type: function_queue } ]再看队列积压与 DLQ 深度——depth是等待消费的条数dlq_depth是已死信的条数消费者跟得上时depth为0对命名队列consumer_count报告活跃的投递槽位数iii trigger engine::queue::topic_stats topicemail-jobs文档示例数值随运行状态变化{ depth: 0, consumer_count: 1, dlq_depth: 0, config: null }也可以打开 Console 的Traces标签页直接观察一条消息从入队到email::send执行的过程见 Queues 文档。观察重试与死信失败投递按指数退避重试默认backoff_ms1000时依次 1 秒、2 秒最多 3 次max_retries默认值仍失败则消息进入 DLQ。想亲眼看到 DLQ 被填充把email::send的 handler 临时改成抛错再入队一条消息几秒内三次尝试失败后消息就会死信worker.registerFunction(email::send, async () { throw new Error(forced failure); });列出有死信消息的队列/主题iii trigger engine::queue::dlq_topics文档示例[{ topic: emails, broker_type: builtin, message_count: 1 }]在什么都没失败之前这类查询返回空——DLQ 函数只在有消息真正耗尽重试后才有内容。查看单条死信消息能拿到错误原因和已重试次数下面的id、时间戳、字节数每次运行都不同属于文档示例iii trigger engine::queue::dlq_messages topicemail-jobs[ { id: 0b9c…, payload: { to: ab.com, subject: hi }, error: ErrorBody { code: \invocation_failed\, message: \forced failure\..., failed_at: 1718900000, retries: 3, size_bytes: 64 } ]重投或丢弃死信消息先把手改回去恢复正常的email::send实现并重启该 worker再把整个队列的死信消息搬回主队列重新处理iii trigger iii::queue::redrive topicemail-jobs文档示例返回{ queue: email-jobs, redriven: 1 }只处理单条消息时把engine::queue::dlq_messages返回的id传进去下面的message_id请替换为查询到的真实 id源文档中即为省略写法0b9c…iii trigger iii::queue::redrive_message topicemail-jobs message_iddlq_messages 返回的 id如果这条消息不值得再试可以直接丢弃——该操作会从 DLQ 永久删除这条消息执行前确认 id 没有拿错iii trigger iii::queue::discard_message topicemail-jobs message_iddlq_messages 返回的 id处理完可以再查一次engine::queue::dlq_messages确认 DLQ 已空。一个需要留意的字段名差异Queues 文档 里 redrive 系列命令用topic传队列/主题名而 queue worker README 对iii::queue::redrive的字段表写的是queue示例为--payload{queue: payment}。两处文档并存按你实际使用的版本以--help输出为准。限制与边界builtinadapter 无外部依赖适合单实例部署它支持重试、DLQ 和 FIFO。redisadapter 只做发布、不实现命名队列的消费、重试和 DLQ多实例生产环境文档推荐rabbitmqadapter字段见 queue worker README。fifo队列一次只处理一个任务prefetch1且重试是阻塞式的——会卡住队列直到该消息成功或死信用于支付、账本这类要求组内严格顺序的场景而不是高吞吐场景。TriggerAction.Enqueue只保证入队成功函数最终执行结果不会回到调用方调用方拿到的只有messageReceiptId。入队调用的路由仍然受 namespace 约束函数 id 在每个 namespace 内唯一跨 namespace 调用要在 trigger 上显式指定namespace否则在调用方所在 namespace 解析找不到时返回function_not_found并在错误里列出该 id 实际存在的 namespace见 Functions 文档。下一步可以结合 Queues 文档 了解 pub/sub 形态的持久队列durable:subscriberiii::durable::publish那是多消费者扇出且要求消息持久的另一种用法。【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考