Hatchet Python SDK Scheduled Client 指南:定时工作流的创建、查询、重调度与批量管理
【免费下载链接】hatchet🪓 An orchestration engine for background tasks, AI agents, and durable workflows项目地址: https://gitcode.com/GitHub_Trending/ha/hatchet
导读
本文围绕 Hatchet 官方 Python SDK 的Scheduled Client(定时任务客户端)展开,讲解如何通过hatchet.scheduled对"一次性定时触发的 workflow run"进行全生命周期管理——包括创建、查询、重调度(reschedule)、删除以及批量操作。文档主体对应仓库中的 scheduled.md,其渲染内容来源于 ScheduledClient 类实现;读完本文,你将掌握定时工作流从创建到清理的完整实战方案,并理解其底层 REST 调用链与异步实现原理。
一、Scheduled Client 是什么
Scheduled Client 是 Hatchet Python SDK 中负责管理一次性定时调度工作流运行(scheduled workflow run)的客户端。它与其他 feature client 一样通过 SDK 根对象挂载:
from hatchet_sdk import Hatchet hatchet = Hatchet() hatchet.scheduled # -> ScheduledClient从源码看,ScheduledClient实例在 client.py 中随Hatchet初始化创建,并由 hatchet.py 以scheduled属性对外暴露。类定义位于 features/scheduled.py,继承自BaseRestClient,本质上是对 Hatchet REST API 中workflow_scheduled_*系列接口的封装。
定位说明:这是一个"逃生舱"
需要特别指出的是,官方在create方法的 docstring 中明确建议:
优先使用
Workflow.run(及其同类方法)来触发工作流,本方法定位为 escape hatch(逃生舱)。
也就是说,常规的按需触发应走 runnables.md 中描述的工作流调用路径;而当你有明确的未来时间点触发一次的需求(如"10 秒后执行""明早 8 点执行")时,才使用 Scheduled Client。若你需要周期性重复执行,应使用 Cron Client(参见 cron.md)或在 workflow 定义中声明 cron 触发器,而不是用本客户端反复创建定时任务。
二、核心方法总览
ScheduledClient共提供 14 个方法,每个同步方法都有对应的aio_异步版本:
| 能力 | 同步方法 | 异步方法 | 底层 REST 接口 |
|---|---|---|---|
| 创建定时运行 | create | aio_create | scheduled_workflow_run_create(WorkflowRunApi) |
| 重调度 | update | aio_update | workflow_scheduled_update(WorkflowApi) |
| 删除单个 | delete | aio_delete | workflow_scheduled_delete(WorkflowApi) |
| 查询单个 | get | aio_get | workflow_scheduled_get(WorkflowApi) |
| 列表查询 | list | aio_list | workflow_scheduled_list(WorkflowApi) |
| 批量删除 | bulk_delete | aio_bulk_delete | workflow_scheduled_bulk_delete(WorkflowApi) |
| 批量重调度 | bulk_update | aio_bulk_update | workflow_scheduled_bulk_update(WorkflowApi) |
所有异步方法的实现都基于asyncio.to_thread包装同步方法(见 features/scheduled.py 中各aio_*方法),因此阻塞的 HTTP 调用不会卡住事件循环。
三、创建定时运行:create / aio_create
方法签名与参数
def create( self, workflow_name: str, trigger_at: datetime.datetime, input: JSONSerializableMapping, additional_metadata: JSONSerializableMapping, ) -> ScheduledWorkflows:| 参数 | 类型 | 说明 |
|---|---|---|
workflow_name | str | 要调度的工作流名称。SDK 会自动通过client_config.apply_namespace(workflow_name)加上命名空间前缀 |
trigger_at | datetime.datetime | 触发时间点,建议使用带时区信息的datetime(UTC) |
input | JSONSerializableMapping | 定时运行时的工作流输入数据(JSON 可序列化字典) |
additional_metadata | JSONSerializableMapping | 与该次未来运行关联的附加元数据键值对,可用于后续过滤查询 |
返回ScheduledWorkflows对象(详见下文"返回模型"一节),其中scheduled_run.metadata.id即该定时运行触发器的 ID,后续所有操作都以它为句柄。
同步示例
来自官方示例 programatic-sync.py:
from datetime import datetime, timedelta, timezone from hatchet_sdk import Hatchet hatchet = Hatchet() scheduled_run = hatchet.scheduled.create( workflow_name="simple-workflow", trigger_at=datetime.now(tz=timezone.utc) + timedelta(seconds=10), input={ "data": "simple-workflow-data", }, additional_metadata={ "customer_id": "customer-a", }, ) id = scheduled_run.metadata.id # the id of the scheduled run trigger要点:
trigger_at使用datetime.now(tz=timezone.utc)生成带时区的时间戳,避免本地时区与服务器时区不一致导致触发时间偏移。additional_metadata建议放入业务维度的标识(如customer_id),后面可以用它做批量筛选。
异步示例
来自官方示例 programatic-async.py:
import asyncio from datetime import datetime, timedelta, timezone from hatchet_sdk import Hatchet hatchet = Hatchet() async def create_scheduled() -> None: scheduled_run = await hatchet.scheduled.aio_create( workflow_name="simple-workflow", trigger_at=datetime.now(tz=timezone.utc) + timedelta(seconds=10), input={ "data": "simple-workflow-data", }, additional_metadata={ "customer_id": "customer-a", }, ) scheduled_run.metadata.id # the id of the scheduled run trigger asyncio.run(create_scheduled())四、查询定时运行:get / list
单个查询 get / aio_get
def get(self, scheduled_id: str) -> ScheduledWorkflows:按"定时运行触发器 ID"精确获取一个定时工作流,返回完整的ScheduledWorkflows实例。
scheduled_run = hatchet.scheduled.get(scheduled_id=scheduled_run.metadata.id)列表查询 list / aio_list
def list( self, offset: int | None = None, limit: int | None = None, workflow_id: str | None = None, parent_workflow_run_id: str | None = None, statuses: list[ScheduledRunStatus] | None = None, additional_metadata: JSONSerializableMapping | None = None, order_by_field: ScheduledWorkflowsOrderByField | None = None, order_by_direction: WorkflowRunOrderByDirection | None = None, ) -> ScheduledWorkflowsList:| 参数 | 说明 |
|---|---|
offset/limit | 分页参数,offset 为跳过的条数,limit 为返回条数上限 |
workflow_id | 按工作流 ID 过滤 |
parent_workflow_run_id | 按父工作流运行 ID 过滤(可用于子流程场景) |
statuses | 按定时运行状态列表过滤,取值见下文ScheduledRunStatus枚举 |
additional_metadata | 按附加元数据键值对过滤,内部经maybe_additional_metadata_to_kv归一化后传给服务端 |
order_by_field | 排序字段(ScheduledWorkflowsOrderByField) |
order_by_direction | 排序方向(WorkflowRunOrderByDirection) |
最简单用法直接不传参数:
scheduled_runs = hatchet.scheduled.list()带过滤的用法:
scheduled_runs = hatchet.scheduled.list( workflow_id="workflow_id", statuses=[ScheduledRunStatus.SCHEDULED], additional_metadata={"customer_id": "customer-a"}, )实现细节:list与get在源码中调用self._wa(client).workflow_scheduled_list/workflow_scheduled_get时,都包裹了tenacity_retry(..., self.client_config.tenacity)(见 features/scheduled.py),即对这两个读操作启用了基于 tenacity 的重试机制,提升在网络抖动下的可用性。
五、重调度:update / aio_update
def update( self, scheduled_id: str, trigger_at: datetime.datetime, ) -> ScheduledWorkflows:将已创建的定时运行改期到新的触发时间,返回更新后的ScheduledWorkflows。
hatchet.scheduled.update( scheduled_id=scheduled_run.metadata.id, trigger_at=datetime.now(tz=timezone.utc) + timedelta(hours=1), )⚠️注意:官方 docstring 明确指出,服务端在以下两种情况下可能拒绝重调度:
- 该定时运行已经触发(已变成实际的 workflow run);
- 该定时运行是通过**代码定义(code definition)**创建的,而非通过 API 创建。
因此update更适合在触发时间尚未到达前、且确认调度源为 API 创建时使用。
六、删除:delete / aio_delete
def delete(self, scheduled_id: str) -> None:按 ID 删除一个定时工作流运行,返回None。
hatchet.scheduled.delete(scheduled_id=scheduled_run.metadata.id)删除是不可逆操作,建议在批量清理场景(如下线某类任务)中配合list+ 过滤条件先确认目标再删除。
七、批量操作:bulk_delete / bulk_update
当定时任务规模变大时,逐个调用效率低下,此时应使用批量接口。
批量删除 bulk_delete
def bulk_delete( self, *, scheduled_ids: list[str] | None = None, workflow_id: str | None = None, parent_workflow_run_id: str | None = None, parent_step_run_id: str | None = None, statuses: list[ScheduledRunStatus] | None = None, additional_metadata: JSONSerializableMapping | None = None, ) -> ScheduledWorkflowsBulkDeleteResponse:两种使用方式(二选一,也可同时提供):
- 显式 ID 列表:直接传入
scheduled_ids; - 过滤条件:提供
workflow_id、parent_workflow_run_id、parent_step_run_id、additional_metadata中的一个或多个。
示例:
# 方式一:显式 ID hatchet.scheduled.bulk_delete(scheduled_ids=[id]) # 方式二:按过滤条件 hatchet.scheduled.bulk_delete( workflow_id="workflow_id", additional_metadata={"customer_id": "customer-a"}, )⚠️限制:
- 若既没有
scheduled_ids也没有任何过滤字段,会抛出ValueError:"bulk_delete requires either scheduled_ids or at least one filter field." statuses过滤目前不被批量删除支持:源码中会记录一条 warning 日志"The 'statuses' filter is not supported for bulk delete and will be ignored."(见 features/scheduled.py),传入也会被忽略,请勿依赖它筛选删除目标。- 返回
ScheduledWorkflowsBulkDeleteResponse,其中包含被删除的 ID 列表及逐项错误信息,便于部分失败时重试。
批量重调度 bulk_update
def bulk_update( self, updates: ( list[ScheduledWorkflowsBulkUpdateItem] | list[tuple[str, datetime.datetime]] ), ) -> ScheduledWorkflowsBulkUpdateResponse:updates支持两种形式:
(scheduled_id, trigger_at)元组列表——最简洁;ScheduledWorkflowsBulkUpdateItem对象列表——需要更精细控制时使用。
hatchet.scheduled.bulk_update( [ (id, datetime.now(tz=timezone.utc) + timedelta(hours=2)), ] )源码中会将元组形式自动转换为ScheduledWorkflowsBulkUpdateItem(id=..., triggerAt=...),随后统一封装进ScheduledWorkflowsBulkUpdateRequest发送(见 features/scheduled.py)。返回ScheduledWorkflowsBulkUpdateResponse,同样包含更新的 ID 与逐项错误。
八、返回模型与状态枚举
ScheduledWorkflows
create、update、get的返回值类型为ScheduledWorkflows(模型定义见 clients/rest/models/scheduled_workflows.py),主要字段:
| 字段 | 说明 |
|---|---|
metadata | 通用资源元信息(APIResourceMeta),其中id即定时运行触发器 ID |
tenant_id/workflow_id/workflow_version_id/workflow_name | 租户、工作流及其版本归属信息 |
trigger_at | 计划触发时间 |
input | 定时运行输入数据 |
additional_metadata | 附加元数据 |
workflow_run_id/workflow_run_created_at/workflow_run_name/workflow_run_status | 触发后对应实际 workflow run 的信息(未触发前为空) |
method | 创建方式(ScheduledWorkflowsMethod枚举) |
priority | 优先级,约束在 1~3 之间 |
注意workflow_run_id字段约束为 36 位 UUID 字符串,若定时任务尚未触发,该字段及其关联字段为None。
ScheduledRunStatus 状态枚举
statuses过滤参数和workflow_run_status使用 scheduled_run_status.py 中的枚举,取值如下:
PENDING, RUNNING, SUCCEEDED, FAILED, CANCELLED, QUEUED, SCHEDULED九、底层实现与调用链
从源码结构可以梳理出完整的调用链(证据见 features/scheduled.py):
hatchet.scheduled.xxx │ ▼ ScheduledClient(继承 BaseRestClient) │ ├─ create → WorkflowRunApi.scheduled_workflow_run_create ├─ update/delete → WorkflowApi.workflow_scheduled_update / _delete ├─ get/list → WorkflowApi.workflow_scheduled_get / _list(tenacity 重试) └─ bulk_* → WorkflowApi.workflow_scheduled_bulk_*(携带 filter 或 ID 列表) │ ▼ OpenAPI 生成的 REST Client(hatchet_sdk/clients/rest) │ ▼ Hatchet API Server(/api/v1 下的 workflow scheduled 端点)几个值得注意的实现事实:
- 命名空间自动注入:
create中工作流名会经过self.client_config.apply_namespace(workflow_name)处理,因此传入不带命名空间的裸名称即可,SDK 保证其落在当前租户/命名空间下。 - 异步只是线程桥接:所有
aio_*方法均通过asyncio.to_thread复用同步实现,没有独立的异步 HTTP 路径,这保证了同步与异步行为完全一致。 - 读操作有重试保护:
get与list使用tenacity_retry包装,配置来自self.client_config.tenacity;写操作(create/update/delete)则不做重试,避免重复创建或重复删除。 - 批量删除的过滤能力有限:
statuses在 bulk delete 中被显式忽略,其余过滤字段通过ScheduledWorkflowsBulkDeleteFilter模型承载。
十、完整示例:同步 + 异步流程串讲
同步流程
from datetime import datetime, timedelta, timezone from hatchet_sdk import Hatchet from hatchet_sdk.clients.rest.models.scheduled_run_status import ScheduledRunStatus hatchet = Hatchet() # 1. 创建:10 秒后触发 scheduled_run = hatchet.scheduled.create( workflow_name="simple-workflow", trigger_at=datetime.now(tz=timezone.utc) + timedelta(seconds=10), input={"data": "simple-workflow-data"}, additional_metadata={"customer_id": "customer-a"}, ) id = scheduled_run.metadata.id # 2. 重调度:改到 1 小时后 hatchet.scheduled.update( scheduled_id=id, trigger_at=datetime.now(tz=timezone.utc) + timedelta(hours=1), ) # 3. 查询 scheduled_runs = hatchet.scheduled.list( workflow_id="workflow_id", statuses=[ScheduledRunStatus.SCHEDULED], additional_metadata={"customer_id": "customer-a"}, ) scheduled_run = hatchet.scheduled.get(scheduled_id=id) # 4. 批量重调度:统一延后 2 小时 hatchet.scheduled.bulk_update([(id, datetime.now(tz=timezone.utc) + timedelta(hours=2))]) # 5. 批量删除:按显式 ID hatchet.scheduled.bulk_delete(scheduled_ids=[id]) # 6. 删除单个 hatchet.scheduled.delete(scheduled_id=id)异步流程
import asyncio from datetime import datetime, timedelta, timezone from hatchet_sdk import Hatchet hatchet = Hatchet() async def manage_scheduled() -> None: scheduled_run = await hatchet.scheduled.aio_create( workflow_name="simple-workflow", trigger_at=datetime.now(tz=timezone.utc) + timedelta(seconds=10), input={"data": "simple-workflow-data"}, additional_metadata={"customer_id": "customer-a"}, ) scheduled_id = scheduled_run.metadata.id await hatchet.scheduled.aio_list() # 列表 await hatchet.scheduled.aio_get(scheduled_id=scheduled_id) # 单个 await hatchet.scheduled.aio_delete(scheduled_id=scheduled_id) # 删除 asyncio.run(manage_scheduled())完整可运行代码可在仓库中查看:programatic-sync.py 与 programatic-async.py。
十一、最佳实践与注意事项
- 时间务必带时区:
trigger_at建议使用datetime.now(tz=timezone.utc),避免本地时区与服务器时区不一致造成触发偏差。 - 善用 additional_metadata:创建时写入业务标识(客户 ID、批次号等),后续
list/bulk_delete即可按此精准筛选,避免遍历全量数据。 - 优先考虑
Workflow.run:官方将ScheduledClient.create定位为 escape hatch;常规按需触发请走 runnables.md 中的工作流调用 API。 - 周期任务别用本客户端:需要 cron 表达式周期性执行时,使用 Cron Client 或在 workflow 定义中声明 cron 触发器。
- update 有前置条件:已经触发、或由代码定义创建的定时运行,重调度可能被服务端拒绝;请在触发前完成改期。
- bulk_delete 的 statuses 无效:该过滤参数会被忽略并打印 warning,请改用
scheduled_ids或其他过滤字段。 - 批量接口具备部分失败语义:
bulk_delete/bulk_update的响应包含逐项错误,建议对失败项记录并重试,保证最终一致。 - 读操作有内置重试:
get/list自带 tenacity 重试,网络抖动下更稳健;写操作无重试,避免副作用重复执行。
相关文档与源码索引
- 文档主体:sdks/python/docs/feature-clients/scheduled.md
- 客户端实现:sdks/python/hatchet_sdk/features/scheduled.py
- 同步示例:sdks/python/examples/scheduled/programatic-sync.py
- 异步示例:sdks/python/examples/scheduled/programatic-async.py
- 返回模型:sdks/python/hatchet_sdk/clients/rest/models/scheduled_workflows.py
- 状态枚举:sdks/python/hatchet_sdk/clients/rest/models/scheduled_run_status.py
- SDK 根对象挂载:sdks/python/hatchet_sdk/client.py 与 sdks/python/hatchet_sdk/hatchet.py
- 工作流触发(首选路径):sdks/python/docs/runnables.md
【免费下载链接】hatchet🪓 An orchestration engine for background tasks, AI agents, and durable workflows项目地址: https://gitcode.com/GitHub_Trending/ha/hatchet
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考