AI流水线实践笔记

折腾 AI 项目有一段时间了,从一开始的单个 Python 脚本到现在的一整套流水线,。

背景

折腾 AI 项目有一段时间了,从一开始的单个 Python 脚本到现在的一整套流水线,中间踩了不少坑。最近团队要上一个大项目,涉及数据预处理、模型训练、评估、部署多个环节,靠之前那种"一个脚本跑到底"的方式显然不够用了。

写这篇也是为了总结一下这段经历,给同样在折腾 AI 流水线的同学一些参考。不会讲太理论化的东西,更多是从实际使用者的角度说说怎么把散落的脚本串起来,变成一个自动化流水线。

需求

先说说我们要解决的几个实际问题:

手动操作太多

  • 数据预处理要手动触发
  • 训练完要手动评估模型
  • 部署要手动复制文件到服务器
  • 出问题了要手动回滚

流程不透明

  • 不知道哪个版本用了什么数据
  • 训练参数散落在各个脚本里
  • 实验结果没有统一管理

复现困难

  • 同样的代码跑两次结果不一样
  • 不知道环境差异在哪里
  • 团队成员之间不好协作

简单说就是需要一个从数据到部署的完整自动化流程,既能减少手动操作,又能保证实验可追踪、可复现。

实现

初期方案:Shell 脚本串联

最开始的想法很简单,写个 Shell 脚本把各个步骤串起来就行了:

#!/bin/bash
# run_pipeline.sh

echo "开始数据预处理"
python preprocess_data.py --input data/raw --output data/processed

echo "开始模型训练"
python train_model.py --data data/processed --output models/checkpoint.pt

echo "开始模型评估"
python evaluate_model.py --model models/checkpoint.pt --output results/

echo "开始模型部署"
python deploy_model.py --model models/checkpoint.pt --target server

这种方式刚开始还行,但很快就暴露了问题:

  • 中间某一步失败了,后续步骤还是会执行
  • 每次都要手动修改参数
  • 没有日志和状态记录
  • 无法并行执行独立任务

进阶方案:Airflow 工作流

后来改用了 Airflow,一个专门用于工作流编排的工具:

# pipelines/training_pipeline.py
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime

default_args = {
    'owner': 'ai-team',
    'start_date': datetime(2024, 1, 1),
    'retries': 1,
}

dag = DAG(
    'ai_training_pipeline',
    default_args=default_args,
    schedule_interval='@daily',
    catchup=False,
)

preprocess_task = PythonOperator(
    task_id='preprocess_data',
    python_callable=preprocess_data,
    op_kwargs={'input_path': 'data/raw', 'output_path': 'data/processed'},
    dag=dag,
)

train_task = PythonOperator(
    task_id='train_model',
    python_callable=train_model,
    op_kwargs={'data_path': 'data/processed', 'model_path': 'models/checkpoint.pt'},
    dag=dag,
)

evaluate_task = PythonOperator(
    task_id='evaluate_model',
    python_callable=evaluate_model,
    op_kwargs={'model_path': 'models/checkpoint.pt', 'output_path': 'results/'},
    dag=dag,
)

# 定义任务依赖
preprocess_task >> train_task >> evaluate_task

这个方案解决了之前的问题:

  • 任务失败后会自动停止后续任务
  • 支持参数配置和模板
  • 完整的日志和监控
  • 支持任务并行和复杂依赖

当前方案:MLflow + Prefect

为了更好地管理实验追踪,现在用的是 MLflow 做实验管理,Prefect 做工作流编排:

# pipelines/pipeline.py
from prefect import flow, task
from prefect.deployments import Deployment
from prefect.orion.schemas.schedules import IntervalSchedule
import mlflow

@task
def preprocess_data(input_path, output_path):
    mlflow.log_param("input_path", input_path)
    mlflow.log_param("output_path", output_path)

    # 预处理逻辑
    processed_data = do_preprocess(input_path)

    mlflow.log_metric("samples_count", len(processed_data))
    return processed_data

@task
def train_model(processed_data, model_config):
    mlflow.log_params(model_config)

    # 训练逻辑
    model = do_train(processed_data, model_config)

    # 记录模型
    mlflow.pytorch.log_model(model, "model")
    mlflow.log_metric("train_loss", model.train_loss)

    return model

@task
def evaluate_model(model, test_data):
    # 评估逻辑
    metrics = do_evaluate(model, test_data)

    for key, value in metrics.items():
        mlflow.log_metric(f"eval_{key}", value)

    return metrics

@flow(name="ai-training-pipeline")
def ai_training_pipeline(config):
    with mlflow.start_run():
        processed_data = preprocess_data(config["input_path"], config["output_path"])
        model = train_model(processed_data, config["model_config"])
        metrics = evaluate_model(model, config["test_data"])

    return metrics

# 部署流水线
deployment = Deployment.build_from_flow(
    flow=ai_training_pipeline,
    name="ai-training-deployment",
    schedule=IntervalSchedule(interval=timedelta(days=1)),
    tags=["ml", "production"]
)
deployment.apply()

流水线架构

整个流水线的架构大概是这样的:

graph TB subgraph 数据层 A[原始数据] --> B[数据预处理] B --> C[特征工程] C --> D[数据集划分] end subgraph 训练层 D --> E[模型训练] E --> F[超参调优] F --> G[模型评估] end subgraph 部署层 G --> H[模型打包] H --> I[部署到服务器] I --> J[监控告警] end subgraph 管理层 K[MLflow 实验追踪] L[Prefect 工作流编排] M[日志和监控] end B -.实验追踪.-> K E -.实验追踪.-> K G -.实验追踪.-> K B -.任务调度.-> L E -.任务调度.-> L I -.任务调度.-> L E -.日志记录.-> M J -.日志记录.-> M

