跳转到内容
搜索文档

Python Workers API 参考

最后更新 查看 MarkdownAgent 设置

本指南涵盖 Python Workflows SDK,包括如何使用 Python 构建和创建 workflow 的说明。

WorkflowEntrypoint

WorkflowEntrypoint 是 Python workflow 的主入口点。它扩展 WorkflowEntrypoint 类并实现 run 方法。

from workers import WorkflowEntrypoint

class MyWorkflow(WorkflowEntrypoint):
    async def run(self, event, step):
        # steps here

WorkflowStep

  • 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 — 步骤名称。
    • timestampdatetime.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 — 要等待的事件类型。
    • timeoutwait_for_event 调用的超时。默认超时为 24 小时。
async def run(self, event, step):
    await step.wait_for_event("my-wait-for-event-step", "my-event-type")

event 参数

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)

NonRetryableError

Python Workflows SDK 提供 NonRetryableError 类,用于表示步骤不应重试。

from workers.workflows import NonRetryableError

raise NonRetryableError(message)

配置 workflow 实例

你可以通过将 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)

如果定义 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"})

这篇文档对您有帮助吗?