news 2026/9/14 17:49:48

iii 实战教程:用队列与发布订阅让短链服务具备持久化执行能力(Durable Execution)

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
iii 实战教程:用队列与发布订阅让短链服务具备持久化执行能力(Durable Execution)

iii 实战教程:用队列与发布订阅让短链服务具备持久化执行能力(Durable Execution)

【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii

导读

本文是 iii 项目 Linkly 短链服务系列教程的第 4 章("Make it durable")。前三章中,每次重定向都会直接触发link::record_click写数据库,导致数据库写入落在重定向的热路径上——一次慢写入就会拖慢整个重定向响应。本章将这一工作迁移到队列上,让重定向立即返回;随后用**发布/订阅(pub/sub)**把链接事件广播给互相解耦的独立订阅者:一个 Python 分析 worker 和一个缓存刷新器。读完本文,你将掌握TriggerAction.Enqueue的用法、queueworker 的队列配置、pubsubiii::durable::publish两种发布订阅模式的取舍,以及如何用 Python 编写一个跨语言、跨 worker 的事件消费者。

原文出处:docs/tutorials/linkly/durable-execution.mdx(本教程的.skill.md渲染版本为 durable-execution.mdx.skill.md)。

本章目标:把写库工作移出热路径

回顾上一章(Ch. 3)的实现:http::redirect直接触发link::record_click,也就是每次用户访问短链时,重定向函数都要同步等待一次数据库写入完成。这带来两个问题:

  1. 延迟耦合:数据库写入慢时,重定向也跟着变慢,用户体验直接受损;
  2. 责任耦合:重定向这个"读路径"上背负了"写日志"的额外职责。

本章的解法是用两个既有的 worker 能力重构:

  • 队列(queue):接受"现在提交、稍后执行"的工作,让link::record_click在后台排空(drain)队列,重定向立刻返回 302;
  • 发布订阅(pub/sub):队列把每条消息投递给一个消费者,而当系统里多个互不相关的部分需要对同一事件做出反应时,改用发布/订阅设计,把事件广播给所有订阅者。

最终得到两个完全解耦的消费者:一个Python analytics worker(按天统计短链创建数)和一个缓存刷新器link::on_link_updated),linkworker 甚至不需要知道它们的存在。

添加 worker:加入 pubsub

本章涉及两个 worker。queueiii project init初始化项目时就已经内置,无需重复添加;pubsub则需要显式加入。在项目目录下执行:

iii worker add pubsub

这条命令会把pubsubworker 注册到当前项目。注意:queue自第 1 章起就在运行,它的设置统一管理在./config/queue.yaml中(详见 Configuration)。

用队列让重定向变快

队列的基本概念

队列保存"现在接受、以后运行"的工作。本章定义了一个名为clicks的队列,专门用于link::record_click

定义 clicks 队列:编辑 config/queue.yaml

./config/queue.yamlvalue:下添加queue_configs并保存。配置热生效,无需重启

id: queue # ... value: queue_configs: clicks: type: standard max_retries: 5 concurrency: 5 adapter: name: builtin

各字段说明(结合 engine/src/workers/queue/README.md 中的配置表):

字段类型说明默认值
typestringstandard(并发处理)或fifo(按消息组有序处理)
max_retriesu32进入死信队列(DLQ)前的最大投递尝试次数3
concurrencyu32同时处理的最大任务数;FIFO 队列会强制为prefetch=110
message_group_fieldstringfifo队列需要,指定决定顺序组的 JSON 字段
backoff_msu64重试基础退避毫秒数,指数退避:backoff_ms × 2^(attempt−1)1000
poll_interval_msu64worker 轮询间隔(毫秒)100
adapter.namestring传输适配器,默认builtinbuiltin

注意:队列名只是对队列的引用,不会对放入该队列的内容施加任何限制——你可以把任何函数调用入队到任何已声明的队列。

让 record_click 变成可排队任务

你在第 3 章已经写好了link::record_click,函数本身一行都不用改——可排队性只取决于触发方式。只需用TriggerAction改变它的触发方式。

第一步,在link/src/index.ts中导入TriggerAction

import { registerWorker, TriggerAction } from "iii-sdk"; import { Logger } from "@iii-dev/helpers/observability";

