Hatchet Python SDK 中 WorkflowsClient 完全指南:工作流声明的程序化查询、版本管理与暂停/恢复

发布时间:2026/9/16 16:15:09
Hatchet Python SDK 中 WorkflowsClient 完全指南:工作流声明的程序化查询、版本管理与暂停/恢复 Hatchet Python SDK 中 WorkflowsClient 完全指南工作流声明的程序化查询、版本管理与暂停/恢复【免费下载链接】hatchet An orchestration engine for background tasks, AI agents, and durable workflows项目地址: https://gitcode.com/GitHub_Trending/ha/hatchet本篇指南聚焦 Hatchet Python SDK 的WorkflowsClient功能客户端之一基于仓库中 features/workflows.py 的源码实现完整讲解其查询、列表、版本、删除、暂停与恢复等 API 的参数、返回类型与底层调用链。读完本文你将能够用 Python 代码对工作流“声明”进行全生命周期管理并理解每个方法背后发送的 REST 请求与重试策略。需要首先明确一个核心概念WorkflowsClient管理的是工作流的声明declaration本身而不是某一次具体执行。源码中的类注释写道workflows are the declaration,notthe individual runs. If youre looking for runs, use theRunsClientinstead.如果你要操作某次运行run的取消、重放、状态查询等应使用hatchet.runsRunsClient而本文的WorkflowsClient负责“这个工作流是什么、有哪些版本、是否处于暂停状态”这类声明层面的操作。如何获取 WorkflowsClientWorkflowsClient在 sdks/python/hatchet_sdk/features/workflows.py 中定义继承自BaseRestClientclass WorkflowsClient(BaseRestClient): The workflows client is a client for managing workflows programmatically within Hatchet. Note that workflows are the declaration, _not_ the individual runs. If youre looking for runs, use the RunsClient instead. 在Hatchet客户端顶层它通过workflows属性暴露。参见 hatchet.pyproperty def workflows(self) - WorkflowsClient: The workflows client is a client for managing workflows programmatically within Hatchet. ... return self._client.workflows因此标准用法是from hatchet_sdk import Hatchet hatchet Hatchet() # 从环境变量读取 tenant / API 地址 / token client hatchet.workflows # WorkflowsClient 实例Hatchet()的构造依赖HATCHET_URL、HATCHET_TENANT_ID、HATCHET_TOKEN等环境变量客户端配置在client_config中下文list方法会用到其中的tenant_id与tenacity重试配置。核心 API 一览WorkflowsClient公开 6 组方法每组都提供同步与aio_前缀的异步两个版本异步版本通过asyncio.to_thread包装同名同步方法行为完全一致方法异步版本返回类型作用get(workflow_id)aio_getWorkflow按 ID 获取单个工作流list(workflow_name, limit, offset)aio_listWorkflowList按名称过滤、分页列出工作流get_version(workflow_id, version)aio_get_versionWorkflowVersion获取工作流的指定版本默认最新版本delete(workflow_id)aio_deleteNone永久删除工作流及其全部数据pause(workflow_id, queue_ttl, ...)aio_pauseWorkflow暂停工作流配置队列 TTL 与触发行为unpause(workflow_id)aio_unpauseWorkflow恢复unpause工作流get按 ID 获取工作流workflow: Workflow client.get(wf-0123456789abcdef) # 或异步 workflow await client.aio_get(wf-0123456789abcdef)参数workflow_id工作流的字符串 ID注意不是名称。返回Workflow对象包含name以及metadata其中metadata.id、metadata.created_at、metadata.updated_at是常用字段见下文实战示例。从源码看workflows.py该方法通过WorkflowApi.workflow_get发起 REST 调用并用tenacity_retry包裹def get(self, workflow_id: str) - Workflow: with self.client() as client: workflow_get tenacity_retry( self._wa(client).workflow_get, self.client_config.tenacity ) return workflow_get(workflow_id)这里有两个值得注意的实现细节self.client()是上下文管理器负责创建生成式 RESTApiClient_wa(client)把ApiClient适配为WorkflowApi来自hatchet_sdk.clients.rest.api.workflow_api该客户端基于 api-contracts/openapi/openapi.yaml 的 OpenAPI 契约生成。tenacity_retry(fn, self.client_config.tenacity)表示所有读/写调用都按客户端配置的 tenacity 重试策略执行瞬时网络错误会被自动重试。list分页列出工作流含 namespace 处理def list( self, workflow_name: str | None None, limit: int | None None, offset: int | None None, ) - WorkflowList:workflow_name按名称过滤不传则返回该 tenant 下全部工作流。limit本次返回的最大条数offset分页偏移量。返回WorkflowList其中rows为Workflow列表可能为None取用时建议写workflow_list.rows or []。源码workflows.py中有一个容易被忽略的细节——tenant 与 namespace 的处理return workflow_list( tenantself.client_config.tenant_id, limitlimit, offsetoffset, nameself.client_config.apply_namespace(workflow_name), )tenant 是隐式的查询范围由客户端配置中的tenant_id决定调用者无法也不应跨 tenant 查询apply_namespace如果客户端配置了 namespace 前缀会将其拼接到workflow_name上再做过滤。这意味着当你通过Hatchet配置了 namespace 时按名称过滤的入参会被自动加上前缀行为与 worker 侧注册工作流时的命名规则保持一致。get_version获取工作流版本def get_version( self, workflow_id: str, version: str | None None ) - WorkflowVersion:version为None时返回最新版本传入具体版本号则返回该历史版本。返回WorkflowVersion对象。Hatchet 的工作流是带版本的声明每次结构变更如新增/修改 task会产生新的 version。理解版本对你排查“旧运行用的是哪一版 DAG”这类问题非常有用。该方法同样带 tenacity 重试底层为workflow_version_get。delete永久删除危险操作def delete(self, workflow_id: str) - None:源码注释明确标注DANGEROUS: This will delete a workflow and all of its data删除是不可逆的会连带工作流的全部相关数据一并删除且该方法没有返回Workflow对象。与get/list不同delete的底层调用没有包裹 tenacity 重试workflows.py失败时会直接抛出异常——这在语义上是合理的删除类操作不应被静默重试。生产环境调用前务必先确认workflow_id正确。pause / unpause暂停与恢复工作流这是WorkflowsClient中最有业务价值的一组方法常用于变更窗口、下游故障隔离等场景。def pause( self, workflow_id: str, queue_ttl: timedelta, paused_workflow_cron_run_queue_behavior: Literal[DROP, QUEUE] QUEUE, paused_workflow_scheduled_run_queue_behavior: Literal[ DROP, QUEUE ] QUEUE, ) - Workflow:参数说明参数类型 / 默认值含义workflow_idstr要暂停的工作流 IDqueue_ttltimedelta必填暂停期间排队中的 run 在被丢弃drop前允许保留的存活时间paused_workflow_cron_run_queue_behaviorDROP \| QUEUE默认QUEUE暂停期间由cron 触发的 run 如何处置入队等待还是直接丢弃paused_workflow_scheduled_run_queue_behaviorDROP \| QUEUE默认QUEUE暂停期间由scheduled 触发的 run 如何处置语义上可以这样理解QUEUE默认暂停期间新触发的 cron/scheduled run 继续进入队列等unpause后执行DROP暂停期间新触发的直接丢弃不进入队列无论哪种行为已经排队且工作流处于暂停状态的 run都会受queue_ttl约束——超过 TTL 仍未被执行就会被丢弃。源码实现workflows.py展示了请求的构造过程ttl_expr timedelta_to_expr(queue_ttl) # timedelta - 时长表达式字符串 workflow_update_pause( workflow_id, WorkflowUpdateRequest( pausePauseWorkflowRequest( PauseWorkflowRequestPause( actionpause, pausedWorkflowCronRunQueueBehavior...(cron 行为), pausedWorkflowScheduledRunQueueBehavior...(scheduled 行为), pausedWorkflowQueueTTLttl_expr, ) ) ), )可以看到pause与unpause复用同一个 REST 端点workflow_update通过actionpause/actionunpause区分操作类型queue_ttl这个timedelta会先经timedelta_to_expr转换为时长表达式字符串再下发避免时区/格式歧义。unpause很简单只接收workflow_id恢复后队列中未过期的 run 会按既有优先级继续执行def unpause(self, workflow_id: str) - Workflow: # 底层发送 PauseWorkflowRequestUnpause(actionunpause)两者都返回更新后的Workflow对象你可以据此确认暂停状态已生效。实战示例来自仓库 examples列举并打印工作流仓库示例 examples/api/api.py 展示了最典型的list用法from hatchet_sdk import Hatchet hatchet Hatchet() def main() - None: workflow_list hatchet.workflows.list() rows workflow_list.rows or [] for workflow in rows: print(workflow.name) print(workflow.metadata.id) print(workflow.metadata.created_at) print(workflow.metadata.updated_at) if __name__ __main__: main()注意rows可能为None这里用or []做了防御。name用于展示metadata.id则是后续get/pause/delete等操作所需的 ID——“先 list 拿 ID再按 ID 操作”是WorkflowsClient的标准工作模式。与 RunsClient 组合批量取消 / 重放声明层面的 ID 也是操作运行数据的必要输入。以 examples/bulk_operations/cancel.py 为例先通过workflows.list()拿到工作流 ID再交给hatchet.runs按 ID 或按过滤条件批量取消workflows hatchet.workflows.list() workflow workflows.rows[0] # 方式一按 run ID 批量取消 workflow_runs hatchet.runs.list(workflow_ids[workflow.metadata.id]) workflow_run_ids [r.metadata.id for r in workflow_runs.rows] hatchet.runs.bulk_cancel(BulkCancelReplayOpts(idsworkflow_run_ids)) # 方式二按过滤条件批量取消时间范围 状态 metadata hatchet.runs.bulk_cancel(BulkCancelReplayOpts( filtersRunFilter( sincedatetime.today() - timedelta(days1), untildatetime.now(tztimezone.utc), statuses[V1TaskStatus.RUNNING], workflow_ids[workflow.metadata.id], additional_metadata{key: value}, ) ))examples/bulk_operations/replay.py 结构完全相同只是把bulk_cancel换成bulk_replay将选中的 run 重新入队执行。这两个例子很好地体现了WorkflowsClient与RunsClient的分工前者定位“对哪个工作流操作”后者执行“对运行做什么操作”。与 Workflow 声明对象的联动SDK 中通过hatchet.workflow装饰器声明的工作流对象BaseWorkflow实现位于 runnables/workflow.py本身就内嵌了对WorkflowsClient的便捷封装无需手写 IDworkflow.idcached_property内部执行self._client.workflows.list(workflow_nameself.name)按名称查到 ID 并缓存查不到会抛出ValueError(No id found for {name})workflow.delete()/aio_delete()源码等价于client.workflows.delete(self.id)workflow.pause(queue_ttl, ...)/aio_pause(...)源码参数与WorkflowsClient.pause完全一致只是省略了workflow_idworkflow.unpause()/aio_unpause()源码同理。因此如果你的代码里已经持有工作流声明对象优先使用这些便捷方法只有在“只有 ID、没有声明对象”的运维/脚本场景例如上面 examples 中的做法才直接走WorkflowsClient。实现细节小结REST 客户端与重试策略从源码结构看WorkflowsClient是典型的“功能客户端 生成式 REST 客户端”分层结构功能层features/workflows.py 定义业务语义暂停行为、版本概念、危险操作警示并处理timedelta转换、namespace 拼接等 SDK 级细节REST 生成层hatchet_sdk.clients.rest.api.workflow_api.WorkflowApi/workflow_run_api.WorkflowRunApi由 OpenAPI 契约 api-contracts/openapi/openapi.yaml 及路径定义api-contracts/openapi/paths/下按资源拆分的 yaml生成方法如workflow_get、workflow_list、workflow_version_get、workflow_delete、workflow_update与上述业务方法一一对应重试层除delete外的方法都经tenacity_retry(fn, client_config.tenacity)包裹重试次数、退避策略由客户端配置client_config.tenacity统一控制。理解这一分层后当你需要排查某次调用失败时可以按“业务方法 → 对应WorkflowApi方法 → REST 端点”的路径逐层定位。使用建议与注意事项ID 与名称不要混用get/pause/unpause/delete只接受 ID只有list支持按名称过滤。声明对象上的workflow.id会自动解析裸调用时请先list拿 ID。delete不可逆且不重试确认 ID 后再调用生产脚本建议先get打印name做二次确认。pause的queue_ttl必填它决定了暂停期间排队 run 的“最长存活时间”DROP/QUEUE两个行为参数按触发来源cron / scheduled分开配置默认均为QUEUE。异步优先在 FastAPI 等 asyncio 服务中请使用aio_*版本避免在事件循环里阻塞调用同步方法。声明 vs 运行涉及具体某次运行的查询、取消、重放一律使用hatchet.runsRunsClientWorkflowsClient不承担该职责。参考文件sdks/python/hatchet_sdk/features/workflows.py ——WorkflowsClient完整实现sdks/python/hatchet_sdk/hatchet.py ——Hatchet.workflows属性入口sdks/python/hatchet_sdk/runnables/workflow.py —— 工作流声明对象上的便捷封装id/delete/pause/unpausesdks/python/examples/api/api.py —— 列举工作流并打印元数据sdks/python/examples/bulk_operations/cancel.py / sdks/python/examples/bulk_operations/replay.py —— 声明 ID 与 RunsClient 的组合用法sdks/python/docs/feature-clients/workflows.md —— 官方 API 参考页mkdocstrings 自动生成锚定features.workflows.WorkflowsClient【免费下载链接】hatchet An orchestration engine for background tasks, AI agents, and durable workflows项目地址: https://gitcode.com/GitHub_Trending/ha/hatchet创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考