任务调度系统选型:Airflow vs Temporal vs Prefect的深度技术对比与选型决策框架

一、任务调度系统的选型困境:为什么不是简单的"哪个好"

任务调度是数据工程和微服务架构中的基础设施组件。三个主流开源方案——Apache Airflow、Temporal、Prefect——各自代表了不同的设计哲学。Airflow起源于Airbnb的DAG(有向无环图)批处理调度需求,核心是"时间驱动";Temporal起源于Uber的微服务编排需求,核心是"工作流即代码";Prefect最初是Airflow的现代化替代品,核心是"动态工作流与易用性"。

选型的困境在于三个系统的能力高度重叠——都能定义DAG、调度任务、处理重试和告警。但它们在架构假设、执行模型、扩展性上的差异,决定了适用场景的本质区别。Airflow的DAG必须在调度前完全确定(静态DAG),Temporal的Workflow可以动态创建子Workflow(动态DAG),Prefect支持运行时改变DAG结构(参数化DAG)。本文从架构设计、执行模型、部署运维、生产级代码四个维度,提供完整的选型决策框架和迁移方案。

二、三者的架构模型对比

三者的核心差异:Airflow的调度器和执行器分离——Scheduler只负责DAG解析和调度决策,Executor负责Task的物理执行。Temporal采用"确定性重放"架构——Workflow代码在Worker端重放执行,所有决策(随机数、时间等)都从Event History中恢复以保证确定性。Prefect采用Agent架构——由Agent主动轮询Prefect Server获取待执行的Task Run,执行完成后上报结果。

三、生产级代码:同一业务逻辑在三个系统中的实现对比

# ============================================
# 业务场景:电商订单处理流水线
# 接收订单 -> 验证库存 -> 支付处理 -> 物流下单 -> 发送通知
# ============================================

from dataclasses import dataclass
from datetime import datetime, timedelta
from typing import Optional
from enum import Enum
import random

class OrderStatus(Enum):
    PENDING = "pending"
    CONFIRMED = "confirmed"
    PAID = "paid"
    SHIPPED = "shipped"
    COMPLETED = "completed"
    CANCELLED = "cancelled"

@dataclass
class Order:
    order_id: str
    user_id: str
    items: list[dict]
    total_amount: float
    status: OrderStatus = OrderStatus.PENDING
    payment_id: Optional[str] = None
    tracking_number: Optional[str] = None
    created_at: datetime = None

# ==================== Airflow实现 ====================

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.dummy import DummyOperator
from airflow.sensors.external_task_sensor import (
    ExternalTaskSensor
)
from airflow.utils.dates import days_ago
from airflow.utils.trigger_rule import TriggerRule

def validate_inventory_airflow(**context):
    """验证库存(Airflow PythonOperator)"""
    order = context['dag_run'].conf.get('order', {})
    order_id = order.get('order_id', 'N/A')

    # 模拟库存检查
    if order_id == "FAIL":
        raise ValueError(f"库存不足: {order_id}")

    print(f"Airflow: 库存验证通过 {order_id}")
    return {"inventory_ok": True, "order_id": order_id}

def process_payment_airflow(**context):
    """处理支付"""
    ti = context['ti']
    result = ti.xcom_pull(
        task_ids='validate_inventory'
    )
    order_id = result['order_id']

    # 模拟支付处理
    payment_result = {
        "payment_id": f"PAY_{order_id}",
        "status": "success",
    }
    print(f"Airflow: 支付处理完成 {payment_result}")
    return payment_result

def ship_order_airflow(**context):
    """物流下单"""
    ti = context['ti']
    payment_result = ti.xcom_pull(
        task_ids='process_payment'
    )
    order_id = payment_result['payment_id'].replace(
        'PAY_', ''
    )

    tracking = f"SF{random.randint(100000, 999999)}"
    print(f"Airflow: 物流下单完成 运单号={tracking}")
    return {"tracking_number": tracking}

