十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

用 iii-stream 构建免轮询的实时后端:三层流式数据模型、CRUD 函数面与响应式触发器实战指南

用 iii-stream 构建免轮询的实时后端:三层流式数据模型、CRUD 函数面与响应式触发器实战指南 用 iii-stream 构建免轮询的实时后端三层流式数据模型、CRUD 函数面与响应式触发器实战指南【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iiiiii-stream 是 iii 引擎内置的实时数据流 worker它以stream_name - group_id - item_id三层层级组织数据对外暴露 CRUD 形态的stream::*函数命名空间与三类响应式触发器stream、stream:join、stream:leave让开发者无需轮询即可搭建实时后端。读完本文你将掌握 iii-stream 的数据模型、8 个核心函数、触发器绑定与授权拦截机制、kv/redis 两种存储适配器的选型以及基于仓库源码的底层运行原理。iii-stream 是什么流式存储 响应式触发的一体化实时方案iii-stream的核心定位是durable real-time streams持久化实时流数据以stream_name流-group_id组-item_id条目的三层结构存入配置好的存储适配器同时暴露两个使用面CRUD 函数面以stream::*为命名空间的读写函数set/get/delete/list/update/send等供任何调用方通过函数触发机制读写数据响应式触发器面stream触发器在数据变更时自动触发stream:join/stream:leave在 WebSocket 订阅者连接/断开时触发。实时后端通过把 handler 绑定到这些触发器上实现而非轮询stream::list。一次写操作例如stream::set的完整流程是持久化新值 → 在 spawn 出的后台任务上分发匹配的stream触发器fire-and-forget调用方在 handler 执行完之前就拿到返回→ 向订阅了该(stream_name, group_id)的所有 WebSocket 客户端广播变更。这一流程可以在 stream.rs 的set/update/delete/send实现中逐行看到函数先通过适配器持久化再调用invoke_triggers分发触发器最后通过adapter.emit_event广播给订阅者。值得强调的是stream:join触发器的双重身份它同时是授权闸门——handler 返回{ unauthorized: true }就会在数据流出之前拒绝该订阅配合 worker 的auth_function配置后者在每个 WebSocket 握手上执行一次并把返回的context盖到每个 join/leave 事件上实现服务端订阅鉴权。适配器方面kv为默认内存或文件持久化仅限单实例无法跨进程扇出多实例集群需要实时广播时使用redis。浏览器与客户端订阅统一走 Browser SDKiii-browser-sdk它通过单条引擎 WebSocket 订阅stream变更直接连接 stream 端口ws://host:{port}/stream/{stream_name}/{group_id}的方式已废弃应以 Browser SDK 为准。何时使用与边界适用场景流数据变更需要触发其他地方的副作用派生投影、审计日志、通知且不希望轮询stream::list需要在服务端对 WebSocket 订阅做授权拦截而不是信任客户端过滤需要对订阅者连接/断开做服务端反应用于在线人数计数、按订阅限流、审计追踪需要并发写者不会竞争破坏的原子部分更新。边界与限制不是通用 KV 存储数据必须是流形态stream_name/group_id/item_idscope/key 这类简单值应使用独立的stateworker默认kv适配器不跨进程扇出多实例集群需要实时广播必须用redisstream:leave不是授权闸门它触发时订阅已经不存在返回值被忽略触发器 handler 在写操作返回后运行handler 失败既不会回滚写入也不会向调用方报错。部署与配置iii-stream是引擎内置 worker在受管的worker-compose.yaml中通过engine.workers.iii-stream配置。参考 README.md 与 iii.worker.yamlengine: workers: iii-stream: port: ${STREAM_PORT:3112} host: 0.0.0.0 adapter: name: redis config: redis_url: ${REDIS_URL:redis://localhost:6379}配置字段字段类型说明portnumber监听端口默认3112hoststring监听主机默认0.0.0.0auth_functionstring用于鉴权 WebSocket 连接的函数 IDadapterAdapter流存储与实时投递的后端适配器iii.worker.yaml中给出了内置默认值port: 3112、host: 127.0.0.1、adapter: kvstore_method: file_basedfile_path: ./data/stream_store。注意 config.rs 中的代码默认值是host: 0.0.0.0、port: 3112worker 声明文件里的127.0.0.1是示例部署采用的收窄值。适配器详解redis以 Redis 作为存储后端并用 Redis Pub/Sub 完成实时投递是多实例部署的唯一内置选择。name: redis config: redis_url: ${REDIS_URL:redis://localhost:6379}kv内置键值存储支持内存或文件持久化无外部依赖。name: kv config: store_method: file_based file_path: ./data/streams_store.dbkv 适配器配置见 config.rs 的KvAdapterConfig字段类型说明store_methodstringin_memory重启丢失或file_based持久化到磁盘file_pathstring文件存储目录路径save_interval_msnumber文件存储的落盘周期毫秒默认5000范围 1003,600,000channel_sizenumber进程内广播通道容量默认256范围 11,048,576事件突发超过慢订阅者消费速度时调大通道满会丢弃最旧的未投递事件此外源码中还内置了bridge适配器BridgeAdapterConfig.bridge_url默认ws://localhost:49134用于把流 pub/sub 转发到远端流后端。适配器以name为判别键的封闭oneOfschema 发布configuration::set会在写入时拒绝未知适配器名与未知字段这可以从 config.rs 的stream_adapter_schema及配套测试如拒绝name: postgres、store_method: weird、过期键url得到验证。运行时热重载iii-stream以iii-stream为 id 在内置configurationworker 中注册自己的配置因此无需重启引擎即可在运行时读取和修改上述字段例如configuration::set { id: iii-stream, value: { ... } }。engine.workers块只是首次启动的种子之后配置条目才是运行时真相来源运行时编辑在引擎重启后仍然生效值在 set 时按 schema 校验读取时展开${VAR:default}占位符。热重载的逐级生效语义详见 README.md 与 stream.rs 的apply_configauth_function对新连接立即生效无需重绑定只在客户端连接时被查询存量连接不受影响host/port触发监听器重绑定——新地址被绑定、WebSocket 服务器在新地址上重建、旧监听器拆除旧地址上的存量连接被断开客户端重连。重绑定有保护绑定失败的值会保留旧服务器继续运行。需要注意只改 host 不改 port且新旧接口重叠如0.0.0.0↔127.0.0.1时不会重绑定旧监听器仍占用端口导致新地址无法绑定此时仅记录日志并保留旧服务器要切换重叠接口需同时改端口或重启adapter触发完整的后端热切换——新 pub/sub 后端被构建并换入事件泵重启新连接使用新后端切换同样有保护无法构建后端则保留旧后端。存量连接在关闭前仍绑定旧后端因此不再收到新事件——多实例部署切换适配器应选择流量低谷。stream::* 函数面CRUD 与枚举八个函数构成了流数据的完整操作面其输入/输出结构定义在 structs.rs实现位于 stream.rs标注#[function(id stream::xxx)]函数参数返回说明stream::setstream_name,group_id,item_id,dataold_value,new_value持久化条目并广播 create/update 事件stream::updatestream_name,group_id,item_id,opsold_value,new_value用有序的set/merge/increment/decrement/append/remove操作原子更新现有条目stream::deletestream_name,group_id,item_idold_value删除条目并广播携带被删值的 delete 事件stream::sendstream_name,group_id,type,data,id(可选)—向组内订阅者广播瞬态事件不持久化如输入中状态、光标位置stream::getstream_name,group_id,item_idvalue按完整三元组读取单个条目stream::liststream_name,group_idgroup(any[])枚举组内所有条目stream::list_groupsstream_namegroups(string[])枚举流内所有组stream::list_all—stream(object[]),count枚举所有流的元数据从实现细节看set在条目已存在时广播Update事件、不存在时广播Create事件old_value.is_some()分支见 stream.rsdelete仅在旧值存在时广播update同样根据old_value决定Create/Update。send构造StreamOutboundMessage::Event { event: EventData { event_type, data } }直接广播完全跳过存储。每个写函数失败时返回对应的错误码STREAM_SET_ERROR、STREAM_GET_ERROR、STREAM_DELETE_ERROR、STREAM_UPDATE_ERROR、STREAM_SEND_ERROR、STREAM_LIST_*_ERROR等这些错误码在 stream.rs 的test_stream_methods_return_adapter_failures测试中被逐项断言。原子更新stream::update与 UpdateOpops是UpdateOp的有序列表定义在 sdk/packages/rust/helpers/src/stream.rs操作字段语义setpath,value在 path 处覆盖写值mergepath(可选),value把对象合并进现有值仅对象path 可省略根合并、为单个一级键、或为嵌套段的字面段数组incrementpath,by数值自增decrementpath,by数值自减appendpath(可选),value向数组追加元素或拼接字符串path 可省略根追加、单个一级键或嵌套段removepath移除字段原子性保证了并发写者不会互相覆盖多个操作按给定顺序一次性应用到现有值上。测试test_stream_module_update_existing_record/test_stream_module_update_new_recordstream.rs验证了更新已有条目广播Update、更新不存在条目则创建并广播Create的行为。函数级覆盖自定义实现一个容易被忽视的扩展点对特定流注册stream::set({stream_name})这类带流名后缀的函数即可覆盖该流的内置实现——源码在调用内置逻辑前先按format!(stream::set({stream_name}))查函数注册表命中则改调自定义 handler。test_stream_custom_functions_override_adapter测试完整验证了get/list/list_groups/set/delete/update六个函数的覆盖路径。响应式触发器数据变更与订阅生命周期绑定stream家族触发器让函数在流活动发生时自动运行——无需轮询stream::list。三种触发器覆盖两类关注点stream响应数据变更stream:join/stream:leave响应 WebSocket 订阅者生命周期。适用场景set/update/delete/send应驱动投影、审计日志或下游通知订阅需要服务端授权stream:join或成对的建立/拆除逻辑stream:joinstream:leave。如果只是按需读取当前值应调用stream::get而非绑定触发器。绑定方式注册 handleriii.registerFunction(presence::on-change, handler)注册触发器iii.registerTrigger({ type: stream, function_id: presence::on-change, config: { stream_name: presence, // stream 必填worker 以它建触发器索引。 group_id: room-1, // 可选。留空/省略 匹配所有组。 item_id: user-123, // 可选。留空/省略 匹配所有条目。 // condition_function_id 同样受支持。 }, })配置语义见 trigger.rs 的StreamTriggerConfig与invoke_triggers的匹配逻辑stream要求非空stream_name注册时触发器按stream_name建立stream_name - trigger_ids索引触发时先按流名查索引再对候选触发器做group_id/item_id精确匹配空值视为通配stream:join与stream:leave不带任何配置字段在 handler 内按事件内容分支即可stream:join返回{ unauthorized: true }拒绝订阅其他任何返回都放行非布尔值如nope会被忽略并视为放行connection.rs 的handle_join_leave及handle_join_leave_ignores_non_boolean_unauthorized_values测试condition_function_id条件函数返回false时跳过 handler返回None视为通过条件函数执行出错也会跳过该 handlerinvoke_triggers中check_condition的处理见 stream.rs 的test_stream_invoke_triggers_covers_conditions_and_errors。从stream_trigger_fires_target_and_condition_in_registering_namespace测试可以看出触发器目标与条件函数都在注册时所属的 namespace中解析执行。触发stream的写操作stream::set、stream::update、stream::delete、stream::send读操作不触发。事件载荷的准确字段可通过iii get function info查询对应触发器类型或 handler 函数 id 获得。参考 README.md两种触发器的载荷字段为streamtypecreate/update/delete、timestamp、streamName、groupId、id、event含type与data的对象stream:join/stream:leavesubscription_id、stream_name、group_id、id可选、context来自 auth。底层消息结构可在 structs.rs 找到StreamWrapperMessagetype/timestamp/streamName/groupId/id/event与StreamOutboundMessage的Unauthorized/Sync/Create/Update/Delete/Event六个变体。触发器与写入的时序关系触发器 handler 通过tokio::spawn在独立任务中执行fire-and-forget原始写调用先返回handler 失败既不会回滚已持久化的写入也不会向调用方报错。这一设计在 stream.rs 的invoke_triggers中体现事件序列化后立即 spawn 任务任务内部对每个匹配触发器执行条件检查与 handler 调用错误仅记录日志与 OTel 状态码。另外针对iii:devtools:*前缀的观测类流触发器投递刻意不生成 trace span避免观测管道自我观测形成无限循环源码注释对此有专门说明。实战实时在线状态Real-Time Presence把数据面与触发器面结合起来一个典型的在线状态场景如下示例来自 README.mdimport { registerWorker, TriggerAction } from iii-sdk const iii registerWorker(ws://localhost:49134) // 写入在线状态 iii.trigger({ function_id: stream::set, payload: { stream_name: presence, group_id: room-1, item_id: user-123, data: { name: Alice, online: true, lastSeen: new Date().toISOString() }, }, action: TriggerAction.Void(), }) // 读取单个用户 const user await iii.trigger({ function_id: stream::get, payload: { stream_name: presence, group_id: room-1, item_id: user-123 }, }) // 列出房间所有成员 const roomMembers await iii.trigger({ function_id: stream::list, payload: { stream_name: presence, group_id: room-1 }, })再叠加订阅生命周期触发器做在线人数统计与授权const fn iii.registerFunction(onJoin, (input) { console.log(Joined ${input.stream_name}/${input.group_id}, input.context) return {} }) iii.registerTrigger({ type: stream:join, function_id: fn.id, config: {}, })stream:join的context来自auth_function的返回。定义鉴权函数接收headers、path、query_params、addr返回{ context: ... }并在配置中指定- name: iii-stream config: auth_function: onAuthTypeScriptiii.registerFunction(onAuth, (input) ({ context: { name: John Doe }, }))Pythondef on_auth(input): return {context: {name: John Doe}} iii.register_function(onAuth, on_auth)鉴权输入与上下文的字段结构对应 structs.rs 的StreamAuthInputheaders/path/query_params/addr与StreamAuthContextcontext。握手时引擎调用auth_function成功解析出 context 后随连接进入 socket 处理stream.rs 的ws_handler。客户端订阅浏览器与客户端订阅使用 Browser SDKiii-browser-sdk它通过单条引擎 WebSocket 订阅stream变更并在每个变更事件上重新渲染连接由 iii-worker-manager 的 RBAC 监听器把关。直接连接 stream 端口ws://host:3112/stream/stream_name/group_id/已废弃应统一走 Browser SDK。源码导航官方功能说明与实战示例engine/src/workers/stream/README.md本技能文档Agent 技能入口engine/src/workers/stream/skills/SKILL.mdworker 声明与默认配置engine/src/workers/stream/iii.worker.yaml配置结构、默认值、适配器 schema 与热重载语义engine/src/workers/stream/config.rs八个函数实现、触发器分发、热应用与单元测试engine/src/workers/stream/stream.rs触发器注册/注销与配置解析engine/src/workers/stream/trigger.rs输入输出结构、消息协议与授权上下文engine/src/workers/stream/structs.rs订阅生命周期、join/leave 授权闸门实现engine/src/workers/stream/connection.rs原子更新操作符 UpdateOp 定义sdk/packages/rust/helpers/src/stream.rs【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表