Hatchet Ruby 示例仓库完全指南:从 Hello World 到生产级工作流实战 📅 发布时间:2026/9/16 14:07:34 👁 浏览次数: 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.rbRSpec 端到端验证基础篇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 WorldREADME 指明的入门入口README 中给出的唯一可运行示例是bundle exec ruby hatchet_client.rb这里的hatchet_client.rb是 sdks/ruby/examples/hatchet_client.rb它演示了 SDK 客户端三个最基础的 APIrequire 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.rbrequire 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.rbrequire_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 高级”分层注释最后创建并启动 WorkerHATCHET 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 Fixtureexamples/worker_fixture.rb 为端到端测试提供进程级基础设施HatchetWorkerFixture.with_worker(command, healthcheck_port:)以子进程方式启动 Worker并通过HATCHET_CLIENT_WORKER_HEALTHCHECK_ENABLEDtrue与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 Tasksimple/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_idwait_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 编程式同步调度配置。Webhookwebhooks/ 与 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. 快速索引按需求直达示例你想解决的问题直接阅读的示例第一次跑通 SDKhatchet_client.rb、quickstart/定义任务并注册 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.rb、worker_fixture.rb8. 从示例到生产与仓库源码的对照阅读建议示例代码的价值在对照底层实现后会被放大数倍建议按以下映射深入任务/工作流 API 契约Ruby 侧HATCHET.task / HATCHET.workflow / HATCHET.worker的完整签名与配置项见 sdks/ruby/srcWorker 与工作流定义入口及 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),仅供参考