Training pipelines: Airflow, Kubeflow, Prefect
- ⬜📋 Experiment tracking + registry: MLflow(MLOps — ML em produção)
Recomendamos completar os pré-requisitos antes de seguir, mas nada te impede de continuar.
Training pipeline como cidadão de primeira classe
| Etapa como componente | Entrada | Saída | Por que separar |
|---|---|---|---|
| Extração | Consulta e período | Dado bruto versionado | Reprocessar sem reextrair |
| Preparação | Dado bruto | Conjunto pronto e dividido | A divisão fica registrada, não sorteada de novo |
| Treino | Conjunto e hiperparâmetros | Artefato e métricas | Repetir com outro parâmetro sem refazer o resto |
| Avaliação | Artefato e conjunto reservado | Relatório e veredito | É este passo que BLOQUEIA a promoção |
| Registro | Artefato aprovado | Versão catalogada | Promoção deixa de ser manual |
Notebook é ótimo para exploração. Produção pede pipeline declarativo: cada passo reproduzível, com inputs versionados, outputs nomeados, logs centralizados e retry automático. O objetivo é que qualquer pessoa do time consiga retreinar o modelo com um comando, e que o CI consiga disparar isso sem humano no loop.
Kubeflow Pipelines — componente a componente
from kfp import dsl
from kfp.dsl import Input, Output, Dataset, Model, Metrics
@dsl.component(base_image="python:3.12", packages_to_install=["pandas", "pyarrow"])
def extract_features(dataset_out: Output[Dataset]):
import pandas as pd
df = pd.read_parquet("s3://ffv-ml/raw/events.parquet")
df.to_parquet(dataset_out.path, index=False)
@dsl.component(base_image="ffv/ml-train:2.1")
def train_model(dataset_in: Input[Dataset], model_out: Output[Model], metrics: Output[Metrics]):
import pandas as pd, joblib
from sklearn.ensemble import GradientBoostingClassifier
df = pd.read_parquet(dataset_in.path)
X, y = df.drop(columns=["label"]), df["label"]
m = GradientBoostingClassifier(n_estimators=300).fit(X, y)
joblib.dump(m, model_out.path)
metrics.log_metric("train_rows", len(df))
@dsl.pipeline(name="churn-training")
def pipeline():
feats = extract_features()
train = train_model(dataset_in=feats.outputs["dataset_out"])
train.set_cpu_limit("4").set_memory_limit("16Gi").set_gpu_limit(1)Configuração de pipeline em YAML
# kubeflow pipeline spec (trecho compilado)
apiVersion: argoproj.io/v1alpha1
kind: Workflow
metadata:
generateName: churn-training-
spec:
entrypoint: pipeline
templates:
- name: train-model
retryStrategy:
limit: 3
backoff:
duration: 60s
factor: 2
container:
image: ffv/ml-train:2.1
resources:
limits:
cpu: "4"
memory: 16Gi
nvidia.com/gpu: "1"
env:
- name: MLFLOW_TRACKING_URI
value: http://mlflow.ffv.internal
timeout: 2hSem retryStrategy, cada falha transitória (node evicted, imagem lenta para puxar) vira pipeline vermelho que precisa ser replanejado à mão.
Por que transformar o treino em pipeline declarado, em vez de um script?
Airflow — quando é a escolha certa
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.cncf.kubernetes.operators.pod import KubernetesPodOperator
from datetime import datetime, timedelta
with DAG(
"churn_retraining",
start_date=datetime(2026, 1, 1),
schedule="0 3 * * 1", # toda segunda 03:00
catchup=False,
default_args={"retries": 2, "retry_delay": timedelta(minutes=10)},
):
extract = PythonOperator(task_id="extract", python_callable=extract_features)
train = KubernetesPodOperator(
task_id="train",
image="ffv/ml-train:2.1",
cmds=["python", "-m", "train"],
container_resources={"limits": {"nvidia.com/gpu": "1"}},
)
extract >> trainPrefect 2.x — DX em Python puro
from prefect import flow, task
@task(retries=3, retry_delay_seconds=30, cache_expiration=3600)
def build_features() -> str:
return "s3://ffv-ml/features/churn.parquet"
@task
def train(path: str) -> str:
return "s3://ffv-ml/models/churn-v7.joblib"
@flow(name="churn-retraining")
def retraining():
feats = build_features()
return train(feats)Prefect brilha quando os steps são dinâmicos (mapping sobre lista de shards, por exemplo) e o time quer escrever Python idiomático sem ceder ao YAML.
Perguntas frequentes
❓ Preciso de orquestrador dedicado para treinar?
❓ Orquestrador nativo de contêiner ou genérico?
❓ Como tornar o treino reproduzível?
Fixando
Quando um orquestrador de fluxo genérico é preferível a uma ferramenta específica de aprendizado de máquina?
Qual característica de um passo de pipeline facilita mais a depuração?
Terminou de ler?
Marcar como concluído registra o XP, mantém sua sequência e coloca 3 cartas deste módulo na fila de revisão espaçada.
Próximos passos sugeridos
Temas deste módulo
Discussão
Carregando comentários…