用 iii-stream 构建实时点击流:在 Linkly 中通过 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/0-19-0/tutorials/linkly/streaming.mdx):在已经具备存储(Chapter 3)与队列/发布订阅(Chapter 4)能力的 URL 短链服务上,引入引擎内置的iii-stream工作进程,新建一个专职的click-streamerworker,把每一次点击在发生的那一刻实时推送给浏览器订阅端。读完本章,你将掌握iii worker add/iii worker init的工作进程脚手架流程、iii-pubsub事件发布与stream::set广播的协作方式,以及stream::list命令行验证手段,为第七章的浏览器实时计数器打好基础。
为什么需要一个专门的流工作进程
iii-stream是 iii 引擎内置的实时数据传输能力:数据一变,立即推送给客户端,而不是让客户端轮询。它天然适用于"仪表盘上的实时点击流"这类场景。
流(Stream)本身是双向的——订阅者既能接收消息,也能往回发送消息;但在 Linkly 的场景里,我们只需要把点击事件向外广播。为了保持关注点分离,教程把"实时广播"这个职责单独放进一个click-streamerworker,让原有的linkworker 继续只负责链接本身的创建与解析,互不干扰。
从引擎实现看,iii-stream是一个engine-owned的 worker(在 stream.rs 中实现,文档见 README.md),数据以三级层级组织:stream_name>group_id>item_id。客户端通过 WebSocket 订阅某个(stream_name, group_id),当条目变化时实时收到更新。这正是后面章节里浏览器端订阅clicks/all并实时计数的基础。
添加两个工作进程
iii-stream是第七章把点击发送给客户端的方式,因此先把它加进项目,再像第一章创建link、第四章创建analytics一样,新建click-streamerworker:
iii worker add iii-stream iii worker init click-streamer --language typescriptiii worker add iii-stream把内置的流工作进程注册进项目配置;iii worker init click-streamer --language typescript生成一个 TypeScript 骨架 worker(生成文件在click-streamer/src/index.ts),稍后我们会替换它的默认内容。
广播:让link只负责"宣布",让click-streamer负责"推送"
继续沿用第四章的解耦思路:linkworker 不直接连接任何客户端,它只负责发布一个"点击发生了"的事件;click-streamer订阅到这个事件后,把它推上实时流。由于一个实时计数器可以容忍偶发的丢事件,这里使用普通的iii-pubsub事件就足够了,不需要队列级别的可靠性保证。
第一步:在linkworker 中发布link.clicked事件
修改link/src/index.ts的link::record_click函数,在写入数据库之后发布事件:
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()让函数立即返回而不等待触发完成,这是针对 pubsub 这类"不要求保证执行"场景的简单性能优化——你不希望一次点击因为等待发布完成而拖慢重定向。
数据库写入仍然使用await(这是必须落地的持久化操作),而事件发布则走 fire-and-forget,两者语义截然不同。
第二步:编写click-streamerworker
click-streamer订阅link.clicked主题,并把每次点击通过stream::set广播到名为clicks的流中。stream::set同时做两件事:存储条目,以及把它推送给订阅了该流与组的每一个 WebSocket 客户端。
用下面的内容替换生成的click-streamer/src/index.ts:
import { registerWorker } from "iii-sdk"; import { Logger } from "@iii-dev/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_name(clicks)、group_id(all)、item_id(用code-clicked_at保证每次点击的唯一性)、data(点击数据本身)。之后订阅clicks/all的所有客户端都能收到这条消息;registerTrigger:注册一个subscribe类型触发器,把link.clicked主题上的每条消息路由到click-streamer::broadcast函数,形成"事件 → 函数 → 流广播"的完整链路;- 连接地址:默认回退到
ws://localhost:49134,即引擎的 WebSocket 总线地址,可通过III_URL环境变量覆盖。
第三步:注册到项目
iii worker add ./click-streamer注册完成后,第七章构建的浏览器端会订阅clicks/all,并把收到的广播实时计数展示出来。
源码视角:stream::set内部发生了什么
iii-stream的stream::set实现在 stream.rs 的StreamWorker::set中(约 stream.rs)。一次写入会依次做三件事:
- 持久化:调用当前配置的 adapter 的
set方法存储数据(stream.rs中adapter.set(&stream_name, &group_id, &item_id, data)); - 触发
stream触发器:构建包含event_type、stream_name、group_id、item_id与event(Create/Update)的StreamWrapperMessage,通过invoke_triggers异步派发给匹配的处理器——注意是 spawn 出的独立任务,写调用先返回,处理器失败不会回滚写入; - 通知所有订阅者:通过 adapter 的
emit_event把变更推给该(stream_name, group_id)上的所有 WebSocket 客户端。
如果条目之前不存在,客户端收到的是Create事件;如果已存在则收到Update事件——这正是浏览器端StreamEvent类型中event.type取"create" | "update" | "delete"的原因(参见 frontend.mdx 中的类型定义)。
iii-stream还提供了完整的stream::*函数面与触发器面(详见 README.md):
| 函数 | 作用 |
|---|---|
stream::set | 存储条目并广播 create/update 事件,返回旧值和新值 |
stream::get | 按(stream_name, group_id, item_id)读取单个条目 |
stream::delete | 删除条目并广播携带被删值的 delete 事件 |
stream::list | 列出某个组内的全部条目 |
stream::list_groups | 列出某个流下的全部组 |
stream::list_all | 列出所有流及其组元数据 |
stream::send | 向组订阅者广播瞬时事件(不持久化,如打字指示器) |
stream::update | 用set/merge/increment/decrement/append/remove原子更新条目 |
触发器方面,除了教程用到的订阅型subscribe(属于iii-pubsub),iii-stream还提供stream(条目变化时触发,可用stream_name/group_id/item_id过滤)、stream:join与stream:leave(WebSocket 订阅者接入/断开时触发,stream:join处理器返回{ unauthorized: true }可在数据流出前拒绝订阅)。
存储后端:kv 与 redis
iii-stream的数据持久化与实时投递都交给可插拔的 adapter:
kv(默认):内置键值存储,支持in_memory(重启丢失)与file_based(落盘持久化)两种模式,无需外部依赖,适合单实例部署;redis:以 Redis 为后端存储流数据,并用 Redis Pub/Sub 做实时投递,多实例部署且需要跨进程实时广播时使用。
默认适配器为kv(见 stream.rs 中DEFAULT_ADAPTER_NAME常量),可在 worker 配置中显式声明:
engine: workers: iii-stream: port: ${STREAM_PORT:3112} host: 0.0.0.0 adapter: name: redis config: redis_url: ${REDIS_URL:redis://localhost:6379}iii-stream还通过内置的configurationworker 注册了自己的运行期配置(id 为iii-stream),port/host/adapter等字段可以在不重启引擎的情况下热更新。
验证:让每一次重定向都出现在实时流里
引擎运行起来后,创建一条测试链接,访问它几次,再读取clicks流的实时内容:
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 iii trigger stream::list stream_name=clicks group_id=all- 第一条命令创建短链
stream-me,对应POST /linksHTTP 端点; - 循环访问三次
http://127.0.0.1:3111/s/stream-me,每次访问都会触发重定向并记录一次点击; iii trigger stream::list读取clicks流all组下的全部条目。
每一次重定向落地,linkworker 都会发布link.clicked,click-streamer随即把它广播进clicks流——所以这条命令的输出里应该能看到 3 条点击记录。整个链路没有轮询:事件从发布到出现在流中,靠的是 pubsub 订阅与流广播两级推送。
收尾:从服务端流到浏览器计数
至此,Linkly 通过一个专职的click-streamerworker,把每一次点击实时流式推送给订阅者,而linkworker 依旧保持纯粹。下一章(Ch. 6: Move bulk data with channels)将用 channel 把 CSV 中的链接一次性流式批量导入。
关于浏览器端如何消费这个流,可以提前预告(详见 frontend.mdx):浏览器 worker 通过registerFunction暴露一个ui::on_click函数,再用registerTrigger注册stream类型触发器,config指向{ stream_name: "clicks", group_id: "all" };每次click-streamer广播新条目,该函数就被调用一次,计数器加一。这正是本教程第五章"服务端广播"与第七章"浏览器订阅"首尾呼应的完整链路。
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考