十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Hatchet Python SDK Scheduled Client 指南:定时工作流的创建、查询、重调度与批量管理

Hatchet Python SDK Scheduled Client 指南:定时工作流的创建、查询、重调度与批量管理 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 接口创建定时运行createaio_createscheduled_workflow_run_createWorkflowRunApi重调度updateaio_updateworkflow_scheduled_updateWorkflowApi删除单个deleteaio_deleteworkflow_scheduled_deleteWorkflowApi查询单个getaio_getworkflow_scheduled_getWorkflowApi列表查询listaio_listworkflow_scheduled_listWorkflowApi批量删除bulk_deleteaio_bulk_deleteworkflow_scheduled_bulk_deleteWorkflowApi批量重调度bulk_updateaio_bulk_updateworkflow_scheduled_bulk_updateWorkflowApi所有异步方法的实现都基于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_namestr要调度的工作流名称。SDK 会自动通过client_config.apply_namespace(workflow_name)加上命名空间前缀trigger_atdatetime.datetime触发时间点建议使用带时区信息的datetimeUTCinputJSONSerializableMapping定时运行时的工作流输入数据JSON 可序列化字典additional_metadataJSONSerializableMapping与该次未来运行关联的附加元数据键值对可用于后续过滤查询返回ScheduledWorkflows对象详见下文返回模型一节其中scheduled_run.metadata.id即该定时运行触发器的 ID后续所有操作都以它为句柄。同步示例来自官方示例 programatic-sync.pyfrom datetime import datetime, timedelta, timezone from hatchet_sdk import Hatchet hatchet Hatchet() scheduled_run hatchet.scheduled.create( workflow_namesimple-workflow, trigger_atdatetime.now(tztimezone.utc) timedelta(seconds10), 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(tztimezone.utc)生成带时区的时间戳避免本地时区与服务器时区不一致导致触发时间偏移。additional_metadata建议放入业务维度的标识如customer_id后面可以用它做批量筛选。异步示例来自官方示例 programatic-async.pyimport 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_namesimple-workflow, trigger_atdatetime.now(tztimezone.utc) timedelta(seconds10), 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_getdef get(self, scheduled_id: str) - ScheduledWorkflows:按定时运行触发器 ID精确获取一个定时工作流返回完整的ScheduledWorkflows实例。scheduled_run hatchet.scheduled.get(scheduled_idscheduled_run.metadata.id)列表查询 list / aio_listdef 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排序字段ScheduledWorkflowsOrderByFieldorder_by_direction排序方向WorkflowRunOrderByDirection最简单用法直接不传参数scheduled_runs hatchet.scheduled.list()带过滤的用法scheduled_runs hatchet.scheduled.list( workflow_idworkflow_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_updatedef update( self, scheduled_id: str, trigger_at: datetime.datetime, ) - ScheduledWorkflows:将已创建的定时运行改期到新的触发时间返回更新后的ScheduledWorkflows。hatchet.scheduled.update( scheduled_idscheduled_run.metadata.id, trigger_atdatetime.now(tztimezone.utc) timedelta(hours1), )⚠️注意官方 docstring 明确指出服务端在以下两种情况下可能拒绝重调度该定时运行已经触发已变成实际的 workflow run该定时运行是通过**代码定义code definition**创建的而非通过 API 创建。因此update更适合在触发时间尚未到达前、且确认调度源为 API 创建时使用。六、删除delete / aio_deletedef delete(self, scheduled_id: str) - None:按 ID 删除一个定时工作流运行返回None。hatchet.scheduled.delete(scheduled_idscheduled_run.metadata.id)删除是不可逆操作建议在批量清理场景如下线某类任务中配合list 过滤条件先确认目标再删除。七、批量操作bulk_delete / bulk_update当定时任务规模变大时逐个调用效率低下此时应使用批量接口。批量删除 bulk_deletedef 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_idworkflow_id, additional_metadata{customer_id: customer-a}, )⚠️限制若既没有scheduled_ids也没有任何过滤字段会抛出ValueErrorbulk_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_updatedef 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(tztimezone.utc) timedelta(hours2)), ] )源码中会将元组形式自动转换为ScheduledWorkflowsBulkUpdateItem(id..., triggerAt...)随后统一封装进ScheduledWorkflowsBulkUpdateRequest发送见 features/scheduled.py。返回ScheduledWorkflowsBulkUpdateResponse同样包含更新的 ID 与逐项错误。八、返回模型与状态枚举ScheduledWorkflowscreate、update、get的返回值类型为ScheduledWorkflows模型定义见 clients/rest/models/scheduled_workflows.py主要字段字段说明metadata通用资源元信息APIResourceMeta其中id即定时运行触发器 IDtenant_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优先级约束在 13 之间注意workflow_run_id字段约束为 36 位 UUID 字符串若定时任务尚未触发该字段及其关联字段为None。ScheduledRunStatus 状态枚举statuses过滤参数和workflow_run_status使用 scheduled_run_status.py 中的枚举取值如下PENDING, RUNNING, SUCCEEDED, FAILED, CANCELLED, QUEUED, SCHEDULED九、底层实现与调用链从源码结构可以梳理出完整的调用链证据见 features/scheduled.pyhatchet.scheduled.xxx │ ▼ ScheduledClient继承 BaseRestClient │ ├─ create → WorkflowRunApi.scheduled_workflow_run_create ├─ update/delete → WorkflowApi.workflow_scheduled_update / _delete ├─ get/list → WorkflowApi.workflow_scheduled_get / _listtenacity 重试 └─ bulk_* → WorkflowApi.workflow_scheduled_bulk_*携带 filter 或 ID 列表 │ ▼ OpenAPI 生成的 REST Clienthatchet_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_namesimple-workflow, trigger_atdatetime.now(tztimezone.utc) timedelta(seconds10), input{data: simple-workflow-data}, additional_metadata{customer_id: customer-a}, ) id scheduled_run.metadata.id # 2. 重调度改到 1 小时后 hatchet.scheduled.update( scheduled_idid, trigger_atdatetime.now(tztimezone.utc) timedelta(hours1), ) # 3. 查询 scheduled_runs hatchet.scheduled.list( workflow_idworkflow_id, statuses[ScheduledRunStatus.SCHEDULED], additional_metadata{customer_id: customer-a}, ) scheduled_run hatchet.scheduled.get(scheduled_idid) # 4. 批量重调度统一延后 2 小时 hatchet.scheduled.bulk_update([(id, datetime.now(tztimezone.utc) timedelta(hours2))]) # 5. 批量删除按显式 ID hatchet.scheduled.bulk_delete(scheduled_ids[id]) # 6. 删除单个 hatchet.scheduled.delete(scheduled_idid)异步流程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_namesimple-workflow, trigger_atdatetime.now(tztimezone.utc) timedelta(seconds10), 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_idscheduled_id) # 单个 await hatchet.scheduled.aio_delete(scheduled_idscheduled_id) # 删除 asyncio.run(manage_scheduled())完整可运行代码可在仓库中查看programatic-sync.py 与 programatic-async.py。十一、最佳实践与注意事项时间务必带时区trigger_at建议使用datetime.now(tztimezone.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.pySDK 根对象挂载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),仅供参考
返回列表