news 2026/9/14 20:48:22

使用 iii SDK 构建跨语言 Worker:Node、Python、Rust、Go 统一开发指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
使用 iii SDK 构建跨语言 Worker:Node、Python、Rust、Go 统一开发指南

使用 iii SDK 构建跨语言 Worker:Node、Python、Rust、Go 统一开发指南

【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii

iii 是一个实时编排引擎,围绕 Function、Trigger、Worker 三个核心原语提供跨语言执行、实时能力发现与可观测性。官方 SDK 将引擎能力封装为四套 API 完全对齐的客户端库(Node.js、Python、Rust、Go),让开发者可以用自己熟悉的语言编写 worker,通过 WebSocket 与引擎通信,注册函数与触发器并响应调用。本文以 sdk/README.md 为骨架,结合 sdk/packages 下各语言 SDK 的 README 与源码实现,系统讲解 SDK 的安装、Hello World、统一 API 面、触发器绑定、调用方式、连接重连机制以及开发测试流程。读完本文,你将能在任意一门受支持语言中独立完成"注册函数 -> 绑定触发器 -> 同步/异步调用 -> 接入可观测性"的完整开发闭环。

一、SDK 全景:一套 API 面,四种语言实现

iii SDK 是引擎的官方客户端集合,分布在仓库 sdk/packages 目录下。每个包都以iii-sdk为名发布到各自语言的生态仓库,但共享完全相同的编程模型:

包名语言安装方式详细文档
iii-sdkNode.js / TypeScriptpnpm add iii-sdknpm install iii-sdkREADME
iii-sdkPythonpip install iii-sdkREADME
iii-sdkRust添加依赖到Cargo.tomlREADME
iii-sdkGogo get github.com/iii-hq/iii/sdk/packages/go/iiiREADME

SDK 采用 Apache 2.0 许可(见 sdk/LICENSE)。从 engine/README.md 可以看到引擎的架构定位:引擎提供持久化编排与跨语言执行能力,一个进程通过打开一条 WebSocket 连接到引擎即成为一个 iiiworker,之后注册函数触发器;引擎按需调用这些函数,worker 在同一条 socket 上回复结果(见 Go SDK README 的说明)。

端口速查(来自 engine/README.md):49134 为 WebSocket worker 连接端口,3111 为 HTTP API,3112 为 Stream API,9464 为 Prometheus metrics。

二、环境准备:先启动引擎,再安装 SDK

2.1 启动引擎

所有语言 SDK 都默认连接ws://localhost:49134。在本地运行 SDK 之前,需要先启动 iii 引擎:

# 安装引擎(含全部 CLI 命令) curl -fsSL https://install.iii.dev/iii/main/install.sh | sh # 验证安装 command -v iii && iii --version # 使用内置默认配置启动(自带 in-memory OpenTelemetry 配置) iii --use-default-config

也可以创建config.yaml后以项目配置启动:iii --config /path/to/config.yaml。启动成功后可用iii console打开控制台。引擎默认 WebSocket 地址即为ws://localhost:49134

2.2 安装各语言 SDK

# Node.js / TypeScript pnpm add iii-sdk # 或 npm install iii-sdk # Python pip install iii-sdk # Rust(Cargo.toml) [dependencies] iii-sdk = "0.11" serde_json = "1" tokio = { version = "1", features = ["full"] } # Go(要求 Go 1.24+) go get github.com/iii-hq/iii/sdk/packages/go/iii

三、四语言 Hello World 全解

3.1 Node.js / TypeScript

import { registerWorker } from 'iii-sdk'; const iii = registerWorker('ws://localhost:49134'); iii.registerFunction('hello::greet', async (input) => { return { message: `Hello, ${input.name}!` }; }); iii.registerTrigger({ type: 'http', function_id: 'hello::greet', config: { api_path: '/greet', http_method: 'POST' }, }); const result = await iii.trigger({ function_id: 'hello::greet', payload: { name: 'world' } });