def send_notification_airflow(**context):
    """发送通知"""
    print("Airflow: 通知已发送")
    return {"notified": True}

def handle_failure_airflow(**context):
    """失败处理"""
    print(f"Airflow: 订单处理失败,执行补偿逻辑")
    return {"compensated": True}

# Airflow DAG定义
dag_airflow = DAG(
    dag_id='order_processing_airflow',
    start_date=days_ago(1),
    schedule_interval='@hourly',
    catchup=False,
    max_active_runs=1,
    default_args={
        'owner': 'data-team',
        'retries': 2,
        'retry_delay': timedelta(minutes=5),
    },
)

with dag_airflow:
    start = DummyOperator(task_id='start')
    end = DummyOperator(
        task_id='end',
        trigger_rule=TriggerRule.ALL_DONE,
    )

    validate = PythonOperator(
        task_id='validate_inventory',
        python_callable=validate_inventory_airflow,
    )

    payment = PythonOperator(
        task_id='process_payment',
        python_callable=process_payment_airflow,
    )

    shipping = PythonOperator(
        task_id='ship_order',
        python_callable=ship_order_airflow,
    )

    notify = PythonOperator(
        task_id='send_notification',
        python_callable=send_notification_airflow,
        trigger_rule=TriggerRule.ALL_SUCCESS,
    )

    fail_handler = PythonOperator(
        task_id='handle_failure',
        python_callable=handle_failure_airflow,
        trigger_rule=TriggerRule.ONE_FAILED,
    )

    # 定义DAG依赖
    start >> validate >> payment >> shipping
    shipping >> notify >> end
    validate >> fail_handler >> end
    payment >> fail_handler

# ==================== Temporal实现 ====================

# Temporal的核心概念:
# Workflow = 确定性业务逻辑(只能调用Activity和做纯逻辑)
# Activity = 非确定性副作用(IO、RPC、随机数等)

# 需要先安装 temporalio
# pip install temporalio

from temporalio import activity, workflow
from temporalio.common import RetryPolicy

# --- Activities定义(非确定性操作)---

@activity.defn(name="validate_inventory_activity")
async def validate_inventory_activity(
    order_id: str
) -> dict:
    """库存验证Activity"""
    print(f"Temporal: 库存验证 {order_id}")

    if "FAIL" in order_id.upper():
        raise activity.ApplicationError(
            f"库存不足: {order_id}",
            details={"order_id": order_id},
            non_retryable=True,
        )

    return {"inventory_ok": True, "order_id": order_id}

@activity.defn(name="process_payment_activity")
async def process_payment_activity(
    order_id: str, amount: float
) -> dict:
    """支付处理Activity"""
    print(f"Temporal: 支付处理 order={order_id}")

    # 生产环境:调用支付网关API
    payment_result = {
        "payment_id": f"PAY_{order_id}",
        "status": "success",
        "amount": amount,
    }
    return payment_result

@activity.defn(name="ship_order_activity")
async def ship_order_activity(
    order_id: str
) -> dict:
    """物流下单Activity"""
    tracking = f"SF{random.randint(100000, 999999)}"
    return {"tracking_number": tracking}

@activity.defn(name="send_notification_activity")
async def send_notification_activity(
    user_id: str, tracking: str
) -> dict:
    """发送通知Activity"""
    print(
        f"Temporal: 发送通知 user={user_id} "
        f"tracking={tracking}"
    )
    return {"notified": True}

# --- Workflow定义(确定性编排)---

