Hatchet Ruby 示例仓库完全指南:从 Hello World 到生产级工作流实战
2026/9/16 14:07:32 网站建设 项目流程

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 官方示例的入口文档。它声明了两个核心事实:

  1. 该目录演示了 Hatchet Ruby SDK 的典型用法
  2. 环境准备只需要一条命令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_TOKENHATCHET_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.createname参数对应工作流注册名,input为 JSON 输入;runs.poll返回带metadata.idstatus的运行句柄。

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 } end

quickstart/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 是所有示例的“总装车间”。它的结构非常值得学习:

  1. require_relative加载每个主题目录的 worker 文件(如simple/workerdag/workerdurable/worker),这些文件在被加载时即完成HATCHET.task(...)/HATCHET.workflow(...)的定义;
  2. 将定义好的常量(SIMPLEDAG_WORKFLOWDURABLE_WORKFLOW……)汇总进ALL_WORKFLOWS数组,并按“Tier 1 基础 → Tier 2 并发 → Tier 3 编排 → Tier 4-5 高级”分层注释;
  3. 最后创建并启动 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_NAME

if __FILE__ == $PROGRAM_NAME保证了该文件被 require 时不启动 Worker、直接运行时才启动——这正是 4.1 节聚合加载的前提。因此你可以单独运行任一示例:

bundle exec ruby simple/worker.rb

4.3 测试基建:Worker Fixture

examples/worker_fixture.rb 为端到端测试提供进程级基础设施:

  • HatchetWorkerFixture.with_worker(command, healthcheck_port:)以子进程方式启动 Worker,并通过HATCHET_CLIENT_WORKER_HEALTHCHECK_ENABLED=trueHATCHET_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_clienthatchet辅助方法访问),以及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_forctx.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' ... end

b)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.first

c) 子任务异步调用与错误传播

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_idctx.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.rbworker.rb,演示任务向客户端推送实时增量结果的能力。

5.7 依赖注入与序列化

  • dependency_injection/ 演示如何向任务上下文注入外部服务依赖(ASYNC_TASK_WITH_DEPSSYNC_TASK_WITH_DEPS等变体),避免任务与外部组件硬编码耦合;
  • serde/ 与 dataclasses/ 演示输入输出数据的序列化/反序列化与结构化类型定义,配合 pkg/repository/jsonb.go 可理解数据在后端的 JSONB 存储形态。

6. 端到端测试:示例即验证

示例仓库不仅是教学代码,更是一套可运行的回归测试集。以 simple/test_simple_spec.rb 为代表的test_*_spec.rb遵循统一范式:

  1. spec_helper在 suite 启动时创建共享 Hatchet 客户端;
  2. HatchetWorkerFixture.with_worker(...)拉起 Worker 子进程并等待健康检查通过;
  3. 通过客户端触发运行,用wait_for_running_status或自定义轮询等待目标状态;
  4. 断言运行状态、任务输出与副作用。

测试主题覆盖 batch_assign/、fanout/、idempotency/(幂等触发)、return_exceptions/(批量任务返回异常而不中断)、runtime_affinity/(运行时亲和性)、unit_testing/(脱离后端对任务做纯单元测试)等场景。运行全套测试:

bundle exec rspec

7. 快速索引:按需求直达示例

你想解决的问题直接阅读的示例
第一次跑通 SDKhatchet_client.rbquickstart/
定义任务并注册 Workersimple/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.rbworker_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 installbundle 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询