实战指南:跨 Worker 流式传输大文件与二进制数据)
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 是运行在引擎上、由 WebSocket 支撑的流式管道用于在 Worker 之间传输大体积或二进制数据而无需把这些数据塞进 JSON 函数载荷function payload。当你的载荷预期超过约16 MB文件、图片、数据集、需要流式传输音频、视频或希望在长任务中持续输出增量进度时就应该使用 Channel小体积 JSON 请继续使用常规的worker.trigger(...)调用。读完本文你将掌握 iii 各语言 SDKNode/TypeScript、Python、Rust中创建、写入、读取 Channel 的完整 API以及如何通过可序列化的 channel ref 把流的另一端交给另一个函数实现跨进程、跨语言的实时数据传输。Channel 是什么协调与数据传输的分离在 iii 中一次函数调用本质上是一条 JSON 消息——这对结构化事件和命令载荷非常合适但对大文件、媒体、流式响应Agent、聊天以及长任务的部分输出来说就成了一种负担。Channel 将协调与数据传输拆开一次函数调用负责协调工作一个 Channel 负责承载数据流引擎负责追踪tracing与路由。一个 Channel 由某个 Worker 创建拥有两个本地流端点writer和reader以及两个可序列化的引用writerRef/readerRef引用可以放进另一个函数的载荷中交接给其他代码路径。核心模型如下表所示详见 docs/0-19-0/understanding-iii/channels.mdx概念作用Channel由引擎管理、基于 WebSocket 的管道Writer向管道发送字节或文本消息Reader从管道接收字节或文本消息Ref通过trigger()载荷传递的小型可序列化令牌关键点在于ref 走常规函数调用而数据本身走 Channel。两个函数甚至可以运行在不同进程、不同语言中。典型的运行时流程是生产者函数调用createChannel()拿到writer、reader、writerRef、readerRef把readerRef通过trigger()交给消费者函数消费者持有 ref 后建立自己的读写端生产者写入数据块消费者实时读取最后生产者关闭 writer消费者端 reader 收到流结束信号可参考 docs/0-19-0/understanding-iii/channels.mdx 中的 sequenceDiagram。如果需要双向通信就创建两个 Channel一个方向一个。载荷上限为什么是 16 MBiii 自身不强制trigger 载荷的最大值实际上限来自引擎与各 SDK 所依赖的 WebSocket 库各自的默认配置。这些库对单帧per-frame和单消息per-message大小有各自的默认值其中最小的常见单帧默认值约为16 MB这就是切换到 Channel 的实际分界线。iii 当前依赖的库均为库默认值会随依赖版本变化iii 不发布硬性保证的上限16 MB 是安全默认线引擎基于 [axum] [tokio-tungstenite]对应WebSocketConfig的max_frame_size与max_message_sizeNode SDK[ws]对应maxPayload选项Python SDK[websockets]对应max_size选项Rust SDK[tokio-tungstenite]与引擎相同见WebSocketConfig。从引擎侧源码看Worker 监听器把 Channel 的 WebSocket 端点挂在/ws/channels/{channel_id}见 engine/src/workers/worker/mod.rs 中的路由挂载与 engine/src/workers/worker/README.md 的说明。同时数据在通道上以帧为单位分片发送Node SDK 的ChannelWriter使用FRAME_SIZE 64 * 102464 KB切分大块数据sdk/packages/node/iii/src/channels.tsPython SDK 与 Rust SDK 也分别用MAX_FRAME_SIZE 64 * 1024做同样的分片sdk/packages/python/iii/src/iii/channels.py、sdk/packages/rust/iii/src/channels.rs。这意味着 16 MB 的分界线更多取决于 WebSocket 库的默认单消息上限而实际写入时 SDK 会自动分片成 64 KB 的帧。使用 Channel本地端点 API一个 Channel 由某个 Worker 创建拥有两个本地流端点writer、reader和两个可序列化 refwriterRef、readerRef。下面依次覆盖本地端 API创建、写入、读取。创建 Channelworker.createChannel()返回一个包含两个本地流对象和两个可序列化 ref 的 Channelwriter与reader是本地流端点writerRef/readerRef是你要传给另一个函数的令牌对方凭它读写另一端。Node / TypeScriptconst channel await worker.createChannel(); // channel.writer // channel.reader // channel.writerRef // channel.readerRefPythonchannel iii_client.create_channel() # channel.writer # channel.reader # channel.writer_ref # channel.reader_refRustlet channel worker.create_channel(None).await?; // channel.writer // channel.reader // channel.writer_ref // channel.reader_ref从引擎侧看create_channel会生成一个channel_id与一个access_key均为 UUID并用 tokio 的mpsc::channel(buffer_size)建立内部缓冲管道随后构造方向分别为Write与Read的两个StreamChannelRef见 engine/src/workers/worker/channels.rs 的ChannelManager::create_channel。channel ref 在 Rust SDK 中对应的结构是StreamChannelRef { channel_id, access_key, direction }sdk/packages/rust/iii/src/channels.rs其中access_key是访问该 Channel WebSocket 的独立能力令牌。向 Channel 写入把载荷写入本地writer写完后关闭它。字节会流经引擎到达持有匹配reader的那个 Worker。Node / TypeScriptconst channel await worker.createChannel(); channel.writer.stream.end(Buffer.from(file contents));Pythonchannel await iii_client.create_channel_async() await channel.writer.write(bfile contents) await channel.writer.close_async()Rustlet channel worker.create_channel(None).await?; channel.writer.write(bfile contents).await?; channel.writer.close().await?;值得注意的底层细节Node 的ChannelWriter在final回调中会延迟约 10 ms 再发送 close 帧让 TCP 栈先把缓冲的数据帧冲刷出去否则 close 帧可能先于数据帧到达引擎造成数据截断见 sdk/packages/node/iii/src/channels.ts 中final与doClose的注释Python 的close_async同样先asyncio.sleep(0.01)再关闭连接sdk/packages/python/iii/src/iii/channels.pyRust 的ChannelWriter::close也是先tokio::time::sleep(10ms)再发送WsMessage::Close并有对应的单元测试验证 close 至少等待 10 mssdk/packages/rust/iii/src/channels.rs。从 Channel 读取读取本地reader直到另一端关闭。字节按持有匹配writer的那个 Worker 写入的顺序到达。Node / TypeScriptconst channel await worker.createChannel(); let bytes 0; for await (const chunk of channel.reader.stream) { bytes Buffer.isBuffer(chunk) ? chunk.length : Buffer.byteLength(chunk); }Pythonchannel await iii_client.create_channel_async() bytes_total 0 async for chunk in channel.reader: bytes_total len(chunk)Rustlet channel worker.create_channel(None).await?; let mut bytes 0; while let Some(chunk) channel.reader.next_binary().await? { bytes chunk.len(); }在 Node 实现中ChannelReader把 WebSocket 收到的二进制帧 push 进一个Readable流当流背压push 返回 false时会 pause 底层 WebSocket连接关闭时 pushnull结束流sdk/packages/node/iii/src/channels.ts。这正对应架构文档中所说的Channel 流是惰性连接的——创建 Channel 只分配 refWebSocket 流在某一端开始读写时才真正连接背压由 SDK 的流实现处理让写入方在读取方跟不上时暂停。跨函数使用 Channelref 的交接一个 Channel 只有两端被不同代码路径持有时才真正发挥作用。典型场景是一个函数持有本地writer/reader另一个函数在自己的载荷中收到匹配的 ref。下面覆盖交接的两个环节如何把 ref 随常规 trigger 调用一起发送以及接收方如何把它还原成可读写的数据流。把 Channel ref 传给另一个函数把readerRef或writerRef作为普通函数调用的一部分传入。接收函数用 ref 去读取或写入这个 Channel。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?;Node 和 Python 会在 handler 运行之前把收到的 channel ref 反序列化成活的ChannelReader/ChannelWriter对象所以 ref 一到手就可以直接迭代或写入。Rust 则是在 JSON 中接收 ref再用ChannelReader::new(...)或ChannelWriter::new(...)显式重建读写端——两者都接受引擎 WebSocket 基址与StreamChannelRefsdk/packages/rust/iii/src/channels.rs。Rust SDK 还提供了extract_channel_refs(input)辅助函数递归地从 JSON 顶层字段含嵌套对象与数组中找出所有形如StreamChannelRef的值并返回其字段路径与反序列化结果。从 Channel 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::{ChannelDirection, ChannelReader, IIIError}; use serde_json::json; let refs iii_sdk::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 }))向 Channel 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::{ChannelDirection, ChannelWriter, IIIError}; use serde_json::json; let refs iii_sdk::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 }))生命周期、背压与安全要点惰性连接创建 Channel 只分配 ref 与内部缓冲WebSocket 连接在某一端开始读写时才建立这一点在引擎侧同样成立——StreamChannel内部持有 tokio 的mpsc::Sender/Receiver由ChannelManager用DashMap按channel_id管理engine/src/workers/worker/channels.rs。背压由各 SDK 的流实现负责。Node 端在Readable背压时 pause WebSocket写入方在读取方跟不上时会暂停。流结束与清理writer 关闭后 reader 收到流结束Worker 断连时其 Channel 连接一并关闭。引擎侧 Channel 设有 TTLCHANNEL_TTL: Duration::from_secs(5 * 60)5 分钟过期资源会被清理engine/src/workers/worker/channels.rs。访问控制Channel WebSocket 的访问由每个StreamChannelRef中携带的access_key能力令牌独立校验Channel 端点挂在 Worker 监听器所在端口/ws/channels/{channel_id}SDK Worker 无需额外配置即可使用createChannel()engine/src/workers/worker/README.md。双向通信Channel 是单向的如需双向创建两个 Channel各管一个方向。各语言完整的 Channel API 表面可继续查阅 SDK 参考文档Node SDK、Python SDK、Rust SDK以及架构说明 docs/0-19-0/understanding-iii/channels.mdx。【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考