ARTICLE DETAIL

资讯详情

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

iii 中的 Channels:跨 Worker 流式传输大负载与二进制数据的实战指南

iii 中的 Channels:跨 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流式通道机制展开它让不同进程、甚至不同语言Node/Python/Rust的 worker 之间能够以字节流方式传输文件、图片、音视频等大负载或二进制数据而无需把数据塞进 JSON 函数 payload。读完本文你将掌握 channel 的创建、本地读写的三语言 API、跨函数传递readerRef/writerRef的完整流程并能结合引擎与 SDK 源码理解其 WebSocket 寻址、鉴权、背压与生命周期回收的底层实现。一、为什么需要 Channels16 MB 的实用分界线函数调用worker.trigger(...)本质是 JSON 消息非常适合结构化事件与命令载荷但对大文件、媒体、流式响应Agent、聊天和长任务的增量输出JSON 是错误的承载方式。Channels 将协调与数据传输分离一次函数调用负责协调工作一条 channel 承载字节流引擎负责路由、追踪与生命周期管理。什么时候该用 channel载荷预期较大文件、图片、数据集超过约16 MB意图是流式传输音频、视频或者希望在长任务中获得增量进度更新。小 JSON 载荷仍然应该直接用普通的worker.trigger(...)调用。16 MB 上限从哪来iii 本身不强制最大 trigger 载荷大小实际上限来自引擎与调用方 SDK 所用 WebSocket 库的每帧/每消息默认值各库默认值不同其中最小的公共每帧默认值约在 16 MB这就是切换到 channel 的实用分界线。当前 iii 依赖的库包括引擎axum tokio-tungstenite参考tungstenite的WebSocketConfig中max_frame_size与max_message_sizeNode SDKws库的maxPayload选项Python SDKwebsockets库的max_size选项Rust SDKtokio-tungstenite与引擎相同。这些是库的默认值会随依赖版本变化iii 不公布硬性保证上限因此 16 MB 是一个安全的经验默认值。二、Channel 的模型流端与可序列化引用一个 channel 由某个 worker 创建包含两个本地流端writer与reader和两个可序列化引用writerRef/readerRef引用可以随 trigger payload 交给另一个函数。核心概念如下概念作用Channel由引擎管理的、基于 WebSocket 的管道Writer向管道发送字节或文本消息的一端Reader从管道接收字节或文本消息的一端Ref通过trigger()载荷传递的小型可序列化令牌关键思想是引用ref通过普通函数调用传递而数据本身走 channel。运行时流程大致为生产方函数调用createChannel()引擎返回writer、reader、writerRef、readerRef生产方通过trigger把readerRef传给消费方函数生产方向本地writer写入分片引擎转发给持有reader的消费方生产方关闭writer消费方reader收到流结束。三、本地端 API创建、写入、读取创建 channelcreateChannel()返回一个 channel包含两个本地流对象和两个可序列化引用writer/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向 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?;从 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(); }四、跨函数使用 channel引用传递与流重建只有当两端由不同代码路径持有时channel 才真正发挥作用通常是一个函数持有本地writer/reader另一个函数在 payload 中收到匹配的 ref。把 channel 引用传给另一个函数把readerRef或writerRef作为普通函数调用 payload 的一部分传出去接收方用该 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 引用反序列化为活的ChannelReader/ChannelWriter对象ref 到达即可迭代或写入Rust 则以 JSON 形式接收 ref需要显式用ChannelReader::new(...)或ChannelWriter::new(...)重建。从 channel 引用读取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 引用写入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 }))五、引擎侧实现ChannelManager 与 WebSocket 路由上面的 API 行为在引擎源码中有清晰的对应。核心实现在 ChannelManager 与 StreamChannel创建即生成一对引用create_channel(buffer_size, owner_worker_id)生成 UUID 形式的channel_id与access_key构造一个tokio::sync::mpsc通道缓冲区大小会被max(1)钳制见 buffer 钳制测试并登记owner_worker_id与创建时间。随后按方向返回StreamChannelRef { channel_id, access_key, direction }两个引用——Write方向给writerRefRead方向给readerRef。这与 Node SDK 的ChannelDirection常量定义 和 Rust SDK 的StreamChannelRef字段完全一致。一次性独占取端take_sender/take_receiver通过MutexOption...::take()保证每个端点只能被连接一次第二次取端返回None见 take_sender 独占测试。这解释了为什么 ref 是能力令牌引用里的access_key不对就取不到端点见 错误密钥返回 None 的测试。鉴权在 WebSocket 升级时执行channel 端点挂载在每个 worker listener 的/ws/channels/{channel_id}路由上见 路由注册升级处理函数channel_ws_upgradews_handler.rs按dir查询参数决定调用take_receiver还是take_senderchannel 不存在时返回 404见 404 行为测试。worker 模块 README 也说明了channels::create属于 RBAC 的 always-allowed 基础设施豁免即使expose_functions为空也能创建 channel而通道数据面的访问则完全由access_key能力令牌独立校验。孤儿回收每条 channel 有 5 分钟 TTLCHANNEL_TTL见 常量定义。后台清扫任务每 60 秒运行一次sweep_stale_channels只删除至少有一端从未被取走的过期 channel两端都已连接在正常使用中的 channel 会被保留到 ws_handler 在流结束后清理见 清扫保留策略测试。worker 断开连接时remove_channels_by_worker会按owner_worker_id清掉它名下的全部 channel——这对应文档中当 worker 断开连接时其 channel 连接随之关闭的行为。数据面保序通道内部是 mpsc 队列ChannelItem::Text/ChannelItem::Binary两类条目按发送顺序进入接收端见 文本/二进制收发测试印证了字节按写入顺序到达的语义。六、SDK 侧实现分片、背压与连接细节Node SDK64 KB 分片与关闭帧延迟Node SDK 的 ChannelWriter 有几个值得注意的工程细节分帧发送FRAME_SIZE 64 * 1024sendChunked把大 buffer 切成 64 KB 片段顺序发送分片实现从而规避底层ws库的每帧大小默认值这也是单条消息可超过 16 MB、但每帧仍在默认限额内的关键。关闭帧延迟final回调中 close frame 会延迟 10 ms 发出延迟关闭实现注释明确说明这是为了让 TCP 栈有时间 flush 所有缓冲的send()数据否则 close frame 可能先于数据帧到达引擎造成数据截断。断连重发队列sendRaw信任 socket 实时的readyState而非内部标志socket 未就绪时把消息压入pendingMessages在下次open时冲刷避免静默丢包sendRaw 实现。读端背压ChannelReader收到二进制帧后stream.push(data)若返回false下游消费不过来则ws.pause()暂停 WebSocket 接收背压实现这正是背压由 SDK 流实现处理写入方可在读取方跟不上时暂停的落地。URL 构造读/写端各自连接/ws/channels/{channelId}?key{accessKey}dir{read|write}buildChannelUrl且连接是懒建立的——创建 channel 只分配 refWebSocket 在某一侧真正开始读或写时才连接。此外Node SDK 提供sendMessage/onMessage走文本消息通道适合在字节流之外发送结构化文本如进度事件见 ChannelWriter.sendMessage 文档注释。Rust SDK显式重建引用Rust SDK 的 channels 模块 定义了ChannelDirection枚举、StreamChannelRef结构体以及ChannelReader::new/ChannelWriter::new构造入口 与 L144读取循环用next_binary()返回OptionVecu8实现。extract_channel_refs实现负责从 handler 的 JSON 输入中扫描出所有 channel 引用返回(key, StreamChannelRef)列表——这与上文 Rust 示例中的用法一一对应。由于 Rust 不做隐式物化handler 必须显式按 key 和方向筛选并构造 reader/writer。Python SDKPython 端实现在 sdk/packages/python/iii/src/channels.py 与 channel.pycreate_channel/create_channel_async返回的 channel 对象上writer/reader支持write、close_asyncreader 支持async for迭代。与 Node 一致Python 在 handler 执行前将 payload 中的 ref 物化为活的 reader/writer 对象。三语言 SDK 的行为均有对应测试覆盖如 Node 数据通道测试、Python 数据通道测试 与 Rust 数据通道测试。七、生命周期、背压与双向通信综合引擎与 SDK 源码channel 的生命周期要点懒连接创建 channel 只分配引用并登记到ChannelManagerWebSocket 流在某一侧开始读或写时才真正建立Node SDK 的ensureConnected即此行为。背压由 SDK 流实现承担读端跟不上时暂停底层 WebSocket 接收写入方随之阻塞。正常结束writer 关闭后引擎侧 sender 释放reader 收到流结束Node 端表现为stream.push(null)。异常清理worker 断开连接时其名下 channel 被整体清除超过 5 分钟 TTL 且至少一端从未连接的孤儿 channel 由 60 秒周期的清扫任务回收。双向通信channel 是单向管道需要双向时创建两条 channel每个方向一条。八、延伸阅读通道架构模型寻址、多路复用与拆除Channels 架构说明各语言 channel API 面的完整参考Node SDK、Python SDK、Rust SDK引擎通道管理核心ChannelManager、WebSocket 升级处理SDK 通道实现Node、Rust、Python【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表