3.2 Python

from iii import register_worker iii = register_worker("ws://localhost:49134") def greet(data): return {"message": f"Hello, {data['name']}!"} iii.register_function({"id": "hello::greet"}, greet) iii.register_trigger({ "type": "http", "function_id": "hello::greet", "config": {"api_path": "/greet", "http_method": "POST"} }) result = iii.trigger({"function_id": "hello::greet", "payload": {"name": "world"}})

注意:Python 的register_function第一个参数是包含"id"键的字典,也可以直接传字符串"hello::greet"(见 Python SDK README 的写法)。

3.3 Rust

use iii_sdk::{register_worker, InitOptions, TriggerRequest, RegisterFunctionMessage, RegisterTriggerInput}; use serde_json::json; #[tokio::main] async fn main() -> Result<(), Box<dyn std::error::Error>> { let iii = register_worker("ws://127.0.0.1:49134", InitOptions::default())?; iii.register_function(RegisterFunctionMessage::with_id("hello::greet".into()), |input| async move { let name = input.get("name").and_then(|v| v.as_str()).unwrap_or("world"); Ok(json!({ "message": format!("Hello, {name}!") })) }); iii.register_trigger(RegisterTriggerInput::new("http", "hello::greet", json!({ "api_path": "/greet", "http_method": "POST" })))?; let result: serde_json::Value = iii .trigger(TriggerRequest::new("hello::greet", json!({ "name": "world" }))) .await?; Ok(()) }

RegisterFunctionMessage::with_id用于按 ID 注册函数,TriggerRequest::new组合了 function_id 与 payload,RegisterTriggerInput::new(type, fn_id, config)绑定触发器。InitOptions::default()负责默认的连接初始化参数。

3.4 Go

package main import ( "context" "encoding/json" "log" iii "github.com/iii-hq/iii/sdk/packages/go/iii" ) func main() { client := iii.RegisterWorker("ws://127.0.0.1:49134") client.RegisterFunction("hello::greet", func(ctx context.Context, data json.RawMessage) (any, error) { var req struct { Body struct { Name string `json:"name"` } `json:"body"` } _ = json.Unmarshal(data, &req) return map[string]any{ "status_code": 200, "body": map[string]string{"message": "Hello, " + req.Body.Name + "!"}, }, nil }) client.RegisterTrigger("hello-http", "http", "hello::greet", json.RawMessage(`{"api_path":"/greet","http_method":"POST"}`), nil) if err := client.Connect(context.Background()); err != nil { log.Fatal(err) } defer client.Close() }

3.5 理解 HTTP 触发器信封(Envelope)

观察以上代码可以发现一个细节:Go 的 handler 返回值带有status_codebody两个字段,而 Node/Python/Rust 示例返回的是纯数据对象。原因在于 Go SDK README 明确说明的HTTP trigger envelope约定:

当一个函数经由http触发器被调用时,引擎会把请求包装为{ path, method, body, headers, … },并期望函数返回{ status_code, body }。仅通过 socket 调用的函数可以使用任意 payload 形状。

也就是说,凡是绑定 HTTP 触发器的函数,入参和返回值都必须遵守 HTTP 信封格式;而通过 SDK 内部trigger()直接调用的函数则不受此约束。这是编写跨语言 worker 时最容易踩坑、也最需要统一的契约点。

四、统一 API 面:四语言操作对照

sdk/README.md 提供了一张核心操作对照表,四个 SDK 暴露的 API 面完全一致:注册函数、注册触发器、发起调用。下表完整继承原文档并补充各语言的实际签名:

