0Pricing
Learn AI with Python · 课时

使用 Airflow 构建可扩展的机器学习 pipeline

基于 DAG 的管道、任务依赖,以及数据摄取 → 训练 → 评估 → 部署。

使用 Airflow 构建可扩展的机器学习 pipeline 是 CoddyKit 上的免费 Learn AI with Python 课时。 这是第 2 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Learn AI with Python 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Learn AI with Python 课程共包含 4 节课。

为什么需要编排

机器学习工作流包含许多步骤:摄取数据、构建特征、训练、评估和部署。手动运行这些步骤很脆弱。Apache Airflow会将它们以代码形式进行编排,并内置调度、重试、依赖关系和监控功能。

DAG

Airflow 将流水线建模为由任务组成的DAG(有向无环图)。无环意味着任何任务都不能通过循环依赖自身,因此执行始终会按照从开始到结束的明确定义顺序进行。

定义 DAG

您可以使用标识符、调度计划和开始日期声明一个 DAG。调度计划控制它运行的频率,例如每天重新训练。

from airflow import DAG
import datetime

with DAG(
    dag_id="ml_pipeline",
    schedule="@daily",
    start_date=datetime.datetime(2024, 1, 1),
    catchup=False,
) as dag:
    ...

操作器就是任务

DAG 中的每个节点都是由操作器创建的任务。不同的操作器会运行不同类型的工作:Python 函数、Bash 命令、SQL 查询等。

使用 PythonOperator 进行训练

PythonOperator会运行一个 Python 可调用对象。它非常适合用于调用训练函数的训练步骤。

from airflow.operators.python import PythonOperator

def train_model():
    # load features, fit model, save artifact
    ...

train = PythonOperator(
    task_id="train",
    python_callable=train_model,
)

使用 BashOperator 进行评估

BashOperator会运行 shell 命令,适合调用评估脚本或 CLI 工具。

from airflow.operators.bash import BashOperator

evaluate = BashOperator(
    task_id="evaluate",
    bash_command="python /opt/ml/evaluate.py --model latest",
)

定义依赖关系

>> 操作符会设置任务顺序:a >> b 表示 a 成功后运行 b。这样便可连接 DAG,确保训练完成后才开始评估。

train >> evaluate
# train must succeed before evaluate runs

更长的依赖链

您可以串联许多任务,以表达完整的流水线顺序。Airflow 会并行运行相互独立的分支,同时遵守您声明的每一项依赖关系。

ingest >> features >> train >> evaluate >> deploy

使用 XCom 传递数据

任务彼此隔离运行,那么一个任务如何将结果传给下一个任务?XCom(跨任务通信)允许一个任务推送一个小型值(例如模型路径或指标),供下游任务提取。

def train_model(ti):
    path = "/models/run_42.pt"
    ti.xcom_push(key="model_path", value=path)

def deploy(ti):
    path = ti.xcom_pull(key="model_path", task_ids="train")
    # deploy the artifact at path

XCom 适用于小数据

XCom 用于传递小型元数据(路径、标识符和指标),而不是大型数据集。大型构件应存放在对象存储(S3)中,只通过 XCom 传递其位置。过度使用 XCom 传递大型负载会给元数据库造成压力。

调度与重试

Airflow 可直接提供生产环境所需的稳健性:调度计划会自动触发运行,失败的任务会通过退避策略进行重试,界面会显示每个任务和每次运行的状态,并在失败时发出警报。

快速检查

测试您对 Airflow 的掌握程度。

回顾

您学习了如何使用 Airflow 构建可扩展的机器学习流水线:

  • 使用包含调度计划的 DAG定义流水线
  • PythonOperator运行训练;BashOperator运行评估
  • >> 操作符设置任务依赖关系
  • XCom在任务之间传递小型构件;大数据则存储在对象存储中

常见问题解答

「使用 Airflow 构建可扩展的机器学习 pipeline」课时是免费的吗?

是的 — 「使用 Airflow 构建可扩展的机器学习 pipeline」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Learn AI with Python 课程的其余内容,请升级到 CoddyKit PRO。 Learn AI with Python 课程共包含 4 节课。

「使用 Airflow 构建可扩展的机器学习 pipeline」这节课中我会学到什么?

基于 DAG 的管道、任务依赖,以及数据摄取 → 训练 → 评估 → 部署。 你通过在浏览器中直接运行的动手代码来练习 Learn AI with Python,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 Learn AI with Python 需要有经验吗?

无需任何先前经验。CoddyKit 上的 Learn AI with Python 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 2 节课,共 4 节。

「使用 Airflow 构建可扩展的机器学习 pipeline」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 Learn AI with Python 课中编写并运行代码吗?

能。每节 Learn AI with Python 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. 人工智能系统架构模式
  2. 使用 Airflow 构建可扩展的机器学习 pipeline
  3. 特征存储:Feast 与 Tecton
  4. 人工智能系统的可观测性与监控
← 返回 Learn AI with Python