@workflow.defn(name="OrderProcessingWorkflow")
class OrderProcessingWorkflow:
    """订单处理Workflow"""

    @workflow.run
    async def run(self, order: dict) -> dict:
        workflow.logger.info(
            f"开始处理订单 {order.get('order_id')}"
        )

        order_id = order["order_id"]
        user_id = order["user_id"]
        amount = order["total_amount"]

        retry_policy = RetryPolicy(
            initial_interval=timedelta(seconds=1),
            maximum_interval=timedelta(minutes=5),
            maximum_attempts=3,
            non_retryable_error_types=[
                "库存不足"
            ],
        )

        try:
            # Step 1: 验证库存
            inventory_result = await (
                workflow.execute_activity(
                    validate_inventory_activity,
                    args=[order_id],
                    start_to_close_timeout=timedelta(
                        seconds=10
                    ),
                    retry_policy=retry_policy,
                )
            )

            # Step 2: 处理支付
            payment_result = await (
                workflow.execute_activity(
                    process_payment_activity,
                    args=[order_id, amount],
                    start_to_close_timeout=timedelta(
                        seconds=30
                    ),
                    retry_policy=retry_policy,
                )
            )

            # Step 3: 物流下单
            shipping_result = await (
                workflow.execute_activity(
                    ship_order_activity,
                    args=[order_id],
                    start_to_close_timeout=timedelta(
                        seconds=15
                    ),
                )
            )

            # Step 4: 发送通知
            notify_result = await (
                workflow.execute_activity(
                    send_notification_activity,
                    args=[
                        user_id,
                        shipping_result[
                            "tracking_number"
                        ],
                    ],
                    start_to_close_timeout=timedelta(
                        seconds=10
                    ),
                )
            )

            return {
                "order_id": order_id,
                "payment_id": payment_result[
                    "payment_id"
                ],
                "tracking": shipping_result[
                    "tracking_number"
                ],
                "status": "completed",
            }

        except activity.ActivityError as e:
            # 补偿逻辑:退款等
            workflow.logger.error(
                f"订单处理失败: {order_id}, 原因: {e}"
            )
            raise workflow.ApplicationError(
                f"订单 {order_id} 处理失败: {e}"
            )

# ==================== Prefect实现 ====================

from prefect import flow, task
from prefect.blocks.system import Secret
from prefect.task_runners import (
    ConcurrentTaskRunner
)
from prefect.cache_policies import NONE

@task(
    name="validate-inventory",
    retries=2,
    retry_delay_seconds=60,
)
def validate_inventory_prefect(order_id: str) -> dict:
    """库存验证Task"""
    print(f"Prefect: 库存验证 {order_id}")

    if "FAIL" in order_id.upper():
        raise ValueError(f"库存不足: {order_id}")

    return {"inventory_ok": True, "order_id": order_id}

@task(name="process-payment", retries=1)
def process_payment_prefect(
    order_id: str, amount: float
) -> dict:
    """支付处理Task"""
    print(f"Prefect: 支付处理 order={order_id}")

    payment_result = {
        "payment_id": f"PAY_{order_id}",
        "status": "success",
        "amount": amount,
    }
    return payment_result

@task(name="ship-order")
def ship_order_prefect(order_id: str) -> dict:
    """物流下单Task"""
    tracking = f"SF{random.randint(100000, 999999)}"
    return {"tracking_number": tracking}

@task(name="send-notification")
def send_notification_prefect(
    user_id: str, tracking: str
) -> dict:
    """发送通知Task"""
    print(
        f"Prefect: 发送通知 user={user_id} "
        f"tracking={tracking}"
    )
    return {"notified": True}

@flow(
    name="order-processing-flow",
    task_runner=ConcurrentTaskRunner(),
    log_prints=True,
)
def order_processing_flow_prefect(
    order: dict
) -> dict:
    """订单处理Flow"""
    order_id = order["order_id"]
    user_id = order["user_id"]
    amount = order["total_amount"]

    print(f"Prefect: 开始处理订单 {order_id}")

    # Step 1: 验证库存
    inventory_result = validate_inventory_prefect(
        order_id
    )

    # Step 2: 处理支付
    payment_result = process_payment_prefect(
        order_id, amount
    )

    # Step 3: 物流下单
    shipping_result = ship_order_prefect(order_id)

    # Step 4: 发送通知
    notify_result = send_notification_prefect(
        user_id,
        shipping_result["tracking_number"],
    )

    return {
        "order_id": order_id,
        "payment_id": payment_result["payment_id"],
        "tracking": shipping_result[
            "tracking_number"
        ],
        "status": "completed",
    }