操作Node.jsPythonRustGo说明
初始化registerWorker(url)register_worker(url, options?)register_worker(url, options)iii.RegisterWorker(url)创建 SDK 实例并自动连接
注册函数iii.registerFunction(id, handler, options?)iii.register_function(id, handler)iii.register_function(id, \|input\| ...)client.RegisterFunction(id, handler)注册可按名字调用的函数
注册触发器iii.registerTrigger({ type, function_id, config })iii.register_trigger({"type": ..., "function_id": ..., "config": ...})iii.register_trigger(type, fn_id, config)?client.RegisterTrigger(id, type, fn, cfg, meta)将触发器(HTTP、cron、queue 等)绑定到函数
同步调用await iii.trigger({ function_id, payload })await iii.trigger({"function_id": id, "payload": data})iii.trigger(TriggerRequest::new(id, data)).await?client.Trigger(ctx, iii.TriggerRequest{...})调用函数并等待结果
异步调用(fire-and-forget)iii.trigger({ function_id, payload, action: TriggerAction.Void() })iii.trigger({"function_id": id, ..., "action": TriggerAction.Void()})iii.trigger(TriggerRequest { action: Some(TriggerAction::Void), ... })client.Trigger(ctx, iii.TriggerRequest{Action: iii.VoidAction()})不等待结果的调用
队列调用(enqueue)iii.trigger({ function_id, payload, action: TriggerAction.Enqueue({ queue }) })iii.trigger({"function_id": id, ..., "action": TriggerAction.Enqueue(queue="name")})iii.trigger(TriggerRequest { action: Some(TriggerAction::Enqueue { queue }), ... })client.Trigger(ctx, iii.TriggerRequest{Action: iii.EnqueueAction("queue")})通过命名队列路由调用

原文档中有一个重要的兼容性说明,写作本文时仍须强调:

callcallVoidtriggerVoid(以及 Python/Rust 的对应变体)已被移除,所有调用统一使用trigger()。需要 fire-and-forget 时,使用trigger({ function_id, payload, action: TriggerAction.Void() })

registerWorker()/register_worker()创建 SDK 实例并自动连接引擎,内部处理 WebSocket 通信、自动重连与 OpenTelemetry 埋点。Node SDK 的导出集中在 sdk/packages/node/iii/src/index.ts:对外暴露registerWorkerInitOptionsTriggerActionInvocationError等类型与函数。

4.1 注册函数:从简单到类型安全

Node 侧支持可选的 options 参数,Python/Rust 直接传入 id 与 handler。以"创建订单"为例:

// Node.js iii.registerFunction('orders::create', async (input) => { return { status_code: 201, body: { id: '123', item: input.body.item } } })
# Python def create_order(data): return {"status_code": 201, "body": {"id": "123", "item": data["body"]["item"]}} iii.register_function("orders::create", create_order)
// Rust iii.register_function("orders::create", |input: Value| async move { let item = input["body"]["item"].as_str().unwrap_or(""); Ok(json!({ "status_code": 201, "body": { "id": "123", "item": item } })) });

Go SDK 更进一步,提供RegisterFunctionTyped[Req, Resp]:通过类型参数推断请求/响应 JSON Schema,并作为request_format/response_format广播给引擎与仪表盘(对应 Rust SDK 的#[derive(JsonSchema)])。Go 没有编译期 derive,因此基于反射推断(依赖invopop/jsonschema),可用jsonjsonschema结构体标签定制 schema:

type CreateOrderRequest struct { Item string `json:"item" jsonschema:"required"` Quantity int `json:"quantity" jsonschema:"minimum=1"` } type OrderResult struct { ID string `json:"id"` } iii.RegisterFunctionTypedCreateOrderRequest, OrderResult (OrderResult, error) { return OrderResult{ID: "ord_123"}, nil })

4.2 注册触发器:把函数暴露给外部事件

触发器把函数与外部事件源(HTTP、cron、queue 等)绑定起来。四语言用法一致:

