拓冰建站拓冰建站
首页 / 资讯中心 / 正文

iii 队列 Worker 实战指南:命名队列、Pub/Sub 持久投递与死信队列(DLQ)全链路操作

iii 队列 Worker 实战指南命名队列、Pub/Sub 持久投递与死信队列DLQ全链路操作【免费下载链接】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 的两种队列形态——命名队列Named Queues与 Pub/Sub 队列durable:subscriber的完整实操流程从 Compose 环境准备、跨语言Node/Python/Rust入队与订阅代码到 topic 检查、消息重试、死信队列DLQ浏览与重投递redrive并结合引擎源码engine/src/workers/queue/印证重试次数、退避策略与命名空间隔离的底层实现读完后可独立在 iii 项目中搭建可重试、可观察、可恢复的异步消息链路。一、开始之前先让引擎与 Compose 守护进程跑起来queueworker 的核心价值在于解耦生产者与消费者一个函数把消息发布到指定主题后立即返回而订阅了该主题的函数在后台处理消息并且自带重试与死信队列。文档明确要求在执行任何后续命令前先保证两个终端分别在运行引擎与 Compose 守护进程# 终端 1启动引擎 iii --config config.yaml # 终端 2在包含 worker-compose.yaml 的目录启动 Compose 守护进程 iii compose --namespace dev --engine ws://127.0.0.1:49134几点说明如果项目还没有 Compose 文件先创建内容为containers: {}的worker-compose.yamliii compose是一个刻意设计的“无动词”守护进程调用方式详见 CLI reference下文出现的-n是iii trigger的--namespace短选项从源码结构看durable:subscriber触发类型的提供方provider在引擎中被固定登记为queueworker见 trigger.rs 中的KNOWN_TRIGGER_TYPE_PROVIDERS表——即队列能力由独立的queueworker 提供而非引擎内建。若未添加该 worker注册durable:subscriber触发器时会收到指向 Compose worker 的“pending”告警engine/src/trigger.rs 的单测验证了这一映射关系。在第三个终端同一项目目录中把queueworker 加入 Composeiii trigger -n dev compose::add workerqueueworker 的来源解析与完整配置参考可查阅仓库内 engine/src/workers/queue/README.md该文件同时声明了一个重要事实队列能力已不再内建于引擎旧的iii-queueworker 不再注册引擎通过独立queueworker 路由TriggerAction.Enqueue与durable:subscriber。二、命名队列Named Queues用 TriggerAction.Enqueue 把函数调用排队当函数执行耗时较长、或需要保证固定的重试次数时可以使用TriggerAction.Enqueue把这次调用放入一个命名队列。入队函数本身不执行只是登记一个待执行的任务并返回。创建命名队列按 queue worker 配置参考 的方式创建一个叫email-jobs的命名队列在iii-config.yaml/compose 的config_override中声明queue_configs之后在入队时使用该队列名。配置字段、默认值与 FIFO 选项以 worker reference 为准仓库内 engine/src/workers/queue/config.rs 可以逐项印证这些默认值字段类型说明源码默认值max_retriesu32进入 DLQ 前的最大投递尝试次数默认3config.rs#L28-L30 的default_max_retries()concurrencyu32同时处理的作业数默认10fifo队列会被改写为prefetch1typestringstandard并发或fifo消息组内有序message_group_fieldstringfifo必填payload 中决定排序分组的 JSON 字段backoff_msu64指数退避基准公式backoff_ms × 2^(attempt−1)默认1000config.rs#L40-L42poll_interval_msu64轮询间隔默认100一份典型的 compose 配置摘自 engine/src/workers/queue/README.mdcontainers: queue: worker: package://api.workers.iii.dev/queue version: 0.21.5 config_name: queue config_override: queue_configs: default: max_retries: 5 concurrency: 5 type: standard payment: max_retries: 10 concurrency: 2 type: fifo message_group_field: transaction_id adapter: name: builtin config: store_method: file_based file_path: ./data/queue_store入队函数Enqueue入队函数与普通的worker.trigger调用注册方式完全相同唯一的区别是触发器携带一个TriggerAction.EnqueueactionNode / TypeScriptimport { 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 }), // queue 指定队列名 }); // messageReceiptId 标识这个被入队的作业Pythonfrom iii import TriggerAction receipt worker.trigger({ function_id: email::send, payload: {to: ab.com, subject: hi}, action: TriggerAction.Enqueue(queueemail-jobs), # queue 指定队列名 }) # receipt[messageReceiptId] 标识这个被入队的作业Rustuse iii_sdk::TriggerAction; use iii_sdk::protocol::TriggerRequest; use serde_json::json; let receipt worker .trigger(TriggerRequest { function_id: email::send.to_string(), payload: json!({ to: ab.com, subject: hi }), action: Some(TriggerAction::Enqueue { queue: email-jobs.to_string() }), // queue 指定队列名 timeout_ms: None, }) .await?; // receipt[messageReceiptId] 标识这个被入队的作业命名队列与 topic 队列的选型对比源自 worker reference见 engine/src/workers/queue/README.md维度Topic-basedPub/SubNamed queue生产者trigger({ function_id: iii::durable::publish, payload: { topic, data } })trigger({ function_id, payload, action: TriggerAction.Enqueue({ queue }) })消费者registerTrigger({ type: durable:subscriber, config: { topic } })无需注册函数即目标投递扇出每个订阅函数都收到消息副本间竞争单次入队指向单个目标函数配置触发器上可选queue_config配置文件的queue_configs适用场景带重试与扇出的持久化 Pub/Sub带重试、FIFO、DLQ 的函数直调三、Pub/Sub 队列durable:subscriber 消费与 iii::durable::publish 发布当多个监听者需要订阅同一份数据、且消息必须持久可靠不丢时队列可以退化为发布/订阅形态消费者按 topic 订阅发布者把data投给每个订阅者。注册消费者消费者通过注册一个durable:subscriber类型的 Trigger 来绑定消息。引擎对每条消息执行一次目标函数并把发布的data作为 payload 传入正常返回即 ack抛出异常即 nack消息随后进入重试、最终进死信队列。Node / TypeScript完整 worker 示例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, }); // 收到每条发布消息的 data worker.registerFunction(email::send, async (msg: { to: string; subject: string }) { // 在这里干活抛异常即 nack让消息重试 return { sent: true }; }); worker.registerTrigger({ type: durable:subscriber, function_id: email::send, config: { topic: emails }, });Pythonimport os from iii import register_worker, InitOptions worker register_worker( os.environ[III_URL], InitOptions(worker_nameemail-worker, namespaceorders), ) # 收到每条发布消息的 data def send(msg: dict) - dict: # 在这里干活raise 即 nack让消息重试 return {sent: True} worker.register_function(email::send, send) worker.register_trigger({ type: durable:subscriber, function_id: email::send, config: {topic: emails}, })Rustuse iii_sdk::{InitOptions, RegisterFunction, register_worker}; use iii_sdk::protocol::RegisterTriggerInput; use serde::Deserialize; use schemars::JsonSchema; use serde_json::json; #[derive(Deserialize, JsonSchema)] struct Email { to: String, subject: String, } let url std::env::var(III_URL).expect(III_URL must be set); let worker register_worker( url, InitOptions { namespace: Some(orders.into()), ..Default::default() }, ); // 收到每条发布消息的 data worker.register_function(email::send, RegisterFunction::new(|_msg: Email| { // 在这里干活返回 error 即 nack让消息重试 Ok(json!({ sent: true })) })); worker.register_trigger(RegisterTriggerInput { trigger_type: durable:subscriber.into(), function_id: email::send.into(), config: json!({ topic: emails }), metadata: None, })?;把 worker 加入 Compose 并启动iii trigger -n dev compose::add worker./email-worker发布消息消费者运行后向它的 topic 发布消息。引擎把data投递给每个订阅者因此email::send对每条消息各执行一次# 向 emails topic 发布一条消息 iii trigger iii::durable::publish --json {topic:emails,data:{to:ab.com,subject:hi}}观察技巧打开 Console 并切到Traces标签页可以追踪消息从 publish 一路流转到email::send执行的全过程。从源码看iii::durable::publish的入参契约是topic必填与data投递给每个订阅函数的任意 payload见 engine/src/workers/queue/README.md 的 “Builtin Functions” 一节。重试与投递语义topic 选择消费什么queue_config 调节怎么消费示例中topic决定消费什么而queue_config用于调节单个订阅者的投递行为。以“严格逐条串行处理”为例把触发器注册为fifo队列Node / TypeScriptworker.registerTrigger({ type: durable:subscriber, function_id: email::send, config: { topic: emails, queue_config: { type: fifo } }, });Pythonworker.register_trigger({ type: durable:subscriber, function_id: email::send, config: {topic: emails, queue_config: {type: fifo}}, })Rustworker.register_trigger(RegisterTriggerInput { trigger_type: durable:subscriber.into(), function_id: email::send.into(), config: json!({ topic: emails, queue_config: { type: fifo } }), metadata: None, })?;queue_config支持哪些字段可以直接从引擎侧的解析代码确认。subscriber_config.rs 的SubscriberQueueConfig::from_value解析出这些键type映射为队列模式fifo→QueueMode::Fifo其余一律Concurrent、maxRetries、concurrency、visibilityTimeout、delaySeconds、backoffType、backoffDelayMs以及 RabbitMQ 专用的maxPriority1–255 优先队列级别。重试语义上有一个容易踩坑的细节在 adapters/builtin/adapter.rs 的to_subscription_config中max_attempts maxRetries 1——即maxRetries: 3意味着首次尝试加 3 次重试共 4 次投递内置单测同文件 L780-L798也显式断言了这一点。默认配置下max_retries3、backoff_ms1000失败的投递会以指数退避1 秒、2 秒重试至多 3 次随后消息进入死信队列。另一个关键语义是命名空间隔离每个订阅者的持久队列按订阅 worker 的 namespace 隔离同一 topic function id 在不同 namespace 下的两个订阅者是两条各自收到全部消息的队列而不是竞争同一条队列的两个消费者。builtin 适配器在启动时还会把旧版无命名空间前缀的订阅队列迁移到默认命名空间adapters/builtin/adapter.rs#L238-L253 的migrate_legacy_subscriber_queues。注意RabbitMQ 适配器的队列命名在 0.23.x 中发生过变更且不会自动迁移跨版本升级时请按仓库 docs/next/upgrading/ 中的升级说明重新声明 durable subscriber 队列。四、检查队列 Topiclist_topics 与 topic_stats以下命令对两种队列都适用。区分标志是返回中的broker_typePub/Sub topic 只有在有函数订阅它之后才会出现向无人订阅的 topic 发布不会注册它也就没有可检查的内容其broker_type为builtin配置好的命名队列则以broker_type: function_queue出现。列出全部 topic上文的emails订阅后应可见iii trigger engine::queue::list_topics[{ name: emails, broker_type: builtin, subscriber_count: 1 }]查询某个 topic 的统计depth为等待消费者处理的积压消息数dlq_depth为已进入死信队列的消息数消费者跟得上时depth保持 0。对命名队列而言consumer_count报告的是它的活跃投递槽位iii trigger engine::queue::topic_stats topicemails{ depth: 0, consumer_count: 1, dlq_depth: 0, config: null }五、死信队列让失败可见并可重投递或丢弃消息只有在订阅函数耗尽全部重试次数后才会进入 DLQ因此在出现失败之前DLQ 相关函数返回空列表。以下流程完整演示“制造失败 → 消息进 DLQ → 检查 → 重投递/丢弃”。1. 让消息进入死信队列把email::send的 handler 改成直接抛错每条消息都会耗尽 3 次重试按指数退避约几秒后进入 DLQ。Node / TypeScriptworker.registerFunction(email::send, async () { throw new Error(forced failure); }); worker.registerTrigger({ type: durable:subscriber, function_id: email::send, config: { topic: emails }, });Pythondef send(_msg: dict) - dict: raise Exception(forced failure) worker.register_function(email::send, send) worker.register_trigger({ type: durable:subscriber, function_id: email::send, config: {topic: emails}, })Rustworker.register_function(email::send, RegisterFunction::new(|_msg: serde_json::Value| { Err::serde_json::Value, _(iii_sdk::Error::Handler(forced failure.into())) })); worker.register_trigger(RegisterTriggerInput { trigger_type: durable:subscriber.into(), function_id: email::send.into(), config: json!({ topic: emails }), metadata: None, })?;然后发布一条消息重试耗尽后它会落入 DLQ具体 id、时间戳、大小每次运行会有差异iii trigger iii::durable::publish --json {topic:emails,data:{to:ab.com,subject:hi}}从源码看nack 之后 builtin 队列在内部根据作业自身的attempts_made与max_attempts决定是重试还是进 DLQadapters/builtin/adapter.rs#L622-L636 的nack_function_queue→BuiltinQueue::nack()。2. 列出存在死信消息的 topiciii trigger engine::queue::dlq_topics[{ topic: emails, broker_type: builtin, message_count: 1 }]3. 浏览死信消息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 } ]4. 整 topic 重投递redrive把代码修回正常逻辑后将该 topic 的 DLQ 消息整批移回主队列重新处理iii trigger iii::queue::redrive topicemails{ queue: emails, redriven: 1 }修复后的函数会重新处理这些消息DLQ 随之清空可以再跑一次dlq_messages验证iii trigger engine::queue::dlq_messages topicemails引擎侧对应实现在 queue.rs#L293-L318 的iii::queue::redrive函数输入RedriveInput同时接受queue与topic两种键名queue.rs#L107-L109 中#[serde(alias topic)]因此上面文档写法与 worker reference 中的{queue: ...}写法等价。5. 单条消息重投递或丢弃只处理一条而不是整个 topic 时把engine::queue::dlq_messages返回的id传给单条操作。重投递回主队列iii trigger iii::queue::redrive_message topicemails message_id0b9c…{ queue: emails, message_id: 0b9c…, redriven: 1 }或者直接丢弃从 DLQ 永久删除iii trigger iii::queue::discard_message topicemails message_id0b9c…{ queue: emails, message_id: 0b9c…, redriven: 1 }相关端到端验证可以参考仓库测试 engine/tests/dlq_redrive_e2e.rsredrive 全流程与 engine/tests/queue_e2e_happy_path.rs、engine/tests/queue_e2e_fanout.rs持久订阅正常路径与扇出行为。六、适配器选择builtin / redis / rabbitmq队列的底层传输由 adapter 决定默认builtin。对照 engine/src/workers/queue/README.md 的适配器对比能力builtinrabbitmqredis重试支持支持不支持死信队列支持支持不支持FIFO 排序支持支持不支持命名队列消费支持支持不支持仅发布Topic 发布/订阅支持支持支持多实例不支持支持支持外部依赖无RabbitMQRedis官方给出的场景建议本地开发用builtinin_memory单实例生产用builtinfile_based落盘到如./data/queue_store多实例生产用rabbitmq。相应配置# redis name: redis config: redis_url: ${REDIS_URL:redis://localhost:6379} # rabbitmq name: rabbitmq config: amqp_url: ${RABBITMQ_URL:amqp://localhost:5672}总结与关键要点两种队列形态各司其职命名队列面向“把某次函数调用排队执行”TriggerAction.Enqueuequeue_configs配置Pub/Sub 队列面向“多消费者、不丢消息的持久事件流”durable:subscriberiii::durable::publishack/nack 语义是全部可靠性机制的基石正常返回即确认抛错即拒绝拒绝触发指数退避重试默认 3 次重试、基准退避 1 秒耗尽后进 DLQ命名空间即队列边界同一 topic 的不同 namespace 订阅者各自持有完整消息流扩容副本之间才是竞争消费关系运维闭环齐全engine::queue::list_topics/topic_stats看积压dlq_topics/dlq_messages查死信iii::queue::redrive/redrive_message/discard_message完成重投递或丢弃且 redrive 类函数对queue、topic两种键名均兼容选型看部署规模单机零依赖选 builtin多实例持久化选 rabbitmqredis 适配器仅适合 topic 发布场景。【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门