从 Airflow 2 到 Airflow 3:Pine59 如何在 Google Cloud 上提速数据与 MLOps 流水线

2026-09-18 20 预计阅读时间: 1 分钟
来源: cloud.google.com AI 摘要 Original link

Disclaimer: This article is an AI-assisted summary. Read it together with the original source when precision matters. The summary may omit context, version differences, or edge cases and is not official documentation.

预计阅读时间:12 分钟

当数据管道每天要处理数百万个复杂数据点时,编排层不再只是“把任务按顺序运行起来”。它必须能够承受流量峰值、快速调度任务、支持机器学习推理,并且让工程师在数百个 DAG 之间定位问题。Pine59 的迁移实践说明,升级 Airflow 的价值不仅在于获得新版本功能,更在于重新划分编排、数据处理和模型计算的边界。

Pine59 提供位置智能数据,管道以小时、季度等不同频率生成分析指标。其中,Daily Foot Traffic 指标一次任务最多需要计算 1,400 万个地点的数据。整套系统运行在 Google Cloud 上,BigQuery 承担主要数据处理,Managed Service for Apache Airflow(原 Cloud Composer)负责编排,并最终迁移到运行 Apache Airflow 3 的 Managed Airflow Gen 3 环境。

先用生产负载验证升级价值

Pine59 长期使用共享 monorepo 管理多个项目的代码、工具和数百个 DAG。随着数据量和机器学习工作负载持续增长,原有环境逐渐暴露出两个典型问题:

  • 高峰期任务长时间停留在 queued 状态,真正开始执行前已经消耗了大量时间。
  • 编排、数据处理和机器学习推理的资源边界不够清晰,扩容和排障都变得复杂。

因此,Pine59 没有只在测试 DAG 上验证新环境,而是直接用生产工作负载对 Managed Airflow Gen 3 进行压力测试。测试结果显示,新环境在处理速度、任务调度和整体稳定性方面都有明显改善,于是团队决定完成全面迁移。

这种迁移顺序值得借鉴:不要先假设升级一定有效,也不要只比较单个任务的 wall-clock time。更可靠的做法是选取具有代表性的 DAG,连续运行足够多的样本,并同时观察排队时间、运行时间、失败率和重试行为。

可以用下面的方式为 DAG 运行记录建立一个简单的对比查询。字段名称需要根据实际 Airflow 元数据库、日志导出表或监控数据进行调整;这里的 SQL 是一个可改造的分析示例:

-- 假设 dag_run_metrics 表由监控或日志导出任务生成
-- engine_version: airflow-2.11 或 airflow-3.1
SELECT
  engine_version,
  dag_id,
  COUNT(*) AS run_count,
  ROUND(AVG(queued_seconds), 2) AS avg_queued_seconds,
  ROUND(AVG(running_seconds), 2) AS avg_running_seconds,
  ROUND(AVG(total_seconds), 2) AS avg_total_seconds,
  SUM(CASE WHEN state = 'failed' THEN 1 ELSE 0 END) AS failed_runs
FROM `your_project.ops.dag_run_metrics`
WHERE dag_id = 'daily_foot_traffic'
  AND started_at >= TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 30 DAY)
GROUP BY engine_version, dag_id
ORDER BY engine_version;

迁移前后使用同一批次、相近的数据规模和相同的 DAG 版本进行比较,才能避免把业务数据变化或 DAG 内部优化误认为平台升级收益。

把编排层和模型推理解耦

Pine59 的管道并不只是搬运数据,还要驱动复杂的机器学习模型。迁移前,模型推理任务主要使用标准 Kubernetes Operator。迁移到 Managed Airflow Gen 3 后,团队将推理计算放到了专门优化过的 Google Kubernetes Engine(GKE)集群中,并由 Airflow 负责调度和协调。

这形成了更清晰的架构分工:

  • Managed Airflow:负责 DAG 调度、依赖管理、重试、状态追踪和运维入口。
  • BigQuery:负责大规模数据处理和分析查询。
  • 专用 GKE 集群:负责模型推理等高消耗计算。
  • 共享 monorepo:统一管理 DAG、Operator、插件和版本兼容逻辑。

下面是一个可改造的 Airflow DAG 示例。它通过 Kubernetes Pod 执行推理任务,同时把输入输出位置作为参数传入。示例假设集群访问、镜像权限和网络配置已经由平台团队准备好;实际使用时应替换项目、集群、命名空间和镜像名称。

from datetime import datetime

from airflow import DAG
from airflow.providers.cncf.kubernetes.operators.pod import KubernetesPodOperator

