资讯动态

iii 框架 Channels 数据通道实战:CSV 批量导入与 worker 间流式传输

发布时间:2026/9/14 13:52:10 来源:尧图企业网站定制
iii 框架 Channels 数据通道实战CSV 批量导入与 worker 间流式传输【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii本篇技术指南围绕 iii 开源仓库中 Linkly 系列教程的第 6 章展开在流stream适合承载实时涓流式事件的前提下通道channel用于一次性移动大批量数据——它是一条在两个端点之间直连的流式管道而不是一次请求一次响应。你将亲手创建一个独立的bulk-importerworker通过通道把一整份 CSV 链接表流式上传到引擎并逐行触发linkworker 的link::create函数完成批量入库最后验证新链接立即可以被解析。读完本文你将掌握 iii SDK 中createChannel的完整用法、readerRef/writerRef的可序列化句柄机制以及通道在引擎侧的真实落地方式。通道与流的定位差异在进入代码之前先明确 iii 中两类数据传输原语的分工流stream面向实时事件的涓流式推送例如 Linkly 第 5 章中click-streamerworker 把每一次点击实时推给订阅者通道channel面向大批量数据的一次性搬运是一条两个端点之间的直连流式管道而不是一次请求 一次响应的普通函数调用。通道天然是双向的——每一端都同时拥有一个 reader 和一个 writer参见 Node SDK 的 Channel 类型定义。不过本章只使用单方向客户端脚本把 CSV 写入 writer 端bulk-importerworker 从 reader 端把整份数据读出来。为了不干扰已有业务教程把导入能力放进一个独立的bulk-importerworker让第 1 章创建的linkworker 始终聚焦于单条链接的创建与解析。第一步脚手架导入 worker与第 1 章初始化linkworker 的方式完全一致在项目目录下执行iii worker init bulk-importer --language typescript这会生成一个标准 TypeScript worker 骨架包含src/index.ts入口与iii.yaml或对应的 worker 清单配置。随后你需要把新 worker 注册进项目iii worker add ./bulk-importer注册完成后引擎便能发现bulk-importerworker 暴露的函数并允许其他 worker 或客户端通过iii trigger/worker.trigger调用它。通过通道导入 CSVworker 端实现bulk-importerworker 只暴露一个函数bulk-importer::import_csv。它的职责是接收通道的 reader 端引用把通道里流过来的数据完整读出并拼装为 UTF-8 文本跳过表头行后逐行触发link::create。用下面的实现替换生成的bulk-importer/src/index.tsimport { registerWorker, Logger } from iii-sdk; const worker registerWorker(process.env.III_URL ?? ws://localhost:49134, { workerName: bulk-importer, }); const logger new Logger(); worker.registerFunction(bulk-importer::import_csv, async (input) { const chunks: Buffer[] []; for await (const chunk of input.reader.stream) { chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)); } const csv Buffer.concat(chunks).toString(utf-8); const rows csv.trim().split(\n).slice(1); // skip the header row let imported 0; for (const row of rows) { const [code, url] row.split(,); if (!url) continue; await worker.trigger({ function_id: link::create, payload: { code: code.trim(), url: url.trim() }, }); imported 1; } logger.info(bulk import complete, { imported }); return { imported }; }); logger.info(bulk-importer ready);几个值得注意的实现细节input.reader.stream是标准的 Node.jsReadable因此可以直接用for await...of逐块消费。每块chunk可能是Buffer也可能不是所以代码统一做了Buffer.from(chunk)兜底SDK 侧的行为可参见 ChannelReader 实现。Buffer.concat(chunks)先把所有分块拼接成完整字节流再统一toString(utf-8)这是处理通道数据按帧分片到达的正确姿势——通道并不保证一次交付整份数据。逐行解析时跳过表头slice(1)并按行内逗号拆出code与url空url的行被continue跳过避免脏数据触发无效创建。每一行通过worker.trigger调用link::create这与客户端直接触发函数走的是同一套调用协议只是调用方变成了另一个 worker。SDK 底层通道是如何创建与寻址的客户端与 worker 使用的createChannel在 SDK 中返回一个包含四件套的Channel对象见 types.tswriter本端的ChannelWriter写数据reader本端的ChannelReader读数据writerRef可序列化的 writer 端句柄可放进普通 trigger payload 发给别的 workerreaderRef可序列化的 reader 端句柄同样可跨进程传递。这正是本教程把 reader 端塞进 payload发给bulk-importer::import_csv的底层依据——readerRef/writerRef不是内存对象而是可跨 WebSocket 连接传输的引用。从实现上看通道数据并不走引擎主端口而是通过 worker 自身所连接的 listener 上的 WebSocket 端点传输。SDK 的buildChannelUrl会拼出如下地址见 channels.ts{engineWsBase}/ws/channels/{channel_id}?key{access_key}dirread|write其中key是随每个StreamChannelRef一起下发的access_key能力令牌dir区分读写方向。引擎侧iii-worker-manager的每个 listener 都会在同一端口挂载该 WebSocket 端点路由注册见 engine/src/workers/worker/mod.rs因此 SDK worker 无需任何额外配置即可直接使用createChannel()通道数据始终流经 worker 所连接的那个 listener而非引擎主端口。同时engine::channels::create属于始终放行的基础设施内置能力即使 RBAC listener 的expose_functions为空也能创建通道并下发通道引用详见 worker/README.md。写入端还有一处容易被忽略的工程细节ChannelWriter内部以64KB 帧大小FRAME_SIZE 64 * 1024对数据做分片发送见 channels.ts并在final阶段延迟约 10ms 发送关闭帧以确保所有数据帧先于关闭帧到达引擎、避免尾部数据截断——这正是大数据量一次性搬运场景下 SDK 为保证完整性所做的处理。验证客户端脚本上传 CSV本章及第 7 章与前几章不同需要在本地准备node 与 npm环境——因为我们开始编写运行在 worker 之外的客户端代码了。准备独立脚本目录上传脚本是一个独立的短脚本它创建通道、把 CSV 写入 writer 端、再把 reader 端句柄交给bulk-importer::import_csv。由于它不是 worker教程建议在项目目录之外给它一个独立的临时目录mkdir test-channels cd test-channels npm init -y npm pkg set typemodule npm install iii-sdknpm pkg set typemodule让脚本以 ESM 方式运行从而可以使用import语法npm install iii-sdk引入 Node SDK浏览器端使用iii-browser这里不需要。上传脚本完整实现把下面的内容保存为test-channels/import-links.jsimport { registerWorker } from iii-sdk; const worker registerWorker(process.env.III_URL ?? ws://localhost:49134, { workerName: uploader, }); const csv [code,url, mylink,https://iii.dev, mydocslink,https://iii.dev/docs].join(\n); const channel await worker.createChannel(); channel.writer.stream.write(Buffer.from(csv)); channel.writer.stream.end(); const result await worker.trigger({ function_id: bulk-importer::import_csv, payload: { reader: channel.readerRef }, }); console.log(result); await worker.shutdown();脚本要点它以registerWorker连接引擎并自报workerName: uploader但它没有注册任何函数只是一个临时客户端worker.createChannel()在引擎侧创建一条通道并返回本地writer/reader与可序列化的readerRefchannel.writer.stream.write(...)写入 CSV 字节channel.writer.stream.end()结束写入端——注意end()内部会负责关闭 WebSocket 帧的时序见上文 SDK 实现触发bulk-importer::import_csvpayload 里直接携带reader: channel.readerRef脚本完成后调用await worker.shutdown()干净地断开连接。运行与预期输出在引擎运行中的前提下执行node import-links.js预期输出{ imported: 2 }两条新链接会立即生效——不需要任何刷新或重建因为linkworker 的存储对link::create的每次调用是即时可见的。立即验证解析iii trigger link::resolve codemydocslink{ url: https://iii.dev/docs }通道能力的更多佐证SDK 测试用例仓库自带的测试可以进一步印证上述用法并展示通道不止支持单向二进制流。data-channels.test.ts 中覆盖了两个典型场景sender → processor 单方向批量传输sender 创建通道把整份 JSON 数据writer.stream.end(payload, cb)写入并等待完成同时把readerRef放进 trigger payloadprocessor 侧for await (const chunk of input.reader.stream)读完全部数据后做统计并返回。这与本教程的 CSV 导入是同一模式见该文件第 7-85 行的第一个用例。worker ↔ coordinator 双向流coordinator 同时创建输入、输出两条通道把readerRef与writerRef一起传给 workerworker 边读边通过writer.sendMessage回传结构化文本消息progress/complete最后再通过writer.stream.end(...)回传二进制结果。这展示了ChannelWriter的sendMessage与ChannelReader的onMessage组合使用方式见该文件第 87-217 行的第二个用例。这些测试从侧面说明通道同时支持二进制流stream.write/for await与结构化文本消息sendMessage/onMessage本教程使用的是前者。单元测试同时验证了分块写入chunkSize: 10与drain背压处理说明 SDK 在写大数据量时按流式语义工作不会一次性把整份数据塞进内存发送缓冲。小结与下一步至此Linkly 已经可以借助一个专职的bulk-importerworker通过一条通道完成一份文件级链接表的一次性流式上传客户端创建通道并写入 CSVworker 从 reader 端读出、逐行触发link::create全程无需把 CSV 塞进单个 trigger payload也无需linkworker 分心处理批量逻辑。下一章 Ch. 7Bring in the browser 会把一个浏览器标签页变成 worker直接在浏览器里调用link::create、订阅实时点击流并注册一个可供服务端回调的浏览器端函数。【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价