iii Channels 完全实战:在 Worker 之间流式传输大文件与二进制数据
发布时间:2026/9/13 10:24:32 作者:尧图编辑部 阅读量:1,286

iii Channels 完全实战在 Worker 之间流式传输大文件与二进制数据【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii本篇技术指南围绕 iii 的 Channels数据通道机制展开当函数调用需要传递文件、图片、数据集或流式音视频等大载荷时如何通过 Channel 把“协调”与“数据搬运”解耦避开 JSON 触发载荷的 WebSocket 尺寸限制。读完后你将掌握 Node / Python / Rust 三种 SDK 下创建通道、读写数据、跨函数传递 channel ref 的完整操作并能结合 iii 引擎与 SDK 源码理解通道寻址、缓冲与背压的实现原理。什么时候该用 Channel16 MB 载荷分界线Channels 让大载荷或二进制数据在 iii 各 worker 之间流动而不必把数据塞进 JSON 函数载荷。当载荷满足以下任一条件时应改用 Channel预期为大文件文件、图片、数据集体积超过约16 MB需要以流式方式传输音频、视频长时间运行的任务希望获得增量式进度输出。对于小型 JSON 结构化事件继续使用普通的worker.trigger(...)调用即可。16 MB 上限从何而来iii 本身并不强制最大 trigger 载荷尺寸。真正的上限来自引擎和调用方 SDK 各自使用的 WebSocket 库——它们各自有每帧/每消息大小的默认值其中最小的常见每帧默认值约在 16 MB 量级这就是应当切换到 Channel 的实用分界线。当前 iii 依赖的库摘自 docs/0-16-0/creating-workers/channels.mdx组件WebSocket 库相关尺寸配置项Engineaxum基于 tokio-tungsteniteWebSocketConfig的max_frame_size/max_message_sizeNode SDKwsmaxPayload选项Python SDKwebsocketsmax_size选项Rust SDKtokio-tungsteniteWebSocketConfig与引擎相同这些是库的默认值会随依赖版本变化iii 不承诺硬性上限。把 16 MB 当作安全默认阈值即可。Channel 的模型Writer、Reader 与 Ref从 Channels 架构文档 看Channel 是引擎管理的、基于 WebSocket 的流式管道核心概念如下概念作用Channel由引擎管理的 WebSocket 管道Writer向管道写入字节或文本消息的一端Reader从管道接收字节或文本消息的一端Ref可序列化的小令牌通过trigger()载荷传递关键思想是ref 走普通的函数调用数据本身走通道。函数调用负责协调工作通道负责承载数据流引擎负责追踪与路由。这解释了为什么 JSON 触发调用适合结构化事件却不适合大文件、媒体、流式响应如 agent/聊天输出或长任务的局部输出。引擎侧的实现ChannelManager从源码结构看引擎在 engine/src/workers/worker/channels.rs 中实现了通道管理StreamChannel持有idUUID、access_keyUUID、一对mpsc通道端tx/rx、属主 worker id 与创建时间ChannelManager::create_channel(buffer_size, owner_worker_id)会强制buffer_size至少为 1用tokio::sync::mpsc建立缓冲管道并分别构造方向为Write与Read的两个StreamChannelRef返回常量CHANNEL_TTL为 5 分钟该文件 L33 附近即长期未用的通道会被清理worker 断开时remove_channels_by_worker会移除它名下所有通道这与架构文档中“worker 断开时其通道连接随之关闭”的描述一致。也就是说一次createChannel在引擎内就是生成 UUID 通道 id access key → 在DashMap中登记 → 返回带方向标记的两个 ref。本地端 API创建通道一个 Channel 由某个 worker 创建包含两个本地流端writer与reader外加两个可序列化 refwriterRef/readerRef可交给另一个函数。createChannel/create_channel助手函数从各 SDK 的helpers模块导入以 worker 客户端作为第一个参数。Node / TypeScriptimport { createChannel } from iii-sdk/helpers; const channel await createChannel(worker); // channel.writer // channel.reader // channel.writerRef // channel.readerRefNode SDK 中createChannel(iii, bufferSize?)是一个自由函数见 sdk/packages/node/iii/src/helpers.ts内部转发给IIIClient的实现最终通过触发引擎内置函数engine::channels::create完成创建见 sdk/packages/node/iii/src/iii.ts 中的__helpers_create_channel。Pythonfrom iii.helpers import create_channel channel create_channel(iii_client) # channel.writer # channel.reader # channel.writer_ref # channel.reader_refRustuse iii_sdk::helpers::create_channel; let channel create_channel(worker, None).await?; // channel.writer // channel.reader // channel.writer_ref // channel.reader_refRust 的create_channel(iii, buffer_size)第二个参数是Optionusize对应引擎侧mpsc管道的缓冲容量见 sdk/packages/rust/iii/src/helpers.rs传入None则使用默认缓冲。向通道写入数据把载荷写入本地writer完成后关闭它。字节经由引擎流向持有匹配reader的 worker。Node / TypeScriptimport { createChannel } from iii-sdk/helpers; const channel await createChannel(worker); channel.writer.stream.end(Buffer.from(file contents));Pythonfrom iii.helpers import create_channel_async channel await create_channel_async(iii_client) await channel.writer.write(bfile contents) await channel.writer.close_async()Rustuse iii_sdk::helpers::create_channel; let channel create_channel(worker, None).await?; channel.writer.write(bfile contents).await?; channel.writer.close().await?;Node 端写入的实现细节从 sdk/packages/node/iii/src/channels.ts 源码看Node 的ChannelWriter有几个值得注意的工程细节64 KB 分帧FRAME_SIZE 64 * 1024L57sendChunked把每次写入的大块数据切成 64 KB 分片顺序发送避免单帧过大撞上 WebSocket 帧上限懒连接WebSocket 在首次真正写入时才建立ensureConnected创建通道本身不建立数据流连接这与架构文档“通道流是懒连接”的描述一致断连重放队列发送失败或 socket 死亡时消息进入pendingMessages在下次open时冲刷避免丢消息关闭时序保护final回调中延迟约 10 ms 再发 close 帧防止 close 帧先于已缓冲的数据帧到达引擎造成数据截断。writer 连接引擎使用的 URL 由buildChannelUrl构造L299-L307{engineWsBase}/ws/channels/{channelId}?key{accessKey}dirwrite|read——channel_id与access_key就是引擎侧StreamChannelRef的两个字段dir区分读写方向。从通道读取数据读取本地reader直到另一端关闭。字节按持有匹配writer的 worker 的写入顺序到达。Node / TypeScriptimport { createChannel } from iii-sdk/helpers; const channel await createChannel(worker); let bytes 0; for await (const chunk of channel.reader.stream) { bytes Buffer.isBuffer(chunk) ? chunk.length : Buffer.byteLength(chunk); }Pythonfrom iii.helpers import create_channel_async channel await create_channel_async(iii_client) bytes_total 0 async for chunk in channel.reader: bytes_total len(chunk)Rustuse iii_sdk::helpers::create_channel; let channel create_channel(worker, None).await?; let mut bytes 0; while let Some(chunk) channel.reader.next_binary().await? { bytes chunk.len(); }Node 的ChannelReader同样懒连接首次read()时才建 socket并内置背压处理ws.on(message)中若stream.push(data)返回false下游读不动了就调用ws.pause()暂停 WebSocket等read()恢复后再resume()——这正是架构文档所说“背压由 SDK 流实现处理writer 可在 reader 跟不上时暂停”的具体实现。跨函数使用 Channel传递 ref通道只有当两端被不同代码路径持有时才有价值通常是一个函数持有本地writer/reader另一个函数在触发载荷中收到匹配的 ref。把 channel ref 随函数调用传过去把readerRef或writerRef作为普通函数调用的一部分传递。接收方用 ref 去读或写通道。Node / TypeScriptconst result await worker.trigger({ function_id: files::process, payload: { filename: report.csv, reader: channel.readerRef, }, });Pythonresult await iii_client.trigger_async({ function_id: files::process, payload: { filename: report.csv, reader: channel.reader_ref.model_dump(), }, })Rustuse iii_sdk::TriggerRequest; use serde_json::json; let result worker .trigger(TriggerRequest { function_id: files::process.to_string(), payload: json!({ filename: report.csv, reader: channel.reader_ref, }), action: None, timeout_ms: None, }) .await?;三种 SDK 对接收侧的处理有差异Node 和 Python 会在 handler 运行前把传入的 channel ref 反序列化成本地可用的ChannelReader/ChannelWriter对象ref 到达 handler 时已经是活的流Rust 则在 JSON 中收到 ref需要显式重建reader 或 writer。在接收函数中从 ref 读取Node / TypeScriptimport type { ChannelReader } from iii-sdk; worker.registerFunction(files::process, async (input: { reader: ChannelReader }) { let bytes 0; for await (const chunk of input.reader.stream) { bytes Buffer.isBuffer(chunk) ? chunk.length : Buffer.byteLength(chunk); } return { bytes }; });Pythonasync def process_file(input: dict) - dict: reader input[reader] total 0 async for chunk in reader: total len(chunk) return {bytes: total} worker.register_function(files::process, process_file)Rustuse iii_sdk::helpers::{ChannelDirection, extract_channel_refs}; use iii_sdk::{ChannelReader, IIIError}; use serde_json::json; let refs extract_channel_refs(input); let (_, reader_ref) refs .iter() .find(|(k, r)| k reader matches!(r.direction, ChannelDirection::Read)) .ok_or_else(|| IIIError::Handler(missing reader channel ref.into()))?; let reader ChannelReader::new(worker.address(), reader_ref); let mut bytes 0; while let Some(chunk) reader.next_binary().await? { bytes chunk.len(); } Ok(json!({ bytes: bytes }))在接收函数中向 ref 写入Node / TypeScriptimport type { ChannelWriter } from iii-sdk; worker.registerFunction(files::generate, async (input: { writer: ChannelWriter }) { input.writer.stream.write(Buffer.from(hello )); input.writer.stream.end(Buffer.from(world)); return { ok: true }; });Pythonasync def generate_file(input: dict) - dict: writer input[writer] await writer.write(bhello ) await writer.write(bworld) await writer.close_async() return {ok: True} worker.register_function(files::generate, generate_file)Rustuse iii_sdk::helpers::{ChannelDirection, extract_channel_refs}; use iii_sdk::{ChannelWriter, IIIError}; use serde_json::json; let refs extract_channel_refs(input); let (_, writer_ref) refs .iter() .find(|(k, r)| k writer matches!(r.direction, ChannelDirection::Write)) .ok_or_else(|| IIIError::Handler(missing writer channel ref.into()))?; let writer ChannelWriter::new(worker.address(), writer_ref); writer.write(bhello ).await?; writer.write(bworld).await?; writer.close().await?; Ok(json!({ ok: true }))Rust 侧extract_channel_refs负责从 handler 入参中识别出通道 refChannelReader::new/ChannelWriter::new则用 worker 地址 ref 重建出可收发的流对象。生命周期、背压与双向通信从 架构文档 与源码交叉印证可以总结通道运行时的完整行为懒连接createChannel只分配 refWebSocket 数据流在某一侧开始读或写时才建立。Node 端ChannelWriter.ensureConnected()与ChannelReader.ensureConnected()都遵循这一模式背压SDK 流实现负责背压Node 端为push返回false时暂停 socketwriter 可因 reader 跟不上而暂停引擎侧mpsc管道的buffer_size提供额外的缓冲容量关闭传播writer 关闭后reader 收到流结束Node 端ws.on(close)时向Readable流 pushnull断连清理worker 断开时引擎移除其名下全部通道remove_channels_by_worker并受 5 分钟CHANNEL_TTL兜底清理双向通信需要双向时创建两个通道每个方向各一个——单个通道本身是单向的dirwrite/dirread的连接是独立的。各语言的完整 Channel API 清单sendMessage/onMessage文本消息、readAll便捷读取、ChannelDirection、extractChannelRefs等可参考该文档版本docs/0-16-0/sdk-reference/目录下的 Node、Python、Rust SDK 参考。测试如何验证这些行为三种 SDK 各自带有针对数据通道的集成测试可作为行为依据Nodesdk/packages/node/iii/tests/data-channels.test.tsPythonsdk/packages/python/iii/tests/test_data_channels.pyRustsdk/packages/rust/iii/tests/data_channels.rs此外 Rust 的 helpers 测试 覆盖create_channel助手函数本身可作为 API 签名的可执行文档。小结载荷超过约 16 MB、二进制流式数据、或需要增量进度输出时用 Channel 代替 JSON 触发载荷心智模型trigger()传 ref协调Channel 传字节数据引擎负责路由、追踪与生命周期Node / Python 的 ref 会在 handler 前被自动物化为活流Rust 需显式extract_channel_refsChannelReader::new/ChannelWriter::new重建源码层面通道的 64 KB 分帧、懒连接、背压暂停、5 分钟 TTL 与断连清理等细节均可在 engine/src/workers/worker/channels.rs 与 sdk/packages/node/iii/src/channels.ts 中逐一对照验证。【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考