# 如果某个Task失败,Prefect自动重试
# 如果需要补偿,可以定义子Flow
@flow(name="order-compensation-flow")
def order_compensation_flow(order_id: str):
    """订单失败补偿Flow"""
    print(f"Prefect: 执行补偿逻辑 order={order_id}")
    # 退款等操作
    return {"rollback": True}

# ==================== 选型决策引擎 ====================

class SchedulerDecisionEngine:
    """任务调度系统选型决策引擎"""

    # 维度权重配置
    DIMENSION_WEIGHTS = {
        "dynamic_dag": 0.20,        # 动态DAG能力
        "operational_simplicity": 0.15,  # 运维简单性
        "scalability": 0.15,        # 扩展性
        "reliability": 0.15,        # 可靠性
        "monitoring": 0.10,         # 监控
        "ecosystem": 0.15,          # 生态系统
        "cost": 0.10,               # 成本
    }

    # 各系统的评分矩阵 (0-10分)
    SCORE_MATRIX = {
        "airflow": {
            "dynamic_dag": 3,   # 静态DAG
            "operational_simplicity": 5,
            "scalability": 6,
            "reliability": 7,
            "monitoring": 8,
            "ecosystem": 10,
            "cost": 9,
        },
        "temporal": {
            "dynamic_dag": 10,  # 原生动态Workflow
            "operational_simplicity": 6,
            "scalability": 9,
            "reliability": 10,
            "monitoring": 7,
            "ecosystem": 6,
            "cost": 6,
        },
        "prefect": {
            "dynamic_dag": 8,
            "operational_simplicity": 9,
            "scalability": 7,
            "reliability": 6,
            "monitoring": 8,
            "ecosystem": 5,
            "cost": 8,
        },
    }

    def __init__(self, requirements: dict = None):
        self.requirements = requirements or {}

    def calculate_scores(self) -> dict[str, float]:
        """计算各系统的综合得分"""
        results = {}

        for system in ["airflow", "temporal", "prefect"]:
            total = 0.0
            detail = {}

            for dim, weight in (
                self.DIMENSION_WEIGHTS.items()
            ):
                score = self.SCORE_MATRIX[system][dim]
                weighted = score * weight
                total += weighted
                detail[dim] = {
                    "raw": score, "weighted": round(
                        weighted, 2
                    )
                }

            results[system] = {
                "total_score": round(total, 2),
                "details": detail,
            }

        return results

    def recommend(self) -> dict:
        """根据需求特征给出推荐"""
        scores = self.calculate_scores()

        # 按场景特征调整权重
        scenario = self.requirements.get(
            "scenario", "batch_etl"
        )

        if scenario == "microservice_orchestration":
            # 微服务编排场景:Temporal优先
            best = "temporal"
            reason = (
                "微服务编排需要动态Workflow和长事务支持,"
                "Temporal的Saga模式天然适合"
            )
        elif scenario == "data_pipeline":
            # 数据管道场景:Airflow优先
            best = "airflow"
            reason = (
                "数据管道需要丰富的Connector生态和"
                "静态DAG的可预测性,Airflow生态最成熟"
            )
        elif scenario == "ml_pipeline":
            # ML管道场景:Prefect优先
            best = "prefect"
            reason = (
                "ML管道需要动态参数化和Pythonic接口,"
                "Prefect的@task/@flow装饰器模式最简洁"
            )
        else:
            # 默认按最高分推荐
            best = max(
                scores, key=lambda k: scores[k]["total_score"]
            )
            reason = "综合评分最高"

        return {
            "recommended": best,
            "reason": reason,
            "scores": scores,
        }

