Hatchet Ruby 示例仓库完全指南:从 Hello World 到生产级工作流实战
【免费下载链接】hatchet🪓 An orchestration engine for background tasks, AI agents, and durable workflows项目地址: https://gitcode.com/GitHub_Trending/ha/hatchet
本指南围绕 Hatchet 开源编排引擎(Go 编写)的 Ruby SDK 示例仓库展开,系统讲解sdks/ruby/examples目录的目录结构、环境搭建与运行方式,并结合仓库内真实源码,覆盖任务定义、DAG 依赖编排、事件触发、持久化(durable)工作流、并发控制、重试、定时调度、Webhook 与单元测试等完整能力。读完本文,你将掌握使用 Hatchet Ruby SDK 从零搭建 Worker、编写工作流并在本地验证其行为的整套实战方法。
1. 示例仓库定位与整体结构
sdks/ruby/examples/README.md 是 Hatchet Ruby SDK 官方示例的入口文档。它声明了两个核心事实:
- 该目录演示了 Hatchet Ruby SDK 的典型用法;
- 环境准备只需要一条命令
bundle install。
虽然 README 本身篇幅精简,但它所指向的示例目录实际承载了 50 余个功能主题、上百个 Ruby 源文件,覆盖了 SDK 的绝大部分能力面。从目录树可以清晰看到示例的横向组织方式——每个主题一个子目录,绝大多数目录下同时包含worker.rb(定义并注册工作流的 Worker 实现)与test_*_spec.rb(RSpec 端到端验证):
- 基础篇:
simple/、quickstart/、dag/、events/ - 生命周期与容错:
retries/、timeout/、on_failure/、on_success/、non_retryable/、cancellation/ - 并发与限流:
concurrency_limit/、concurrency_limit_rr/、concurrency_cancel_*、concurrency_shared/、rate_limit/、priority/ - 编排模式:
child/、fanout/、bulk_fanout/、durable/、durable_event/、durable_sleep/、conditions/ - 触发与调度:
cron/、scheduled/、trigger_methods/、idempotency/ - 扩展能力:
streaming/、webhooks/、sticky_workers/、affinity_workers/、dependency_injection/、logger/、serde/ - 工程实践:
unit_testing/、bulk_operations/、migration_guides/
本文将以 README 的骨架(Setup → Examples → Hello World)为主线,把上述目录中的真实源码作为深化素材,逐层展开。
2. 环境搭建:两条命令跑通
2.1 安装依赖
进入示例目录后执行:
bundle install依赖声明位于 Gemfile,其中最关键的一行是:
gem "hatchet-sdk", path: "../src"这意味着示例仓库直接以本地源码路径引用 SDK(而非发布到 RubyGems 的版本),因此你可以在修改 SDK 源码后立即验证效果,适合 SDK 贡献者与深度调试场景。Gemfile 同时还引入了三个辅助依赖:
base64:编解码基础库;rspec ~> 3.0:运行示例自带的端到端测试;net-http:在测试 fixture 中发起 Worker 健康检查 HTTP 请求。
2.2 运行前的前置条件
示例运行依赖一个可用的 Hatchet 后端。Hatchet 支持多种部署形态(Docker Compose、hatchet-lite 等),仓库根目录提供了一键编排配置,例如 docker-compose.yml 与 docker-compose.infra.yml。SDK 默认通过环境变量(如HATCHET_CLIENT_TOKEN、HATCHET_CLIENT_SERVER_URL)连接后端,具体配置项可在 pkg/config/client 的加载逻辑中查看。本地开发时也可参考 hack/dev/start-api.sh 与 hack/dev/start-engine.sh 拉起 API 与 Engine 进程。
3. Hello World:README 指明的入门入口
README 中给出的唯一可运行示例是:
bundle exec ruby hatchet_client.rb这里的hatchet_client.rb是 sdks/ruby/examples/hatchet_client.rb,它演示了 SDK 客户端三个最基础的 API:
require 'hatchet-sdk' # 初始化客户端(全局单例,避免重复连接) HATCHET = Hatchet::Client.new() unless defined?(HATCHET) # 1) 创建事件 result = HATCHET.events.create( key: "test-event", data: { message: "test" } ) puts "Event created: #{result.inspect}" # 2) 直接触发一个名为 simple 的工作流运行 run = HATCHET.runs.create( name: "simple", input: { Message: "test workflow run" }, ) puts "TriggeredRun ID: #{run.metadata.id}" # 3) 轮询该运行的最终状态 result = HATCHET.runs.poll(run.metadata.id) puts "Run status: #{result.status}"这段代码本身就是一份完整的“三连”教学:发事件 → 触运行 → 查状态。其中runs.create的name参数对应工作流注册名,input为 JSON 输入;runs.poll返回带metadata.id与status的运行句柄。
3.1 更轻量的 Quickstart 入口
与hatchet_client.rb互补的是 quickstart/ 目录,它提供了一个“先定义任务、再同步调用”的最小闭环:
quickstart/workflows/first_task.rb:
require "hatchet-sdk" HATCHET = Hatchet::Client.new unless defined?(HATCHET) FIRST_TASK = HATCHET.task(name: "first-task") do |input, ctx| puts "first-task called" { "transformed_message" => input["message"].downcase } endquickstart/run.rb:
require_relative "workflows/first_task" result = FIRST_TASK.run({ "message" => "Hello World!" }) puts "Finished running task: #{result['transformed_message']}"HATCHET.task(name: ...) do |input, ctx| ... end是 Ruby SDK 定义任务的核心 DSL:块接收input(哈希输入)与ctx(上下文),返回值会序列化为任务的输出。FIRST_TASK.run(input)是同步运行入口,返回任务输出哈希;而worker.rb中还会看到run_no_wait(异步触发并返回引用)等变体。
4. Worker 体系:注册与启动所有示例
4.1 单 Worker 聚合注册
examples/worker.rb 是所有示例的“总装车间”。它的结构非常值得学习:
- 用
require_relative加载每个主题目录的 worker 文件(如simple/worker、dag/worker、durable/worker),这些文件在被加载时即完成HATCHET.task(...)/HATCHET.workflow(...)的定义; - 将定义好的常量(
SIMPLE、DAG_WORKFLOW、DURABLE_WORKFLOW……)汇总进ALL_WORKFLOWS数组,并按“Tier 1 基础 → Tier 2 并发 → Tier 3 编排 → Tier 4-5 高级”分层注释; - 最后创建并启动 Worker:
HATCHET = Hatchet::Client.new(debug: true) unless defined?(HATCHET) worker = HATCHET.worker("all-examples-worker", slots: 40, workflows: ALL_WORKFLOWS) worker.start这里slots: 40表示该 Worker 进程同时提供 40 个执行槽位(slot),是 Hatchet 并发资源模型的核心参数之一(槽位语义可参见 pkg/repository/slot_types.go 与 pkg/worker 目录的实现);debug: true打开 SDK 调试日志,方便观察调度与执行链路。
4.2 单主题独立 Worker
每个主题目录的worker.rb末尾通常都有:
def main worker = HATCHET.worker("test-worker", workflows: [SIMPLE, SIMPLE_DURABLE]) worker.start end main if __FILE__ == $PROGRAM_NAMEif __FILE__ == $PROGRAM_NAME保证了该文件被 require 时不启动 Worker、直接运行时才启动——这正是 4.1 节聚合加载的前提。因此你可以单独运行任一示例:
bundle exec ruby simple/worker.rb4.3 测试基建:Worker Fixture
examples/worker_fixture.rb 为端到端测试提供进程级基础设施:
HatchetWorkerFixture.with_worker(command, healthcheck_port:)以子进程方式启动 Worker,并通过HATCHET_CLIENT_WORKER_HEALTHCHECK_ENABLED=true与HATCHET_CLIENT_WORKER_HEALTHCHECK_PORT两个环境变量打开健康检查;- 子进程启动后轮询
http://localhost:<port>/health,返回 200 才认为就绪(wait_for_worker_health,默认最多 25 次、每次间隔 1 秒); - 测试结束时以进程组为单位发送
TERM,超时则KILL,保证不残留孤儿进程。
examples/spec_helper.rb 则定义了测试公共设施:session 级共享的Hatchet::Client.new(debug: true)(通过RSpec.configuration.hatchet_client与hatchet辅助方法访问),以及wait_for_running_status轮询辅助函数——它持续调用client.runs.get_details(run_id),直到运行进入RUNNING状态,并在404时静默重试(因为运行记录可能尚未可见)。这套"共享客户端 + 轮询辅助"的模式是所有test_*_spec.rb的通用骨架。
5. 核心能力逐个击破:示例源码深度解读
5.1 任务与持久化任务(Task / Durable Task)
simple/worker.rb 同时展示了普通任务与持久化任务的差异:
SIMPLE = HATCHET.task(name: "simple") do |input, ctx| { "result" => "Hello, world!" } end SIMPLE_DURABLE = HATCHET.durable_task(name: "simple_durable") do |input, ctx| result = SIMPLE.run(input) # 在 durable 任务内部同步调用另一个任务 { "result" => result["result"] } end- 普通任务:执行期间依赖 Worker 进程存活,进程中断则执行状态丢失;
- 持久化任务(durable_task):执行历史被写入后端(见 pkg/repository/durable_events.go 的事件存储实现),进程崩溃后可恢复,且支持
ctx.sleep_for、ctx.wait_for等长时间等待原语。
5.2 DAG 依赖编排
dag/worker.rb 展示了用parents:声明任务依赖、由引擎自动推导执行顺序的典型 DAG 写法:
DAG_WORKFLOW = HATCHET.workflow(name: "DAGWorkflow") STEP1 = DAG_WORKFLOW.task(:step1, execution_timeout: 5) do |input, ctx| { "random_number" => rand(1..100) } end STEP2 = DAG_WORKFLOW.task(:step2, execution_timeout: 5) do |input, ctx| { "random_number" => rand(1..100) } end # step3 依赖 step1、step2 都完成 DAG_WORKFLOW.task(:step3, parents: [STEP1, STEP2]) do |input, ctx| one = ctx.task_output(STEP1)["random_number"] two = ctx.task_output(STEP2)["random_number"] { "sum" => one + two } end # parents 既可用任务对象,也可用符号 :step3 DAG_WORKFLOW.task(:step4, parents: [STEP1, :step3]) do |input, ctx| puts ctx.task_output(STEP1).inspect, ctx.task_output(:step3).inspect { "step4" => "step4" } end要点归纳:
HATCHET.workflow(name: ...)返回工作流对象,workflow.task(...)在其上定义步骤;parents:接受任务对象或符号名,两种引用方式等价;ctx.task_output(step)按依赖读取上游任务输出,是 DAG 数据传递的标准手段;execution_timeout以秒为单位限制单步执行时长,超时行为可配合 timeout/ 示例验证。
5.3 事件触发与过滤
events/worker.rb 演示事件驱动的任务触发;events/event.rb 与 events/filter.rb 则分别演示事件创建与按表达式过滤订阅。事件是 Hatchet 解耦生产端与消费端的核心机制(底层消息队列实现见 internal/msgqueue 的 NATS/Postgres/RabbitMQ 三种适配器),结合hatchet_client.rb中的HATCHET.events.create(key:, data:),即可构成"事件生产者 + 事件驱动 Worker"的完整链路。
5.4 持久化工作流:Sleep 与事件等待
durable/worker.rb 是示例仓库中信息量最大的文件之一,展示了三个关键原语:
a)ctx.sleep_for—— 可恢复的定时休眠
DURABLE_WORKFLOW.durable_task(:durable_task, execution_timeout: 60) do |_input, ctx| ctx.sleep_for(duration: DURABLE_SLEEP_TIME) # DURABLE_SLEEP_TIME = 5 puts 'Sleep finished' ... endb)ctx.wait_for+ 条件组合 —— 等待事件或定时器
ctx.wait_for( 'event', Hatchet::UserEventCondition.new(event_key: DURABLE_EVENT_KEY, expression: 'true') )以及“或”条件组(任意一个满足即返回,返回值中可拿到命中的 key 与 event_id):
wait_result = ctx.wait_for( SecureRandom.hex(16), Hatchet.or_( Hatchet::SleepCondition.new(DURABLE_SLEEP_TIME), Hatchet::UserEventCondition.new(event_key: DURABLE_EVENT_KEY) ) ) key = wait_result.keys.first event_id = wait_result[key].keys.firstc) 子任务异步调用与错误传播
ERROR_RAISING_DURABLE_PARENT = HATCHET.durable_task(name: 'error-raising-durable-parent', execution_timeout: 30) do |input, ctx| ref = ERROR_RAISING_TASK.run_no_wait(input) # 异步触发子任务 begin ref.result # 阻塞获取结果,子任务异常在此抛出 rescue StandardError => e child_raised = true child_error_str = e.message end { 'child_raised' => child_raised, 'child_error_str' => child_error_str, 'child_run_external_id' => ref.workflow_run_id, 'parent_run_external_id' => ctx.workflow_run_id } end这展示了一个非常重要的工程模式:父任务对子任务失败的显式捕获与结构化上报。run_no_wait返回任务引用,ref.result同步等待并重抛子任务异常,ref.workflow_run_id与ctx.workflow_run_id可关联父子运行。配套的 durable_event/、durable_sleep/、durable_eviction/ 示例分别深化了事件驱动恢复、sleep 恢复与事件驱逐策略。
5.5 并发控制、重试与限流
示例仓库对并发语义的覆盖最为细致,全部位于concurrency_*系列目录:
| 示例目录 | 并发策略 | 说明 |
|---|---|---|
concurrency_limit/ | 并发上限 | 同一策略键下最多 N 个运行并行 |
concurrency_limit_rr/ | 轮询限流 | 多个键之间轮流分配槽位 |
concurrency_cancel_in_progress/ | 取消进行中 | 新运行到达时取消正在执行的运行 |
concurrency_cancel_newest/ | 取消最新 | 保留最旧、取消后到者 |
concurrency_cancel_queued_except_newest/ | 取消排队(留最新) | 取消排队项但保留最新运行 |
concurrency_cancel_queued_except_oldest/ | 取消排队(留最旧) | 取消排队项但保留最早运行 |
concurrency_multiple_keys/ | 多策略键 | 同一任务同时按多个维度限流 |
concurrency_workflow_level/ | 工作流级并发 | 在整条工作流层面施加限制 |
concurrency_shared/、concurrency_dynamic/ | 共享 / 动态键 | 跨工作流共享限制或运行时动态计算键 |
配合 rate_limit/、priority/、retries/、non_retryable/ 与 timeout/(含REFRESH_TIMEOUT_WF超时续期演示),构成了完整的"并发 + 优先级 + 重试 + 超时"质量保障组合。这些策略最终在服务端由 pkg/scheduling/v1 的调度器与 pkg/repository/scheduler_concurrency.go 落地执行。
5.6 调度、Webhook 与流式输出
- 定时调度:cron/worker.rb 与 scheduled/worker.rb 定义按 CRON 表达式与固定时间点触发的任务;cron/programatic_sync.rb 与 scheduled/programatic_sync.rb 演示通过 API 编程式同步调度配置。
- Webhook:webhooks/ 与 webhook_with_scope/ 演示任务如何被外部 HTTP 回调触发(含作用域限定与静态载荷两种模式),对应服务端实现 pkg/repository/webhooks.go。
- 流式输出:streaming/ 提供
async_stream.rb与worker.rb,演示任务向客户端推送实时增量结果的能力。
5.7 依赖注入与序列化
- dependency_injection/ 演示如何向任务上下文注入外部服务依赖(
ASYNC_TASK_WITH_DEPS、SYNC_TASK_WITH_DEPS等变体),避免任务与外部组件硬编码耦合; - serde/ 与 dataclasses/ 演示输入输出数据的序列化/反序列化与结构化类型定义,配合 pkg/repository/jsonb.go 可理解数据在后端的 JSONB 存储形态。
6. 端到端测试:示例即验证
示例仓库不仅是教学代码,更是一套可运行的回归测试集。以 simple/test_simple_spec.rb 为代表的test_*_spec.rb遵循统一范式:
spec_helper在 suite 启动时创建共享 Hatchet 客户端;HatchetWorkerFixture.with_worker(...)拉起 Worker 子进程并等待健康检查通过;- 通过客户端触发运行,用
wait_for_running_status或自定义轮询等待目标状态; - 断言运行状态、任务输出与副作用。
测试主题覆盖 batch_assign/、fanout/、idempotency/(幂等触发)、return_exceptions/(批量任务返回异常而不中断)、runtime_affinity/(运行时亲和性)、unit_testing/(脱离后端对任务做纯单元测试)等场景。运行全套测试:
bundle exec rspec7. 快速索引:按需求直达示例
| 你想解决的问题 | 直接阅读的示例 |
|---|---|
| 第一次跑通 SDK | hatchet_client.rb、quickstart/ |
| 定义任务并注册 Worker | simple/、worker.rb |
| 任务之间存在依赖 | dag/ |
| 用事件驱动任务 | events/ |
| 需要长时间等待/可恢复执行 | durable/、durable_event/、durable_sleep/ |
| 限制并发/取消策略 | concurrency_*全系列、rate_limit/ |
| 失败重试与超时 | retries/、non_retryable/、timeout/ |
| 定时/计划触发 | cron/、scheduled/ |
| 外部 HTTP 回调 | webhooks/、webhook_with_scope/ |
| 订阅实时输出 | streaming/ |
| 写测试验证工作流 | unit_testing/、各test_*_spec.rb、worker_fixture.rb |
8. 从示例到生产:与仓库源码的对照阅读建议
示例代码的价值在对照底层实现后会被放大数倍,建议按以下映射深入:
- 任务/工作流 API 契约:Ruby 侧
HATCHET.task / HATCHET.workflow / HATCHET.worker的完整签名与配置项,见 sdks/ruby/src(Worker 与工作流定义入口)及 sdks/typescript/src 同构 API; - 服务端执行引擎:调度与并发策略落点在 pkg/scheduling/v1 与 pkg/repository/scheduler_concurrency.go;运行状态机见 pkg/statusutils/status.go;
- 持久化事件存储:durable 示例背后的日志存储见 pkg/repository/durable_events.go,服务端侧的事件处理入口在 internal/services/ingestor;
- 消息队列抽象:事件/任务分发依赖 internal/msgqueue,生产环境通常走 NATS 或 Postgres 适配器;
- 配置加载:SDK 连接参数(token、server URL、健康检查端口等)的解析逻辑见 pkg/config/client 与 pkg/config/loader。
9. 小结
sdks/ruby/examples用一份极简 README 挂载了一整座 Ruby SDK 能力图谱:从bundle install到bundle exec ruby hatchet_client.rb的 Hello World,再到覆盖 DAG、持久化、并发策略、调度、Webhook 与测试的上百个真实示例。它既是新手最快上手的路径,也是老手校验 SDK 行为、贡献代码时的回归测试集。按本文索引对照源码逐例运行,你即可在最短时间内掌握 Hatchet Ruby SDK 的全部核心用法,并将其迁移到自己的生产工作流中。
【免费下载链接】hatchet🪓 An orchestration engine for background tasks, AI agents, and durable workflows项目地址: https://gitcode.com/GitHub_Trending/ha/hatchet
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考