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 的队列配置、pubsub与iii::durable::publish两种发布订阅模式的取舍,以及如何用 Python 编写一个跨语言、跨 worker 的事件消费者。
原文出处:docs/tutorials/linkly/durable-execution.mdx(本教程的.skill.md渲染版本为 durable-execution.mdx.skill.md)。
本章目标:把写库工作移出热路径
回顾上一章(Ch. 3)的实现:http::redirect直接触发link::record_click,也就是每次用户访问短链时,重定向函数都要同步等待一次数据库写入完成。这带来两个问题:
- 延迟耦合:数据库写入慢时,重定向也跟着变慢,用户体验直接受损;
- 责任耦合:重定向这个"读路径"上背负了"写日志"的额外职责。
本章的解法是用两个既有的 worker 能力重构:
- 队列(queue):接受"现在提交、稍后执行"的工作,让
link::record_click在后台排空(drain)队列,重定向立刻返回 302; - 发布订阅(pub/sub):队列把每条消息投递给一个消费者,而当系统里多个互不相关的部分需要对同一事件做出反应时,改用发布/订阅设计,把事件广播给所有订阅者。
最终得到两个完全解耦的消费者:一个Python analytics worker(按天统计短链创建数)和一个缓存刷新器(link::on_link_updated),linkworker 甚至不需要知道它们的存在。
添加 worker:加入 pubsub
本章涉及两个 worker。queue在iii project init初始化项目时就已经内置,无需重复添加;pubsub则需要显式加入。在项目目录下执行:
iii worker add pubsub这条命令会把pubsubworker 注册到当前项目。注意:queue自第 1 章起就在运行,它的设置统一管理在./config/queue.yaml中(详见 Configuration)。
用队列让重定向变快
队列的基本概念
队列保存"现在接受、以后运行"的工作。本章定义了一个名为clicks的队列,专门用于link::record_click。
定义 clicks 队列:编辑 config/queue.yaml
在./config/queue.yaml的value:下添加queue_configs并保存。配置热生效,无需重启:
id: queue # ... value: queue_configs: clicks: type: standard max_retries: 5 concurrency: 5 adapter: name: builtin各字段说明(结合 engine/src/workers/queue/README.md 中的配置表):
| 字段 | 类型 | 说明 | 默认值 |
|---|---|---|---|
type | string | standard(并发处理)或fifo(按消息组有序处理) | — |
max_retries | u32 | 进入死信队列(DLQ)前的最大投递尝试次数 | 3 |
concurrency | u32 | 同时处理的最大任务数;FIFO 队列会强制为prefetch=1 | 10 |
message_group_field | string | 仅fifo队列需要,指定决定顺序组的 JSON 字段 | — |
backoff_ms | u64 | 重试基础退避毫秒数,指数退避:backoff_ms × 2^(attempt−1) | 1000 |
poll_interval_ms | u64 | worker 轮询间隔(毫秒) | 100 |
adapter.name | string | 传输适配器,默认builtin | builtin |
注意:队列名只是对队列的引用,不会对放入该队列的内容施加任何限制——你可以把任何函数调用入队到任何已声明的队列。
让 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.Enqueue与durable:subscriber路由到独立的queueworker(见 engine/src/workers/queue/README.md),如果你遇到enqueue_error: engine::queue::enqueue not found,说明项目里没有这个 worker,需要把它加进 Compose。
用发布订阅广播事件
队列把每条消息投递给一个消费者。当系统里多个互不相关的部分需要对同一事件做出反应时,应该改用发布/订阅(publish/subscribe)设计。
queue 与 pubsub 的职责划分
项目同时提供了queue和pubsub两个 worker,官方文档给出了清晰的选择标准:
- 需要保证成功(或失败进入 DLQ)的发布订阅流:使用
queue提供的iii::durable::publish与durable:subscriber触发类型; - 不需要保证的发布订阅流:使用
pubsub提供的publish与subscribe。
本章会实现两个主题:创建链接时发布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::update到PUT /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_worker与InitOptions的实现见 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.start用watchfiles启动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\"}" donePython 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)可以看到,适配器是决定队列行为的关键一环,值得在选择时留意:
| 适配器 | 重试 | DLQ | FIFO | 命名队列消费 | 主题 pub/sub | 多实例 | 外部依赖 |
|---|---|---|---|---|---|---|---|
builtin | 是 | 是 | 是 | 是 | 是 | 否 | 无 |
rabbitmq | 是 | 是 | 是 | 是 | 是 | 是 | RabbitMQ |
redis | 否 | 否 | 否 | 否(仅发布) | 是 | 是 | Redis |
经验法则:本地开发用builtin(in_memory存储);单实例生产用builtin(file_based存储);多实例生产用rabbitmq。builtin适配器可配置store_method: in_memory | file_based与file_path;rabbitmq需要amqp_url。多实例场景下若只做主题广播,redis适配器也可作为轻量选择。
总结
本章完成了一次典型的"持久化执行"重构:
- 重定向不再等待数据库写入:点击记录乘上
clicks队列,在后台排空,配合重试与死信队列兜底; - 链接事件经 pub/sub 扇出:
link.created走普通pubsub,link.updated走queue提供的持久化发布订阅,各自按"错过事件是否会破坏状态"这一标准选型; - 两个消费者完全解耦:Python analytics worker(按天计数)和缓存刷新器(
link::on_link_updated)都只关心事件主题,linkworker 不知道它们的存在; - 跨语言协作:同一套触发机制同时服务 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),仅供参考