使用 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 pathXCom 适用于小数据
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 反馈 — 无需本地设置。
此课程中的所有课时
- 人工智能系统架构模式
- 使用 Airflow 构建可扩展的机器学习 pipeline
- 特征存储:Feast 与 Tecton
- 人工智能系统的可观测性与监控