本指南涵盖 Python Workflows SDK,包括如何使用 Python 构建和创建 workflow 的说明。
WorkflowEntrypoint 是 Python workflow 的主入口点。它扩展 WorkflowEntrypoint 类并实现 run 方法。
from workers import WorkflowEntrypoint
class MyWorkflow(WorkflowEntrypoint):
async def run(self, event, step):
# steps here-
step.do(name=None, *, concurrent=False, config=None)— 用于在 workflow 中定义步骤的装饰器。name— 步骤的可选名称。如果省略,则使用函数名(func.__name__)。concurrent— 可选布尔值,指示此步骤的依赖项是否可以并发运行。config— 可选WorkflowStepConfig,用于配置步骤特定重试行为。作为 Python 字典传递,然后类型转换为WorkflowStepConfig对象。
除
name外的所有参数均为仅关键字参数。
依赖项通过参数名称隐式解析。如果步骤函数参数名称与先前声明的步骤函数匹配,其结果将注入步骤。
如果定义 ctx 参数,步骤上下文将注入该参数。
from workers import WorkflowEntrypoint
class MyWorkflow(WorkflowEntrypoint):
async def run(self, event, step):
@step.do()
async def my_first_step():
# do some work
return "Hello World!"
await my_first_step()请注意,装饰器不会调用步骤,它只是返回可用于调用步骤的可调用对象。你必须调用该可调用对象才能使步骤运行。
从步骤返回状态时,必须确保返回值可序列化。
-
step.sleep(name, duration)name— 步骤名称。duration— 休眠时长,以秒数或WorkflowDuration兼容字符串表示。
async def run(self, event, step):
await step.sleep("my-sleep-step", "10 seconds")-
step.sleep_until(name, timestamp)name— 步骤名称。timestamp—datetime.datetime对象或自 UNIX 纪元以来的秒数,workflow 实例将休眠至此时间。
import datetime
async def run(self, event, step):
await step.sleep_until("my-sleep-step", datetime.datetime.now() + datetime.timedelta(seconds=10))-
step.wait_for_event(name, event_type, timeout="24 hours")name— 步骤名称。event_type— 要等待的事件类型。timeout—wait_for_event调用的超时。默认超时为 24 小时。
async def run(self, event, step):
await step.wait_for_event("my-wait-for-event-step", "my-event-type")event 参数是包含传递给 workflow 实例的 payload 及其他元数据的字典:
payload- 传递给 workflow 实例的 payload。timestamp- workflow 触发的时间戳。instanceId- 当前 workflow 实例的 ID。workflowName- workflow 的名称。
Workflows 语义允许用户捕获传播到顶层的异常。
在 except 块中捕获特定异常可能不起作用,因为某些 Python 错误在通过 RPC 层传递时不会重新实例化为相同类型的错误。
async def run(self, event, step):
async def try_step(fn):
try:
return await fn()
except Exception as e:
print(f"Successfully caught {type(e).__name__}: {e}")
@step.do("my_failing")
async def my_failing():
print("Executing my_failing")
raise TypeError("Intentional error in my_failing")
await try_step(my_failing)Python Workflows SDK 提供 NonRetryableError 类,用于表示步骤不应重试。
from workers.workflows import NonRetryableError
raise NonRetryableError(message)你可以通过将 WorkflowStepConfig 对象传递给 step.do 装饰器的 config 参数,将步骤绑定到特定重试策略。
使用 Python Workflows 时,需要确保 dict 符合 WorkflowStepConfig 类型。
from workers import WorkflowEntrypoint
class DemoWorkflowClass(WorkflowEntrypoint):
async def run(self, event, step):
@step.do('step-name', config={"retries": {"limit": 1, "delay": "10 seconds"}})
async def first_step():
# do some work
pass如果定义 ctx 参数,步骤上下文 将注入该参数。上下文是具有以下键的字典:
| 键 | 类型 | 描述 |
|---|---|---|
step |
dict |
包含 name(步骤名称)和 count(使用此名称调用 step.do 的次数)。 |
attempt |
int |
当前尝试次数(从 1 开始)。 |
config |
dict |
此步骤已解析的重试和超时配置。 |
from workers import WorkflowEntrypoint
class CtxWorkflow(WorkflowEntrypoint):
async def run(self, event, step):
@step.do()
async def read_context(ctx):
print(ctx["step"]["name"]) # step name
print(ctx["step"]["count"]) # step count
print(ctx["attempt"]) # attempt number
print(ctx["config"]) # resolved step config
return ctx["attempt"]
return await read_context()请注意,env 是通过 JsProxy ↗ 暴露给 Python 脚本的 JavaScript 对象。你可以像在 JavaScript worker 上一样访问绑定。请参阅 Workflow 绑定文档 了解可用方法。
考虑之前名为 MY_WORKFLOW 的绑定。创建新实例的方式如下:
from workers import Response, WorkerEntrypoint
class Default(WorkerEntrypoint):
async def fetch(self, request):
instance = await self.env.MY_WORKFLOW.create()
return Response.json({"status": "success"})