ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

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

使用 iii SDK 构建跨语言 Worker:Node、Python、Rust、Go 统一开发指南 使用 iii SDK 构建跨语言 WorkerNode、Python、Rust、Go 统一开发指南【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iiiiii 是一个实时编排引擎围绕 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-sdk或npm install iii-sdkREADMEiii-sdkPythonpip install iii-sdkREADMEiii-sdkRust添加依赖到Cargo.tomlREADMEiii-sdkGogo get github.com/iii-hq/iii/sdk/packages/go/iiiREADMESDK 采用 Apache 2.0 许可见 sdk/LICENSE。从 engine/README.md 可以看到引擎的架构定位引擎提供持久化编排与跨语言执行能力一个进程通过打开一条 WebSocket 连接到引擎即成为一个 iiiworker之后注册函数和触发器引擎按需调用这些函数worker 在同一条 socket 上回复结果见 Go SDK README 的说明。端口速查来自 engine/README.md49134 为 WebSocket worker 连接端口3111 为 HTTP API3112 为 Stream API9464 为 Prometheus metrics。二、环境准备先启动引擎再安装 SDK2.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 # RustCargo.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 / TypeScriptimport { 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 Pythonfrom iii import register_worker iii register_worker(ws://localhost:49134) def greet(data): return {message: fHello, {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 Rustuse iii_sdk::{register_worker, InitOptions, TriggerRequest, RegisterFunctionMessage, RegisterTriggerInput}; use serde_json::json; #[tokio::main] async fn main() - Result(), Boxdyn 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 与 payloadRegisterTriggerInput::new(type, fn_id, config)绑定触发器。InitOptions::default()负责默认的连接初始化参数。3.4 Gopackage 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_code与body两个字段而 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-forgetiii.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()})不等待结果的调用队列调用enqueueiii.trigger({ function_id, payload, action: TriggerAction.Enqueue({ queue }) })iii.trigger({function_id: id, ..., action: TriggerAction.Enqueue(queuename)})iii.trigger(TriggerRequest { action: Some(TriggerAction::Enqueue { queue }), ... })client.Trigger(ctx, iii.TriggerRequest{Action: iii.EnqueueAction(queue)})通过命名队列路由调用原文档中有一个重要的兼容性说明写作本文时仍须强调call、callVoid、triggerVoid以及 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对外暴露registerWorker、InitOptions、TriggerAction、InvocationError等类型与函数。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可用json与jsonschema结构体标签定制 schematype CreateOrderRequest struct { Item string json:item jsonschema:required Quantity int json:quantity jsonschema:minimum1 } 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(queuename)})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 源码还可以看到两个关键设计地址解析优先级显式传入的地址 环境变量III_URL由iii compose、容器运行时、systemd 等 supervisor 注入 默认值ws://127.0.0.1:49134。源码特意使用 IPv4 回环地址而非localhost避免主机只监听 IPv4 时localhost解析到::1导致连不上。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 流式数据通道ChannelsGo 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_namegroup_iditem_id定位数据项配合UpdateOp::increment(total, 100)、UpdateOp::set(status, json!(processing))实现原子更新操作。这套能力在四个 SDK 中均有等价实现用于跨 worker 共享实时状态。6.3 可观测性OpenTelemetry 一体化SDK 内置 OpenTelemetry 埋点registerWorker()自动完成 OTel 初始化invocation 以 span 形式记录Node 源码中通过withSpan、recordSpanEvent、injectTraceparent、injectBaggage等工具实现 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.InvocationErrorhandler 返回任意非InvocationError错误会被统一报告为invocation_failedpanic 也会被 recover 并同样上报保证调用方得到错误而不是超时。Node SDK 对应导出InvocationError、RegistrationRejectedErrorindex.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 与 pnpmNode SDKPython 3.10 与 uvPython SDKRust 1.85 与 CargoRust 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 --releasePython SDK 的开发模式安装、类型检查与 lintpip install -e . mypy src ruff check src7.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、RBACrbac-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 在异常与边界情况下的真实行为。八、示例与下一步仓库为每种语言都提供了可运行的示例工程Nodeiii-example含 HTTP、队列、DLQ、流、状态、中间件示例——sdk/packages/node/iii-examplePythoniii-example函数、状态、流、触发器类型示例——sdk/packages/python/iii-exampleRustiii-exampleHTTP、cron、自定义触发器、Logger 示例——sdk/packages/rust/iii-exampleGoiii-example——sdk/packages/go/iii-example仓库中还提供了 helpers 包sdk/packages/node/helpers、sdk/packages/python/helpers、sdk/packages/rust/helpers封装 HTTP、队列、流、worker 连接管理与可观测性工具供 SDK 内部复用并在示例中展示。引擎架构细节可参阅 engine/README.mdSDK 与引擎的 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),仅供参考
返回列表