第二步,给http::redirect中现有的link::record_click调用加上action,让queueworker 把它入队而不是内联执行:

worker.registerFunction("http::redirect", async (req) => { // ...previous code... await worker.trigger({ function_id: "link::record_click", payload: { code, clicked_at: new Date().toISOString() }, action: TriggerAction.Enqueue({ queue: "clicks" }), }); return { status_code: 302, headers: { Location: url } }; });

从源码看,TriggerAction.Enqueue({ queue })会序列化为{ type: "enqueue", queue }的线上格式,SDK 的契约测试对此有明确锁定(见 sdk/packages/node/iii/tests/trigger-action.test.ts),引擎的TriggerAction反序列化依赖这个格式。

效果

改造后,重定向在点击被接受进队列的瞬间就返回 302,不再等待数据库写入完成。link::record_click在后台排空队列,并且自带重试机制;如果写入持续失败,任务最终进入死信队列(DLQ)。值得一提的是,队列能力不再是引擎内置:从源码看,引擎会把TriggerAction.Enqueuedurable:subscriber路由到独立的queueworker(见 engine/src/workers/queue/README.md),如果你遇到enqueue_error: engine::queue::enqueue not found,说明项目里没有这个 worker,需要把它加进 Compose。

用发布订阅广播事件

队列把每条消息投递给一个消费者。当系统里多个互不相关的部分需要对同一事件做出反应时,应该改用发布/订阅(publish/subscribe)设计。

queue 与 pubsub 的职责划分

项目同时提供了queuepubsub两个 worker,官方文档给出了清晰的选择标准:

  • 需要保证成功(或失败进入 DLQ)的发布订阅流:使用queue提供的iii::durable::publishdurable:subscriber触发类型;
  • 不需要保证的发布订阅流:使用pubsub提供的publishsubscribe

本章会实现两个主题:创建链接时发布link.created,更新链接时发布link.updated。由于目前还没有链接更新功能,我们会一并补上它和一个 HTTP 端点。

在 link.created 上发布

link::create内部,完成数据库写入和state::set之后,触发内置的publish函数:

worker.registerFunction("link::create", async (payload: { url: string; code?: string }) => { // ...previous code... await worker.trigger({ function_id: "publish", payload: { topic: "link.created", data: { code, url } }, }); logger.info("link created", { code, url }); return { code, url }; });

在 link.updated 上发布

首先,给linkworker 增加更新路径,让链接目标可以改变并对外广播。下面是领域函数:它更新数据库行,并通过持久化发布订阅iii::durable::publish,由queueworker 提供)发布link.updated事件。把它放在link/src/index.ts末尾:

worker.registerFunction("link::update", async (payload: { code: string; url: string }) => { const url = /^https?:\/\//i.test(payload.url) ? payload.url : `https://${payload.url}`; await worker.trigger({ function_id: "database::execute", payload: { db: DB, sql: "UPDATE links SET url = ? WHERE code = ?", params: [url, payload.code], }, }); await worker.trigger({ function_id: "iii::durable::publish", payload: { topic: "link.updated", data: { code: payload.code, url } }, }); return { code: payload.code, url }; });

通过 HTTP 暴露链接更新

接着用 HTTP handler 连接link::update,负责校验输入并调用领域函数:

worker.registerFunction("http::update", async (req) => { const code = req.path_params.code; const url = req.body?.url; if (!url) { return { status_code: 400, body: { error: 'missing "url"' }, headers: { "Content-Type": "application/json" }, }; } const link = await worker.trigger<{ code: string; url: string }, { code: string; url: string }>({ function_id: "link::update", payload: { code, url }, }); return { status_code: 200, body: link, headers: { "Content-Type": "application/json" } }; });

最后是绑定http::updatePUT /links/:code的触发器:

worker.registerTrigger({ type: "http", function_id: "http::update", config: { api_path: "/links/:code", http_method: "PUT" }, });

这段代码本身与 pub/sub 无关,但它是本章末尾验证新功能所必需的。

响应式状态:不耦合地保持缓存正确

link::update只改了数据库,没有改状态缓存,因此查询可能读到过期数据。与其在link::update内部处理缓存刷新,不如用持久化订阅者订阅link.updated事件:

worker.registerFunction("link::on_link_updated", async (data: { code: string; url: string }) => { await worker.trigger({ function_id: "state::set", payload: { scope: "links", key: data.code, value: { url: data.url } }, }); }); worker.registerTrigger({ type: "durable:subscriber", function_id: "link::on_link_updated", config: { topic: "link.updated" }, });

持久化 vs 普通发布订阅

这里有一个值得反复琢磨的设计决策:

  • link.updated使用持久化发布订阅iii::durable::publish搭配durable:subscriber触发器,两者都由queueworker 提供。像缓存刷新器这样的消费者必须收到每一次更新——丢失一个事件,缓存就会一直指向过期 URL;
  • link.created留在普通发布订阅(pubsub)上,因为它的唯一消费者是一个"尽力而为"的每日计数器,偶尔漏掉一条无伤大雅。

通用原则:当错过事件会破坏状态时用持久化发布订阅;当只是 fire-and-forget 的扇出时用普通发布订阅。从queueworker 的源码文档看,持久化发布订阅本质上是基于主题的队列:每个订阅了某主题的函数都会收到每条消息的一份拷贝,支持重试与 DLQ(见 engine/src/workers/queue/README.md 的 "Topic-based queues" 说明)。

用 Python 创建分析 worker

队列和事件在单个 worker 内部有用,跨 worker 同样有用。此前所有代码都是 TypeScript,但 iii 的 worker 不限制语言或运行时——这次我们用 Python 写一个统计链接数量的 analytics worker。

创建新 worker

用与第 1 章创建linkworker 相同的方式创建 Python worker。它会生成一个analytics/worker,包含src/main.py示例和iii.worker.yaml清单:

iii worker init analytics --language python

订阅 link.created 事件

把示例src/main.py替换为下面这份代码:订阅link.created事件,统计每个新短链的创建次数:

import os from datetime import datetime, timezone from iii import register_worker, InitOptions from iii_helpers.observability import Logger worker = register_worker( os.environ.get("III_URL", "ws://localhost:49134"), InitOptions(worker_name="analytics"), ) logger = Logger() DB = "analytics" def ensure_schema() -> None: """The analytics worker owns its own table, in its own database.""" worker.trigger( { "function_id": "database::execute", "payload": { "db": DB, "sql": "CREATE TABLE IF NOT EXISTS daily_link_counts (day TEXT PRIMARY KEY, count INTEGER NOT NULL)", }, } ) def on_link_created(data: dict) -> dict: """Runs whenever link publishes `link.created`. Counts links per day.""" day = datetime.now(timezone.utc).strftime("%Y-%m-%d") worker.trigger( { "function_id": "database::execute", "payload": { "db": DB, "sql": "INSERT INTO daily_link_counts (day, count) VALUES (?, 1) " "ON CONFLICT(day) DO UPDATE SET count = count + 1", "params": [day], }, } ) logger.info(f"counted new link {data.get('code')} for {day}") return {"ok": True} ensure_schema() worker.register_function("analytics::on_link_created", on_link_created) worker.register_trigger( { "type": "subscribe", "function_id": "analytics::on_link_created", "config": {"topic": "link.created"}, } ) print("Analytics worker started")

这段代码值得注意的几点:

  • 连接建立register_worker默认连接ws://localhost:49134,可用III_URL环境变量覆盖(Python SDK 中register_workerInitOptions的实现见 sdk/packages/python/iii/src/iii/iii.py);
  • 自持数据:analytics worker 在自己的数据库、自己的表里记账,linkworker 完全不需要知道它的存在;
  • 幂等建表ensure_schema()在启动时执行CREATE TABLE IF NOT EXISTS
  • 按天计数:用ON CONFLICT(day) DO UPDATE SET count = count + 1实现 upsert 式累加。

配置 worker

现有清单analytics/iii.worker.yaml已经够用:

name: analytics runtime: # Base OCI image used as the worker rootfs. base_image: docker.io/iiidev/python:latest scripts: install: pip install -e . start: watchfiles 'python src/main.py'

字段含义:runtime.base_image是 worker 的根文件系统基础 OCI 镜像;scripts.install在 worker 运行时内安装 SDK、可观测性辅助库与源码监视器;scripts.startwatchfiles启动python src/main.py并支持热重载。

analytics 的数据放在自己的数据库里,因此需要给databaseworker 的配置文件追加一个analytics数据库,与第 3 章的primary并列。编辑config/database.yaml,在value: databases:下新增第二项,保存后自动生效,无需重启

id: database name: Database value: databases: primary: pool: acquire_timeout_ms: 5000 idle_timeout_ms: 30000 max: 10 url: sqlite:./data/iii.db analytics: url: sqlite:./data/analytics.db

最后,把新的 analytics worker 加入项目:

iii worker add ./analytics

验证效果

创建五个链接、跟随其中一个几次,并修改它的目标:

# Make some new links for n in $(seq 1 5); do curl -s -X POST http://127.0.0.1:3111/links \ -H 'Content-Type: application/json' -d "{\"url\":\"https://iii.dev/$n\",\"code\":\"analyticslink$n\"}" done

Python worker 会如预期地记录新建链接的数量。用iii trigger直接查询 analytics 数据库验证:

iii trigger database::query db=analytics sql="SELECT day, count FROM daily_link_counts"
{ "rows": [{ "day": "2026-05-27", "count": 5 }], "row_count": 1 }

iii trigger也可用于iii::queue::redrive之类的内置函数,例如把某命名队列的死信消息重新投回主队列,见 engine/src/workers/queue/README.md。)

适配器选择:从本地开发到多实例生产

虽然本章的config/queue.yaml只用了默认的builtin适配器,但从queueworker 的源码文档(engine/src/workers/queue/README.md)可以看到,适配器是决定队列行为的关键一环,值得在选择时留意:

适配器重试DLQFIFO命名队列消费主题 pub/sub多实例外部依赖
builtin
rabbitmqRabbitMQ
redis否(仅发布)Redis

经验法则:本地开发用builtinin_memory存储);单实例生产用builtinfile_based存储);多实例生产用rabbitmqbuiltin适配器可配置store_method: in_memory | file_basedfile_pathrabbitmq需要amqp_url。多实例场景下若只做主题广播,redis适配器也可作为轻量选择。

总结

本章完成了一次典型的"持久化执行"重构:

  1. 重定向不再等待数据库写入:点击记录乘上clicks队列,在后台排空,配合重试与死信队列兜底;
  2. 链接事件经 pub/sub 扇出link.created走普通pubsublink.updatedqueue提供的持久化发布订阅,各自按"错过事件是否会破坏状态"这一标准选型;
  3. 两个消费者完全解耦:Python analytics worker(按天计数)和缓存刷新器(link::on_link_updated)都只关心事件主题,linkworker 不知道它们的存在;
  4. 跨语言协作:同一套触发机制同时服务 TypeScript 与 Python worker,印证了 iii "worker 不绑定语言或运行时"的设计。

下一步是第 5 章:实时流式点击,届时会用专门的click-streamerworker 实时广播每一次点击。

【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/14 17:48:33

为什么不推荐走 Agent 开发?聊聊我的真实踩坑经历

1. 引言 最近两年&#xff0c;Agent&#xff08;智能体&#xff09;概念被炒得火热&#xff0c;各大厂商都在推自己的 Agent 框架和平台。很多开发者看到 Demo 里 Agent 能自动规划、自动调用工具、自动完成任务&#xff0c;觉得这就是未来&#xff0c;于是跃跃欲试&#xff0c…

作者头像 李华
网站建设 2026/9/14 17:47:23

2026高热密度CPU散热决策指南:风冷水冷选型与安装避坑

1. 这不是“买个风扇就完事”的事&#xff1a;为什么2026年选散热器比三年前更烧脑你拆开新买的AMD Ryzen 9 7950X3D或Intel Core i9-14900KS&#xff0c;手心冒汗——不是因为CPU贵&#xff0c;而是因为你突然意识到&#xff1a;这颗芯片的峰值功耗能冲到300W以上&#xff0c;…

作者头像 李华
网站建设 2026/9/14 17:47:15

微信小程序端侧人脸漫画风格迁移实战

简介&#xff1a;本资源是一套基于微信小程序平台的AI人脸漫画化转换源码&#xff0c;面向具备前端开发基础与图像处理兴趣的开发者&#xff0c;解决将真实人脸照片实时转化为卡通风格图像的技术实践需求。压缩包共121个文件&#xff0c;包含19个JS逻辑文件&#xff08;如index…

作者头像 李华