用 iii 的 iii-stream 为 Linkly 实现实时点击流推送:从 pubsub 事件到 WebSocket 广播的完整实战
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
本篇是 Linkly 短链服务系列教程(见 docs/tutorials/linkly)第 5 章的实战技术指南,讲解如何利用引擎内置的iii-stream实时流能力,把每一次点击事件在毫秒级推送到订阅端。读完你将掌握iii worker add/iii worker init搭建流式 worker、用TriggerAction.Void()做解耦的事件发布、用stream::set完成"存储 + 广播"、以及用iii trigger stream::list验证实时链路,为第 7 章的浏览器实时计数器打好基础。
为什么需要 iii-stream:实时数据与普通调用的区别
在前面的章节中,Linkly 的linkworker 已经可以通过database::execute把点击记录写入数据库。但"记录点击"和"让仪表盘实时看到点击"是两件不同的事:
- 数据库写入面向持久化,查询靠轮询,无法主动推送;
- 实时仪表盘需要数据一变就到达客户端,不能靠浏览器反复刷新。
iii-stream就是为这种实时传输而生的。正如 engine/src/workers/stream/README.md 所描述的:它是一个由引擎托管的 worker(engine-owned),把实时数据组织成stream_name→group_id→item_id的三层层级结构,客户端通过 WebSocket 订阅,一旦条目变化就立即收到更新。
流本身是双向的——订阅者既能接收消息也能回发消息。但本章的需求很简单:只把点击事件向外广播。因此教程采用了一个清晰的职责划分:
linkworker 专注于链接本身(创建、解析、记录点击);- 新建一个专门的
click-streamerworker,独占"实时广播"这一职责。
这种解耦延续了第 1 章以来"一个 worker 只做一件事"的设计,参见 第 1 章 Foundations。
第一步:添加 iii-stream 并脚手架 click-streamer
iii-stream是引擎自带的系统 worker(在 crates/iii-worker/src/cli/builtin_defaults.rs 的BUILTIN_NAMES中与iii-sandbox并列),所以不用手写,直接把它加进项目:
iii worker add iii-stream iii worker init click-streamer --language typescript两条命令的含义分别对应第 1 章建立的工作流:
iii worker add <name>把内置 worker 写入项目的config.yaml(也可以在注册本地 worker 时指向目录,如iii worker add ./click-streamer);iii worker init click-streamer --language typescript在项目内脚手架出一个 TypeScript 的 worker 目录click-streamer/,其运行方式由click-streamer/iii.worker.yaml清单描述。
iii-stream的默认监听端口是3112,默认适配器是无需外部依赖的kv(支持in_memory与file_based两种存储方式);需要多实例跨进程实时扇出时才切换到redis适配器,相关配置字段见 engine/src/workers/stream/README.md 的 Adapters 一节。
让 link worker 发布 link.clicked 事件
保持linkworker 解耦的关键是:它只宣布"发生了一次点击",而把"如何推送给订阅者"留给click-streamer。做法是在link::record_click里,写完数据库后顺带向 pubsub 主题link.clicked发布一条事件:
worker.registerFunction( "link::record_click", async (payload: { code: string; clicked_at: string }) => { await worker.trigger({ function_id: "database::execute", payload: { db: DB, sql: "INSERT INTO clicks (code, clicked_at) VALUES (?, ?)", params: [payload.code, payload.clicked_at], }, }); worker.trigger({ function_id: "publish", payload: { topic: "link.clicked", data: payload }, action: TriggerAction.Void(), }); return { recorded: true }; }, );注意这里的两个细节:
- 没有
await,并且显式设置了action: TriggerAction.Void()。TriggerAction.Void()是 SDK 提供的 fire-and-forget 路由方式:调用立即返回、不等结果(见 sdk/packages/node/iii/README.md 中 "Invoke (fire-and-forget)" 一行)。对 pubsub 这类"尽力而为"的场景,这能避免发布动作拖慢点击记录的主链路,是一种简单的性能优化; - 之所以可以用普通 pubsub 而不是可靠队列,是因为实时计数能容忍极少数事件丢失——掉一两次点击计数,仪表盘不会崩。引擎内置的 topic 型 pubsub 采用"每个订阅了该主题的函数各收到一份消息副本"的扇出语义(参见 engine/src/workers/queue/README.md 对 topic-based publish 的描述),正好满足一对多的广播需求。
搭建 click-streamer:订阅主题并广播到流
现在编写click-streamerworker。它做两件事:通过subscribe触发器订阅link.clicked主题;收到事件后用stream::set把点击写入clicks流。替换脚手架生成的click-streamer/src/index.ts:
import { registerWorker } from "iii-sdk"; import { Logger } from "@iii-dev/helpers/observability"; const worker = registerWorker(process.env.III_URL ?? "ws://localhost:49134", { workerName: "click-streamer", }); const logger = new Logger(); worker.registerFunction( "click-streamer::broadcast", async (data: { code: string; clicked_at: string }) => { await worker.trigger({ function_id: "stream::set", payload: { stream_name: "clicks", group_id: "all", item_id: `${data.code}-${data.clicked_at}`, data, }, }); return { streamed: true }; }, ); worker.registerTrigger({ type: "subscribe", function_id: "click-streamer::broadcast", config: { topic: "link.clicked" }, }); logger.info("click-streamer ready");stream::set 的语义:一次调用,三步动作
stream::set是iii-stream的核心写操作(源码定义于 engine/src/workers/stream/stream.rs#L991-L1044,描述为 "Set a value in a stream")。按 engine/src/workers/stream/README.md 的说明,一次stream::set会依次完成:
- 持久化:通过当前适配器(默认
kv,可换redis)保存条目; - 广播:通知所有订阅了该
(stream_name, group_id)的 WebSocket 客户端; - 触发:评估并触发注册的
stream系列触发器。
参数与返回值的完整定义同样记录在 engine/src/workers/stream/README.md:
| 参数 | 类型 | 说明 |
|---|---|---|
stream_name | string | 流名称,如clicks |
group_id | string | 组标识,如all(第 7 章浏览器就订阅clicks/all) |
item_id | string | 条目 ID,本章用${code}-${clicked_at}保证每次点击唯一 |
data | any | 要存储并广播的数据负载 |
返回old_value与new_value。从 stream.rs 的源码可以看出,当没有注册自定义覆盖函数时,stream::set最终落到adapter.set(&stream_name, &group_id, &item_id, data)这一行——存储与广播的职责被收敛在适配器层。
触发器的选择:subscribe 类型
worker.registerTrigger把click-streamer::broadcast绑定到link.clicked主题。这样linkworker 发布事件后,引擎会把消息扇出给click-streamer::broadcast,由它转写成stream::set调用。第 1 章已经提到"每个已注册的函数自带可调用触发器",而这里显式注册的subscribe触发器则是"外部事件驱动函数执行"的典型用法。
注册到项目
iii worker add ./click-streamer验证实时链路:curl 触发,stream::list 读取
引擎运行中,先创建一条带自定义短码的链接,再连续访问 3 次,制造 3 次点击:
curl -s -X POST http://127.0.0.1:3111/links \ -H 'Content-Type: application/json' -d '{"url":"https://iii.dev","code":"stream-me"}' for n in $(seq 1 3); do curl -s -o /dev/null http://127.0.0.1:3111/s/stream-me; done然后读取clicks流的实时内容:
iii trigger stream::list stream_name=clicks group_id=all链路全貌如下:
- 每次
GET /s/stream-me命中linkworker 的 HTTP 重定向函数,内部调用link::record_click; link::record_click写库后向link.clicked发布事件(fire-and-forget);click-streamer的 subscribe 触发器收到事件,调用click-streamer::broadcast;click-streamer::broadcast执行stream::set,点击条目落入clicks/all流并广播给所有订阅者;iii trigger stream::list把流内已广播的条目枚举出来,作为验证依据。
iii trigger的key=value参数格式与第 1 章调用link::create的方式一致(参见 第 1 章 的 "Call the functions" 一节),因此这条命令也可以直接理解为"以命令行调用引擎函数"。
向前看:第 7 章浏览器如何消费这个流
本章埋下的clicks/all流,正是第 7 章浏览器端实时计数器的数据源。从 frontend.mdx 可以看到后续的接入方式:
- 浏览器通过
iii-worker-manager的 RBAC 门控监听器连接引擎,配置中用expose_functions白名单放行stream::*系列函数,浏览器才能订阅流; - 浏览器 SDK 通过单一引擎 WebSocket 订阅
stream变更并随事件重渲染——这正是stream::set第 2 步"广播给所有订阅该 stream/group 的 WebSocket 客户端"的消费端。
(注:直接连流端口ws://host:3112/stream/<stream_name>/<group_id>/的方式已被官方标记为 deprecated,推荐使用 Browser SDK,见 engine/src/workers/stream/README.md 的 Client Subscriptions 一节。)
小结
至此,Linkly 拥有了完整的实时点击流链路:linkworker 只负责宣布事件,click-streamer独占"订阅 + 广播"职责,iii-stream的三层流模型把每次点击同时送达持久化存储与所有 WebSocket 订阅者。下一章 Ch. 6:用 channels 批量迁移数据 将延续"职责分离"的思路,用一次流式上传把 CSV 里的链接批量导入。
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考