// Node.js iii.registerTrigger({ type: 'http', function_id: 'orders::create', config: { api_path: '/orders', http_method: 'POST' }, })
# Python iii.register_trigger({ "type": "http", "function_id": "orders::create", "config": {"api_path": "/orders", "http_method": "POST"}, })
// Rust iii.register_trigger("http", "orders::create", json!({ "api_path": "/orders", "http_method": "POST" }))?;
// Go(注意 Go 的签名多一个触发器 ID 参数,并可选携带 metadata) client.RegisterTrigger("orders-http", "http", "orders::create", json.RawMessage(`{"api_path":"/orders","http_method":"POST"}`), nil)

Go SDK 还支持RegisterTriggerType(id, description, handler)实现自定义触发器类型,这与 Node SDK 中registerTriggerType的能力对齐——四个 SDK 都允许把引擎的触发器能力延伸到自定义事件源。

4.3 调用函数:三种调用语义

统一调用入口trigger()支持三种语义:

// Node.js import { registerWorker, TriggerAction } from 'iii-sdk' const iii = registerWorker('ws://localhost:49134') // 1. 同步等待结果 const result = await iii.trigger({ function_id: 'orders::create', payload: { item: 'widget' } }) // 2. fire-and-forget(如埋点、日志) iii.trigger({ function_id: 'analytics::track', payload: { event: 'page_view' }, action: TriggerAction.Void() }) // 3. 通过命名队列异步路由 iii.trigger({ function_id: 'orders::process', payload: { order_id: '456' }, action: TriggerAction.Enqueue({ queue: 'payments' }) })

Python 对应写法为iii.trigger({"function_id": ..., "action": TriggerAction.Enqueue(queue="name")});Rust 通过TriggerRequest { action: Some(TriggerAction::Enqueue { queue: "payments".to_string() }), .. }表达;Go 使用iii.EnqueueAction("jobs")。队列语义由引擎侧 queue 模块落地,SDK 只负责把调用请求路由进命名队列。

五、连接、重连与生命周期(源码视角)

以 Node SDK 为例,sdk/packages/node/iii/src/iii-constants.ts 定义了完整连接参数,这些默认值体现了 SDK 的健壮性设计:

常量默认值说明
DEFAULT_BRIDGE_RECONNECTION_CONFIG初始 1000ms、上限 30000ms、倍率 2、抖动 0.3、无限重试指数退避 + 随机抖动
DEFAULT_INVOCATION_TIMEOUT_MS30000默认调用超时
WS_HANDSHAKE_TIMEOUT_MS10000WebSocket 握手超时(与 Rust SDK 的 connect_timeout 对齐)
WS_PING_INTERVAL_MS20000保活 ping 间隔(与 Rust SDK 的 ping_interval 对齐)
WS_IDLE_TIMEOUT_MS60000长时间无入站帧强制重连(与 Rust SDK 的 idle_timeout 对齐)

从 sdk/packages/node/iii/src/iii.ts 源码还可以看到两个关键设计:

  1. 地址解析优先级:显式传入的地址 > 环境变量III_URL(由iii compose、容器运行时、systemd 等 supervisor 注入)> 默认值ws://127.0.0.1:49134。源码特意使用 IPv4 回环地址而非localhost,避免主机只监听 IPv4 时localhost解析到::1导致连不上。
  2. worker 身份与环境注入III_WORKER_NAME携带编排器分配的 worker 名(Compose 托管 worker 由 supervisor 注入),其优先级高于代码内嵌的显式名称,因为引擎按名字匹配在线注册;默认名形如${os.hostname()}:${process.pid}

Go SDK 的连接行为与 Node 完全对齐(Go SDK README):指数退避重连(起始 1s、×2、封顶 30s、±30% 抖动、无限重试,可用iii.WithReconnectConfig覆盖);离线缓冲——断线期间发出的调用被缓冲并在重连后冲刷,注册信息则从内存注册表重放;每次连接最后注册 worker metadata(标记runtime: "go")。