# 使用示例
if __name__ == "__main__":
    # 选型决策
    engine = SchedulerDecisionEngine({
        "scenario": "data_pipeline",
        "team_size": 5,
        "use_dynamic_dag": True,
    })

    result = engine.recommend()
    print("=== 任务调度系统选型推荐 ===")
    print(f"推荐: {result['recommended']}")
    print(f"原因: {result['reason']}")
    print("\n评分详情:")
    for system, score_data in (
        result['scores'].items()
    ):
        print(
            f"\n{system}: "
            f"总分={score_data['total_score']}"
        )
        for dim, detail in (
            score_data['details'].items()
        ):
            print(
                f"  {dim}: "
                f"{detail['raw']}/10 "
                f"(加权={detail['weighted']})"
            )

    # 执行Airflow版本(通过PythonOperator)
    print("\n=== Airflow DAG结构 ===")
    print(
        "start -> validate_inventory -> process_payment"
        " -> ship_order -> send_notification -> end"
    )
    print(
        "validate_inventory -> handle_failure -> end # 失败路径"
    )
    print(
        "process_payment -> handle_failure -> end     # 失败路径"
    )

四、工程落地中的关键决策:从Airflow迁移到Temporal的陷阱

从Airflow迁移到Temporal的最大挑战不是代码改写,而是心智模型的转变。Airflow的Task是"按DAG顺序执行的无状态函数"——Task之间通过XCom传递少量数据,不保持任何内部状态。Temporal的Workflow是"有状态的长期运行对象"——Workflow可以持续数天甚至数月,内部状态由Event History持久化。

迁移过程中的三个关键陷阱:一是确定性约束——Airflow的PythonOperator可以调用任何外部API,但Temporal的Workflow必须是确定性的(所有外部调用必须封装为Activity)。如果Airflow代码中有random.random()datetime.now()、HTTP请求等非确定性操作,必须重构为Activity。二是XCom大对象——Airflow中通过XCom传递几MB的数据是常态,但Temporal的Workflow输入/输出限制为2MB(更大的数据需通过Activity直接写入外部存储,Workflow只传递引用)。三是补偿逻辑——Airflow通过trigger_rule=ONE_FAILED定义失败补偿路径,Temporal通过workflow.continue_as_new或Saga模式实现补偿(每个正向操作对应一个补偿操作)。

迁移的推荐路径是先迁移最简单的DAG(3-5个Task),验证确定性约束和补偿逻辑的正确性后,再逐步迁移复杂DAG。迁移过程中的双跑策略:Airflow和Temporal并行运行2周,通过diff对比两个系统的输出一致性,确认无误后正式切换。

五、总结

任务调度系统选型的关键是匹配架构假设与业务场景。Airflow适合数据管道(静态DAG+丰富Connector生态+成熟社区),Temporal适合微服务编排(动态Workflow+确定性重放+Saga分布式事务),Prefect适合现代数据栈(Pythonic API+动态参数化+云原生部署)。维度加权评分为:动态DAG 20%、运维简单性 15%、扩展性 15%、可靠性 15%、监控 10%、生态 15%、成本 10%。三个系统的执行模型本质区别:Airflow是"Scheduler+Executor"分离,Temporal是"Workflow确定性重放+Activity副作用隔离",Prefect是"Agent主动轮询"。迁移过程中的关键约束是Temporal的Workflow确定性要求(禁止rand/time/HTTP等非确定性操作)和2MB输入输出限制。迁移的推荐策略是先迁移简单DAG并行双跑2周,验证一致性后逐步迁移复杂DAG。对于中小团队(<10人),Prefect的易用性和Pythonic接口是最大优势;对于需要长事务(数天级别)和补偿逻辑的场景,Temporal的Saga模式是必选项;对于已有成熟Airflow基础设施的团队,迁移成本是首要考量。

更多推荐