with DAG(
    dag_id="daily_foot_traffic_inference",
    start_date=datetime(2024, 1, 1),
    schedule="0 2 * * *",
    catchup=False,
    tags=["bigquery", "ml", "gke"],
) as dag:
    run_inference = KubernetesPodOperator(
        task_id="run_model_inference",
        name="daily-foot-traffic-inference",
        namespace="ml-inference",
        image="REGION-docker.pkg.dev/PROJECT/ml/foot-traffic:VERSION",
        cmds=["python", "-m", "inference_job"],
        arguments=[
            "--input-table",
            "PROJECT.DATASET.features_{{ ds_nodash }}",
            "--output-table",
            "PROJECT.DATASET.predictions_{{ ds_nodash }}",
        ],
        labels={"workload": "foot-traffic", "managed-by": "airflow"},
        get_logs=True,
        is_delete_operator_pod=True,
        startup_timeout_seconds=600,
    )

实际部署时还需要根据模型特点配置 CPU、内存、GPU、节点选择器、污点容忍度和并行度。一个重要原则是:不要让 Airflow Worker 直接承担长时间、高资源消耗的推理计算。编排器应该发起任务并追踪结果,而不是成为模型运行时本身。

用插件和兼容层改善开发体验

当一个 monorepo 中包含数百个互相依赖的 DAG 时,可观测性和排障效率会直接影响交付速度。Pine59 利用 Airflow 3 改进后的插件开发和用户界面能力,快速构建了面向内部开发者的扩展。

其中两个实践尤其具体:

BigQuery 自动链接

Airflow 日志和 XCom 中经常会出现内部 BigQuery 表名。插件可以识别这些引用,并动态生成指向 BigQuery Studio 的链接。工程师不必复制表名、切换浏览器再手动搜索,排查数据问题的路径明显缩短。

DAG Run 配置搜索

DAG Run 的配置通常以 key-value 形式保存,例如日期范围、客户标识或模型版本。通过在 DAG 概览页添加搜索表单,工程师可以直接查询特定配置并定位匹配的运行记录。这比逐个打开 DAG Run 更适合处理批量任务故障。

除了 UI 插件,Pine59 还在 monorepo 中加入了兼容层。这个 compat 模块根据 Airflow 版本抽象差异,让 Operator 的迁移不必一次性修改所有 DAG。可以采用类似下面的项目结构:

repo/
├── dags/
│   ├── daily_foot_traffic.py
│   └── feature_pipeline.py
├── compat/
│   ├── __init__.py
│   └── airflow_api.py
├── plugins/
│   ├── bigquery_links.py
│   └── dag_config_search.py
└── tests/
    └── test_compat.py

兼容模块的核心目标不是永远支持所有旧版本,而是把版本差异集中到少数文件中。迁移完成后,应逐步删除不再需要的分支,避免兼容层演变成新的隐性负担。

性能收益要拆成可解释的指标

Pine59 在迁移后观察到,任务排队延迟显著下降,任务能够更快进入运行状态。对超过 300 次同一 DAG 的运行进行比较时,Managed Airflow Gen 3 与 Airflow 3.1 的排队时间相较旧的 Gen 2 与 Airflow 2.11 有明显改善。

Daily Foot Traffic 是更直观的例子:原先完成一次任务接近 38 分钟,迁移到新环境并结合 DAG 内部优化后,耗时降至不到 26 分钟,处理时间减少接近 32%。

这里需要区分三个指标:

  1. Queued time:任务等待调度或执行资源的时间。
  2. Running time:任务已经启动后实际执行的时间。
  3. End-to-end time:从 DAG Run 开始到全部任务完成的总时间。

平台升级通常最先改善排队时间;SQL、分区、批量大小和模型代码优化则可能改善运行时间。将两者拆开,才能知道收益来自哪里,也才能避免在下一轮扩容时做出错误判断。

迁移前可以照着检查什么

如果团队也在考虑从 Airflow 2 迁移到 Airflow 3 或 Managed Airflow Gen 3,可以按下面的清单推进:

  • 选择一到三个真正受高峰流量影响的生产 DAG,而不是只选最简单的示例。
  • 记录迁移前后的排队时间、运行时间、总耗时、失败率和重试次数。
  • 盘点 Operator、Provider、插件和自定义 Hook 的版本兼容性。
  • 将模型推理、批量计算等重资源任务放到独立的 GKE 或其他计算环境。
  • 在 monorepo 中集中处理 Airflow 版本差异,并为兼容层编写测试。
  • 优先建设能减少排障步骤的 UI 扩展,例如数据表自动链接和配置搜索。
  • 迁移后持续观察资源利用率和成本,不要只看单次任务的速度。

Pine59 的经验并不是“升级版本就能自动获得 32% 的加速”。真正有效的组合是:用生产负载验证新编排环境,重新划分 Airflow 与 ML 计算的职责,同时优化 DAG 和开发工具。对于基础设施维护已经开始挤压数据与 AI 交付时间的团队,这种迁移可以成为一次架构整理的机会,而不只是一次版本切换。


相关推荐