Register*类调用可以在Connect之前或之后执行:注册保存在内存中,并在每次(重)连接时(重新)发送给引擎。这一点在 Node/Python/Rust/Go 四个 SDK 中行为一致,是断线重连后 worker 功能不丢失的基石。

六、进阶能力:数据通道、流操作与可观测性

sdk/README.md 明确指出:语言相关的进阶细节(modules、streams、OpenTelemetry)以各语言 SDK 的 README 为准。这里汇总仓库内可验证的四个进阶能力。

6.1 流式数据通道(Channels)

Go SDK 提供CreateChannel(ctx, bufferSize)打开双向数据通道(writer + reader 两端,各自独立 WebSocket)。通道上的约定是:文本帧即消息SendMessage/OnMessage),二进制帧即流数据Write/ReadAll),按 WebSocket opcode 区分,无需额外信封;writer 以Close()结束流,reader 将其视为 EOF。通道引用(WriterRef/ReaderRef)可以放进 trigger payload 传给另一个 worker,由对方打开对端——适合大体积或持续流式数据的传输,而非把所有数据塞进单次调用结果。Node/Python/Rust SDK 同样内置 channel 与 stream 模块(见 node/iii/src/channels.ts、rust/iii/src/channels.rs 等)。

6.2 流(Streams)与原子更新

Rust SDK README 展示了通过内置函数stream::set/stream::update操作流:以stream_name+group_id+item_id定位数据项,配合UpdateOp::increment("total", 100)UpdateOp::set("status", json!("processing"))实现原子更新操作。这套能力在四个 SDK 中均有等价实现,用于跨 worker 共享实时状态。

6.3 可观测性:OpenTelemetry 一体化

SDK 内置 OpenTelemetry 埋点:registerWorker()自动完成 OTel 初始化,invocation 以 span 形式记录(Node 源码中通过withSpanrecordSpanEventinjectTraceparentinjectBaggage等工具实现 W3C trace context 的注入与透传)。Rust 侧Logger结构体(iii_helpers::observability::Logger)发射 OTelLogRecord,未初始化 OTel 时回退到tracingcrate。引擎以--use-default-config启动时自带 in-memory OTel 配置,因此本地开发即可直接观察 traces、metrics、logs。

6.4 调用元数据与错误语义

  • 元数据侧车:Go handler 通过iii.MetadataFromContext(ctx)读取每次调用可选的元数据(如{"tenant":"acme"}),注册函数时也可用RegisterFunctionOptions{Metadata: ...}附加静态元数据。元数据随调用上下文传递,不改变 handler 签名。
  • 错误类型:Go 中errors.Is(err, iii.ErrTimeout)判断超时、errors.Is(err, iii.ErrNotConnected)判断关闭引起的取消、errors.As(err, &ie)拿到携带远端Code/Message/Stacktrace*iii.InvocationError;handler 返回任意非InvocationError错误会被统一报告为invocation_failed,panic 也会被 recover 并同样上报,保证调用方得到错误而不是超时。Node SDK 对应导出InvocationErrorRegistrationRejectedError(index.ts)。

七、开发与测试:从源码构建各语言 SDK

仓库内每个 SDK 包都自带构建与测试管线(sdk/packages/node/iii/package.json、sdk/packages/python/iii/pyproject.toml、sdk/packages/rust/iii/Cargo.toml)。

7.1 前提条件

  • Node.js 20+ 与 pnpm(Node SDK)
  • Python 3.10+ 与 uv(Python SDK)
  • Rust 1.85+ 与 Cargo(Rust SDK)
  • 运行在ws://localhost:49134的 iii 引擎

7.2 构建

cd packages/node && pnpm install && pnpm build cd packages/python/iii && python -m build cd packages/rust/iii && cargo build --release

Python SDK 的开发模式安装、类型检查与 lint:

pip install -e . mypy src ruff check src

7.3 测试

cd packages/node && pnpm test cd packages/python/iii && pytest cd packages/rust/iii && cargo test

