iii 项目实战:用队列与持久化 Pub/Sub 解耦 Linkly 重定向热路径(durable-execution 第 4 章)
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
本教程章节聚焦 Linkly 短链服务的「持久化执行」改造:把点击记录的数据库写入从重定向热路径搬上队列,让302立即返回;再用发布/订阅(pub/sub)把链接创建、更新事件广播给独立的 Python 分析 worker 与缓存刷新订阅者,实现linkworker 对下游的一无所知。读完本章,你将掌握queue_configs队列配置、TriggerAction.Enqueue异步路由、iii::durable::publish与durable:subscriber持久化事件流,以及如何在同一个项目中混合使用 TypeScript 与 Python worker。
本章是 Linkly 系列教程的第 4 章,前置内容为 第 3 章:持久化。在第 3 章里,http::redirect通过worker.trigger直接调用link::record_click,数据库写入落在重定向的同步调用链上——一次慢写就会拖慢整个跳转。本章的目标是把这段工作搬上队列,让重定向立即返回,再用pub/sub把链接事件广播给相互独立的订阅方:一个 Python 分析 worker 与一个缓存刷新器,二者都与linkworker 完全解耦。
添加 worker:iii-pubsub
本章需要两个 worker。iii-queue在iii project init时已经存在于你的项目中,负责提供标准队列与持久化的发布/订阅能力;iii-pubsub需要现在手动添加,本章稍后要向它发布第一个事件:
iii worker add iii-pubsub关于 worker 的安装机制:
iii worker add会从 worker registry 解析并写入项目的config.yaml与iii.lock。iii.lock记录 worker 依赖与版本,保证部署可复现,类似于其他包管理器的 lock 文件(参见 第 1 章:基础)。
用队列让重定向变快
队列的本质是「现在接收、稍后执行」:clicks队列中的每一条消息,都代表一次「在 T 时刻有人跟进了短码 X」的插入操作。定义队列的位置是iii-queueworker(iii project init已内置),只需要在config.yaml的queue_configs中补充几项:
workers: # ... - name: iii-queue config: queue_configs: clicks: type: standard max_retries: 5 concurrency: 5 adapter: name: builtinqueue_configs的完整参数说明
上述配置对应引擎侧queueworker 的配置结构(字段定义与校验见 config.rs 与 configuration.rs,参数明细可对照 queue worker 文档):
| 字段 | 类型 | 说明 | 默认值 |
|---|---|---|---|
type | string | standard(并发处理)或fifo(消息组内严格有序) | standard |
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决定队列的传输后端。默认的builtin是进程内实现,无外部依赖,适合单实例部署;生产多实例场景可切换rabbitmq(完整支持重试、DLQ、FIFO 与命名队列消费)或redis(仅支持主题发布/订阅,不支持命名队列消费与重试),适配器对比详见 queue worker 文档。
改写http::redirect:从同步调用到入队
你在第 3 章已经写好了link::record_click,http::redirect直接触发它。函数本身一行都不用改——变的只是调用方式。
首先把TriggerAction导入link/src/index.ts:
import { registerWorker, TriggerAction } from "iii-sdk"; import { Logger } from "@iii-dev/observability";然后给http::redirect里现有的link::record_click调用加上action,让iii-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 } }; });现在重定向在点击被队列接受的那一刻就立即返回。link::record_click在后台排空队列,即使写入持续失败,也有重试机制与死信队列兜底。
从 SDK 层面看,TriggerAction是{ type: 'enqueue'; queue: string } | { type: 'void' }的联合类型(定义见 iii-types.ts),TriggerAction.Enqueue({ queue })与TriggerAction.Void()是引擎侧路由语义的构造器(见 iii.ts):
- 省略
action:同步 request/response,等待被调用函数返回; TriggerAction.Enqueue(...):经命名队列异步路由,引擎确认入队后即返回Promise<EnqueueResult>;TriggerAction.Void():fire-and-forget,不等待响应。
该路由契约由 SDK 自带的契约测试覆盖,见 trigger-action.test.ts,其中明确断言TriggerAction.Enqueue({ queue: 'orders' })序列化结果与{ type: 'void' }的表示形式。
队列的可靠性:重试与死信队列
max_retries: 5意味着每条消息最多投递 5 次;超过后消息进入该队列的 DLQ。引擎提供内置函数iii::queue::redrive可以把某个命名队列 DLQ 中的消息全部重新放回主队列:
iii trigger \ --function-id='iii::queue::redrive' \ --payload='{"queue": "clicks"}'返回结构为{ "queue": "clicks", "redriven": <数量> },语义与参数定义见 queue worker 文档 的 Builtin Functions 一节。
用发布/订阅广播事件
队列把每条消息投递给一个消费者;当系统中多个互不相关的部分需要响应同一个事件时,应该改用发布/订阅(publish/subscribe)设计。
两种发布/订阅怎么选项目同时提供
iii-queue与iii-pubsub两个 worker。iii-queue在提供标准队列的同时,也自带持久化的发布与订阅能力。
- 当发布/订阅流程必须保证成功(失败则进入 DLQ)时,使用
iii-queue的iii::durable::publish与durable:subscriber;- 当发布/订阅流程不需要保证时,使用
iii-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.created走的是iii-pubsub的常规publish——它的唯一消费者是一个尽力而为(best-effort)的每日计数器,偶尔漏掉一条事件无伤大雅。
发布link.updated
先补充领域函数link::update:更新数据库行,并通过持久化pub/sub(iii::durable::publish,由iii-queue提供)发布link.updated事件:
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 }; });iii::durable::publish在引擎侧的定位是主题型队列的生产端,会向每个订阅了该主题的不同函数扇出(fan-out)消息副本;其载荷约定为{ topic: string, data: any },实现定义见 queue.rs 的#[function(id = "iii::durable::publish")]。
通过 HTTP 暴露链接更新
先写校验输入并调用领域函数的 HTTP handler:
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" } }; });再绑定触发器,让它响应PUT /links/:code:
worker.registerTrigger({ type: "http", function_id: "http::update", config: { api_path: "/links/:code", http_method: "PUT" }, });响应式状态:不耦合地保持缓存正确
link::update改了数据库,却没有改状态缓存,查询可能读到过期数据。与其在link::update内部顺手刷新缓存,不如用持久化订阅者订阅link.updated事件:
持久化 vs 常规 pub/sub 的取舍。
link.updated使用持久化 pub/sub:iii::durable::publish配durable:subscriber触发器,二者都由iii-queueworker 提供。像缓存刷新器这样的消费者必须收到每一次更新——丢一个事件,缓存就会指向过期 URL。而link.created停留在常规 pub/sub(iii-pubsub),它的消费者只是尽力而为的每日计数器,偶尔漏一次无关紧要。规则:当错过事件会污染状态时用持久化 pub/sub;纯扇出、可容忍丢失时用常规 pub/sub。
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" }, });durable:subscriber是iii-queue提供的消费端触发器类型(枚举与注册路径见 queue.rs 与 builtin 适配器):每有一个订阅同一主题的独立函数,就为其投递一份消息副本,投递带重试与 DLQ 保障。订阅者竞争模式下,同一函数的多个副本之间互为竞争消费者。
创建 Python 分析 worker
队列与事件不仅在一个 worker 内部有用,也能跨 worker 生效。到目前为止所有代码都是 TypeScript,但 worker 并不限定语言或运行时——这次就用 Python 创建一个统计链接数的分析 worker。
创建新 worker
用与第 1 章创建linkworker 相同的方式脚手架一个 Python worker,会生成带src/main.py示例与iii.worker.yamlmanifest 的analytics/目录:
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_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")Python 侧 SDK 的入口是register_worker(address, options),可自动连接引擎(地址优先取显式参数,其次取III_URL环境变量);InitOptions(worker_name=...)用于声明 worker 名称。register_trigger({...})则把触发器绑定到已注册函数上,API 形态与实现见 iii.py 与 Python SDK README。
关键设计点:分析 worker 在自己的数据库里拥有自己的表——daily_link_counts按天累计link.created次数,插入语句利用ON CONFLICT(day) DO UPDATE实现幂等累加。这样linkworker 永远不需要知道analytics的存在。
配置 worker 的 manifest
生成的 manifest 还没有运行脚本,修改iii.worker.yaml补上 install 与 start 脚本:
scripts: install: "pip install watchfiles && pip install -e ." start: "watchfiles 'python src/main.py'"watchfiles负责监听源码变化并自动重启 Python 进程(对应 TypeScript worker 侧tsx watch的体验)。分析数据放在独立数据库中,因此需要在databaseworker 里、在第 3 章primary数据库旁边再加一个analytics数据库:
workers: # ... - name: database config: databases: primary: # ... url: sqlite:./data/iii.db analytics: url: sqlite:./data/analytics.db最后把新 worker 加进config.yaml:
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\":\"link$n\"}" donePython worker 如预期地统计到了新链接创建次数:
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 <function_id> key=value形式的 CLI 调用是 iii 手动触发函数的方式,函数名后也支持--help查看参数(参见 第 3 章 结尾的提示)。
小结:链路全景
完成本章后,Linkly 的调用拓扑发生了本质变化:
- 重定向不再等待数据库写入:
http::redirect通过TriggerAction.Enqueue({ queue: "clicks" })把点击记录投递到clicks队列,302立即返回,link::record_click在后台排空队列(max_retries: 5,超限进 DLQ); - 链接事件经 pub/sub 扇出:
link.created走iii-pubsub常规发布,link.updated走iii::durable::publish持久化发布; - 消费者彼此独立:Python 分析 worker 通过
subscribe触发器统计每日建链数,缓存刷新器通过durable:subscriber触发器在link.updated后刷新state::set——linkworker 完全不知道这两者存在,新增消费者无需改动生产者代码。
这种「队列解耦写路径 + 持久化事件流解耦读路径」的组合,正是 iii 里组合、扩展与实时观测每个服务的基本范式。下一章 第 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),仅供参考