ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

用 iii-stream 构建实时点击流:在 Linkly 中通过 WebSocket 向订阅者广播每一次点击

2026/9/14 2:34:01 拓冰建站 浏览量
用 iii-stream 构建实时点击流:在 Linkly 中通过 WebSocket 向订阅者广播每一次点击 用 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_namegroup_iditem_id。客户端通过 WebSocket 订阅某个(stream_name, group_id)当条目变化时实时收到更新。这正是后面章节里浏览器端订阅clicks/all并实时计数的基础。添加两个工作进程iii-stream是第七章把点击发送给客户端的方式因此先把它加进项目再像第一章创建link、第四章创建analytics一样新建click-streamerworkeriii 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-streamerworkerclick-streamer订阅link.clicked主题并把每次点击通过stream::set广播到名为clicks的流中。stream::set同时做两件事存储条目以及把它推送给订阅了该流与组的每一个 WebSocket 客户端。用下面的内容替换生成的click-streamer/src/index.tsimport { 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_nameclicks、group_idall、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与eventCreate/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-pubsubiii-stream还提供stream条目变化时触发可用stream_name/group_id/item_id过滤、stream:join与stream:leaveWebSocket 订阅者接入/断开时触发stream:join处理器返回{ unauthorized: true }可在数据流出前拒绝订阅。存储后端kv 与 redisiii-stream的数据持久化与实时投递都交给可插拔的 adapterkv默认内置键值存储支持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-streamport/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_nameclicks group_idall第一条命令创建短链stream-me对应POST /linksHTTP 端点循环访问三次http://127.0.0.1:3111/s/stream-me每次访问都会触发重定向并记录一次点击iii trigger stream::list读取clicks流all组下的全部条目。每一次重定向落地linkworker 都会发布link.clickedclick-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),仅供参考