各 SDK 的测试目录覆盖了丰富的场景,是理解 API 语义的最佳教材:Node 侧 sdk/packages/node/iii/tests 包含连接握手超时、心跳、重连(connection-reattach.test.ts)、触发器动作(trigger-action.test.ts)、触发器注册错误(trigger-registration-error.test.ts)、流(stream.test.ts)、RBAC(rbac-workers.test.ts)等用例;Python 侧 sdk/packages/python/iii/tests 覆盖同步/异步 API、重连发送 reattach、环境契约、引擎常量等;Rust 侧 sdk/packages/rust/iii/tests 提供 mock engine 与集成测试;Go 侧 sdk/packages/go/iii/tests 覆盖触发器注册、数据通道、worker metadata、注册去重等。阅读这些测试可以帮助你理解 SDK 在异常与边界情况下的真实行为。

八、示例与下一步

仓库为每种语言都提供了可运行的示例工程:

  • Node:iii-example(含 HTTP、队列、DLQ、流、状态、中间件示例)——sdk/packages/node/iii-example
  • Python:iii-example(函数、状态、流、触发器类型示例)——sdk/packages/python/iii-example
  • Rust:iii-example(HTTP、cron、自定义触发器、Logger 示例)——sdk/packages/rust/iii-example
  • Go:iii-example——sdk/packages/go/iii-example

仓库中还提供了 helpers 包(sdk/packages/node/helperssdk/packages/python/helperssdk/packages/rust/helpers),封装 HTTP、队列、流、worker 连接管理与可观测性工具,供 SDK 内部复用并在示例中展示。引擎架构细节可参阅 engine/README.md,SDK 与引擎的 WebSocket 协议交互可从 engine/src/protocol.rs 与各 SDK 的protocol.*源文件对照理解。

总结

iii SDK 的价值在于用一条 WebSocket 把任意语言进程接入统一编排引擎:四个语言包共享相同的 API 面(registerWorker->registerFunction->registerTrigger->trigger),统一处理连接、重连、离线缓冲、OpenTelemetry 埋点与错误语义;在此之上,channels 提供流式数据传输、streams 提供共享实时状态、队列触发器提供异步路由。无论团队技术栈是 TypeScript、Python、Rust 还是 Go,都可以用同样的心智模型开发 worker,并让不同语言的服务在同一引擎内相互调用、统一观测。

【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/14 20:46:54

Twitter自动化运营全攻略:从手动发帖到系统化引流涨粉

做Twitter运营的朋友&#xff0c;应该都体会过那种“一个人活成一支队伍”的疲惫感。内容要写、帖子要发、留言要回、竞品要盯、数据要记&#xff0c;一天下来真正花在“思考策略”上的时间反而不多。我也试过靠闹钟提醒自己凌晨爬起来发帖&#xff0c;结果人是起来了&#xff…

作者头像 李华
网站建设 2026/9/14 20:46:45

计算机二级WPS考试核心考点与备考策略

1. 计算机二级WPS考试概述作为国内办公软件应用能力的重要认证&#xff0c;计算机二级WPS考试近年来报考人数持续攀升。根据官方数据统计&#xff0c;2023年全国报考WPS科目的人数较2022年增长了47%&#xff0c;这主要得益于国产办公软件的快速发展和企事业单位对WPS技能要求的…

作者头像 李华
网站建设 2026/9/14 20:46:18

Linux驱动自动加载全解析:从udev到modprobe的匹配链路

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/14 20:45:39

在线绘图工具实测:替代Visio的10款流程图/架构图协作方案

先说个背景。我自己用了十年的Visio&#xff0c;从2007一路用到2019。以前画网络拓扑、泳道流程图、机房机柜图&#xff0c;基本都靠它。但这两年我越来越不想打开Visio了——倒不是画图水平退步&#xff0c;而是“用Visio”这个动作本身就变得很烦&#xff1a;公司电脑要申请授…

作者头像 李华