news 2026/9/14 10:35:09

iii 项目实战:用队列与持久化 Pub/Sub 解耦 Linkly 重定向热路径(durable-execution 第 4 章)

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
iii 项目实战:用队列与持久化 Pub/Sub 解耦 Linkly 重定向热路径(durable-execution 第 4 章)

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::publishdurable: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-queueiii project init时已经存在于你的项目中,负责提供标准队列与持久化的发布/订阅能力;iii-pubsub需要现在手动添加,本章稍后要向它发布第一个事件:

iii worker add iii-pubsub

关于 worker 的安装机制:iii worker add会从 worker registry 解析并写入项目的config.yamliii.lockiii.lock记录 worker 依赖与版本,保证部署可复现,类似于其他包管理器的 lock 文件(参见 第 1 章:基础)。

用队列让重定向变快

队列的本质是「现在接收、稍后执行」:clicks队列中的每一条消息,都代表一次「在 T 时刻有人跟进了短码 X」的插入操作。定义队列的位置是iii-queueworker(iii project init已内置),只需要在config.yamlqueue_configs中补充几项:

workers: # ... - name: iii-queue config: queue_configs: clicks: type: standard max_retries: 5 concurrency: 5 adapter: name: builtin

queue_configs的完整参数说明

上述配置对应引擎侧queueworker 的配置结构(字段定义与校验见 config.rs 与 configuration.rs,参数明细可对照 queue worker 文档):

字段类型说明默认值
typestringstandard(并发处理)或fifo(消息组内严格有序)standard
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决定队列的传输后端。默认的builtin是进程内实现,无外部依赖,适合单实例部署;生产多实例场景可切换rabbitmq(完整支持重试、DLQ、FIFO 与命名队列消费)或redis(仅支持主题发布/订阅,不支持命名队列消费与重试),适配器对比详见 queue worker 文档。

改写http::redirect:从同步调用到入队

你在第 3 章已经写好了link::record_clickhttp::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-queueiii-pubsub两个 worker。iii-queue在提供标准队列的同时,也自带持久化的发布与订阅能力。

  • 当发布/订阅流程必须保证成功(失败则进入 DLQ)时,使用iii-queueiii::durable::publishdurable:subscriber
  • 当发布/订阅流程不需要保证时,使用iii-pubsubpublishsubscribe

本章将实现两个主题:链接创建时发布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::publishdurable: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:subscriberiii-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\"}" done

Python 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 的调用拓扑发生了本质变化:

  1. 重定向不再等待数据库写入http::redirect通过TriggerAction.Enqueue({ queue: "clicks" })把点击记录投递到clicks队列,302立即返回,link::record_click在后台排空队列(max_retries: 5,超限进 DLQ);
  2. 链接事件经 pub/sub 扇出link.creatediii-pubsub常规发布,link.updatediii::durable::publish持久化发布;
  3. 消费者彼此独立: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),仅供参考

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

绿色证书与综合能源系统优化调度模型解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/14 10:28:42

LlamaIndex文档处理与知识检索框架实战指南

1. LlamaIndex核心功能解析LlamaIndex是一个专注于文档处理与知识检索的开源框架&#xff0c;其核心价值在于将非结构化文档转化为可被大语言模型高效利用的知识库。我在实际项目中用它处理过技术手册、财务报告等复杂文档&#xff0c;最直观的感受是它解决了传统OCR工具对表格…

作者头像 李华
网站建设 2026/9/14 10:25:55

OpenSandbox 实战:在沙箱内启动 OpenClaw Gateway 并暴露 HTTP 端点

OpenSandbox 实战&#xff1a;在沙箱内启动 OpenClaw Gateway 并暴露 HTTP 端点 【免费下载链接】OpenSandbox Secure, Fast, and Extensible Sandbox runtime for AI agents. 项目地址: https://gitcode.com/GitHub_Trending/ope/OpenSandbox 本指南以 OpenSandbox 官方…

作者头像 李华
网站建设 2026/9/14 10:24:37

Python虚拟环境实战:从原理到企业级应用

1. Python虚拟环境核心价值解析在Python开发领域&#xff0c;虚拟环境&#xff08;venv&#xff09;是项目依赖管理的基石工具。我经历过多个Python项目因缺乏环境隔离导致的"依赖地狱"——不同项目对同一包有冲突版本要求时&#xff0c;系统级的Python环境会陷入混乱…

作者头像 李华
网站建设 2026/9/14 10:22:54

AI论文写作工具测评与继续教育应用指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华