跳转至

OOP Workflow 示例

examples/oop_workflow.py 展示如何把工作流状态放在普通 Python 对象上,同时继续使用 Astrum 的装饰器任务。

运行方式

python examples/oop_workflow.py

预期输出:

completed
A-001 in APAC: 129.60 USD

写法

先创建 registry,在 class 内部修饰实例方法,然后实例化对象,并通过 AstrumConfig(class_instances=[...]) 提供给调度器。

workflow = SchedulerRegistry("oop_workflow")


class OrderService:
    def __init__(self, region: str) -> None:
        self.region = region
        self.tax_rate = 0.08

    @workflow.task("load_order")
    async def load_order(self) -> dict:
        return {"order_id": "A-001", "region": self.region, "subtotal": 120}

下游实例方法可以像模块级函数一样使用 Ref/F

@workflow.task("price_order")
async def price_order(
    self,
    subtotal: Ref[int, F("load_order", "subtotal")],
) -> dict:
    total = round(subtotal * (1 + self.tax_rate), 2)
    return {"total": total, "currency": "USD"}

调度器会为被装饰的未绑定实例方法自动注入 self

service = OrderService(region="APAC")
report = await workflow.run(
    target_tasks=["format_receipt"],
    config=AstrumConfig(
        class_instances=[service],
        skip_type_check=True,
        silence_warnings=True,
    ),
)

手动 DAG 是另一种情况:如果你直接把 service.method 传给 DynamicScheduler,Python 已经完成了 self 绑定,不需要配置 class_instances