踩坑

环境管理问题

一开始没有用容器化,结果在不同机器上跑同一个脚本,环境差异导致各种诡异问题。比如这个机器装的 PyTorch 1.12,那个机器装的是 2.0,API 不兼容。

现在统一用 Docker:

# Dockerfile
FROM pytorch/pytorch:2.0.0-cuda11.7-cudnn8-runtime

WORKDIR /app

COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

COPY . .

CMD ["python", "main.py"]

配合 docker-compose:

# docker-compose.yml
version: '3.8'

services:
  mlflow:
    image: mlflow/mlflow:v2.0
    ports:
      - "5000:5000"
    volumes:
      - ./mlruns:/mlflow/mlruns

  training:
    build: .
    environment:
      - MLFLOW_TRACKING_URI=http://mlflow:5000
    volumes:
      - ./data:/app/data
      - ./models:/app/models
    depends_on:
      - mlflow

数据依赖问题

有一次训练跑了一晚上,结果第二天发现数据源在训练中途更新了,导致训练用的数据不一致。

现在的做法是数据快照:

# 保存数据快照
import hashlib
import shutil

def get_data_hash(data_path):
    """计算数据文件的哈希值"""
    hasher = hashlib.md5()
    with open(data_path, 'rb') as f:
        while chunk := f.read(8192):
            hasher.update(chunk)
    return hasher.hexdigest()

def save_data_snapshot(data_path, snapshot_dir):
    """保存数据快照"""
    data_hash = get_data_hash(data_path)
    snapshot_path = os.path.join(snapshot_dir, f"{data_hash}.pkl")

    if not os.path.exists(snapshot_path):
        shutil.copy(data_path, snapshot_path)

    return snapshot_path

资源竞争问题

同时跑多个实验时,GPU 显存不够导致训练失败。

用 Prefect 的资源限制:

from prefect import flow, task
from prefect.orion.schemas.schedules import CronSchedule

@task(
    tags=["gpu"],
    cache_key_fn=lambda **kwargs: f"train-{kwargs['model_config']}",
    retries=2,
    retry_delay_seconds=60
)
def train_model(data, model_config):
    # 检查 GPU 可用性
    import torch
    if not torch.cuda.is_available():
        raise RuntimeError("GPU not available")

    # 训练逻辑
    model = do_train(data, model_config)
    return model

配合 Kubernetes 做资源限制:

# k8s-pod.yaml
apiVersion: v1
kind: Pod
metadata:
  name: training-pod
spec:
  containers:
  - name: trainer
    image: ai-training:latest
    resources:
      limits:
        nvidia.com/gpu: 1
        memory: "16Gi"
        cpu: "8"
      requests:
        nvidia.com/gpu: 1
        memory: "8Gi"
        cpu: "4"

实验追踪问题

刚开始没做好实验追踪,想对比不同超参的效果时,发现参数和结果都没记录。

现在用 MLflow 做完整追踪:

import mlflow
from mlflow.tracking import MlflowClient

def log_experiment(params, metrics, artifacts, tags=None):
    """记录实验信息"""
    with mlflow.start_run():
        # 记录参数
        mlflow.log_params(params)

        # 记录指标
        mlflow.log_metrics(metrics)

        # 记录模型和文件
        for artifact_name, artifact_path in artifacts.items():
            mlflow.log_artifact(artifact_path, artifact_name)

        # 添加标签
        if tags:
            mlflow.set_tags(tags)

        # 返回 run ID 用于后续查询
        return mlflow.active_run().info.run_id

def compare_experiments(experiment_name, metric_name):
    """对比实验结果"""
    client = MlflowClient()
    exp = client.get_experiment_by_name(experiment_name)

    runs = client.search_runs(
        experiment_ids=[exp.experiment_id],
        order_by=[f"metrics.{metric_name} DESC"]
    )

    results = []
    for run in runs:
        results.append({
            'run_id': run.info.run_id,
            'params': run.data.params,
            'metrics': run.data.metrics,
            'start_time': run.info.start_time
        })

    return results

结果

折腾了这么一套下来,效果还是很明显的:

效率提升

  • 完整流水线从手动 2 小时缩减到自动 15 分钟
  • 实验迭代速度提升了 3 倍
  • 团队成员可以专注于算法优化,不用管基础设施

质量提升

  • 实验结果可追踪、可复现
  • 环境差异导致的问题减少了 80%
  • 回滚和调试变得很简单

协作提升

  • 实验配置集中管理
  • 团队成员可以快速理解彼此的工作
  • 新人上手时间从一周缩短到两天

当然,这套东西也不是完美无缺的,维护成本还是有的。但对于需要长期迭代、多人协作的项目来说,还是值得投入的。

结语

AI 流水线这东西,说复杂也复杂,说简单也简单。核心就是把散落的脚本有序地组织起来,配合合适的工具做好追踪和管理。

不一定非要用我说的这套,关键还是要根据自己的项目规模和团队情况来选择。小项目用 Shell 脚本就够了,大项目再上 Airflow、Prefect 这些工具。

最重要的还是要有这个意识:一开始就要想着"以后怎么办",而不是先把功能堆上去再说。不然等到项目大了再重构,那个成本就太高了。

希望这篇能给正在折腾 AI 流水线的同学一些参考。有啥问题或者更好的方案,也欢迎交流讨论。

版权声明: 本文首发于 指尖魔法屋-AI流水线实践笔记https://blog.thinkmoon.cn/post/285-ai-pipeline-script-automated-practice/) 转载或引用必须申明原指尖魔法屋来源及源地址!