当每个机器学习团队都维护一套 Spark 训练脚本时,训练平台很快会出现参数格式不一、依赖关系隐蔽、失败处理重复、运行过程难以审计等问题。Yelp 推出的内部框架 Training Orchestrator,试图把这些分散脚本收敛为配置驱动、基于 DAG 的统一执行模型。
这项变化的重点并不是简单地“少写几个 Spark 脚本”,而是把训练任务从团队私有代码,转化为平台能够解析、调度和治理的工作流定义。
从独立脚本转向声明式训练任务
传统训练脚本往往同时承担多种职责:读取参数、准备数据、调用 Spark、训练模型、验证指标、上传产物和发送通知。随着业务增长,不同团队会复制这些逻辑,再分别修改。脚本数量增加后,平台很难统一升级运行时、补充监控或修改重试策略。
Training Orchestrator 采用配置驱动模式后,团队描述的是“训练什么”,编排框架负责“如何执行”。一份训练配置通常可以表达以下信息:
- 任务使用的程序、镜像或入口命令
- 输入数据与模型产物的位置
- CPU、内存和 Spark executor 等资源参数
- 任务之间的依赖关系
- 超时、重试及失败处理策略
来源摘要没有披露 Yelp 内部配置格式或具体 API。下面的 YAML 是一个可以这样实践的简化示例,并非 Training Orchestrator 的真实接口:
name: restaurant-ranking-training
run_date: "2025-01-15"
jobs:
prepare_features:
command: "python jobs/prepare_features.py"
args:
input: "data/events.jsonl"
output: "artifacts/features.json"
train_model:
needs: [prepare_features]
command: "python jobs/train_model.py"
args:
input: "artifacts/features.json"
output: "artifacts/model.json"
evaluate_model:
needs: [train_model]
command: "python jobs/evaluate_model.py"
args:
model: "artifacts/model.json"
threshold: "0.80"
配置把入口、参数和依赖放在稳定的数据结构中。平台可以在真正启动 Spark 作业之前完成 schema 校验、依赖检查和权限检查,也可以为所有团队统一注入日志、指标及运行元数据。
DAG 让训练依赖变得可计算
机器学习训练通常不是单一作业,而是一条包含特征生成、训练、评估和发布的流水线。DAG,也就是有向无环图,把每个阶段表示为节点,把依赖关系表示为有向边。
这种模型带来几个直接收益:
- 没有依赖关系的节点可以并行执行。
- 上游失败时,下游节点不会使用不完整产物继续运行。
- 编排器可以检测循环依赖,避免工作流永远无法启动。
- 单个节点可以独立重试,不必重新执行整条训练流水线。
- 运行记录可以精确到节点,便于定位数据准备、训练或评估阶段的问题。
不过,DAG 只描述执行顺序,并不会自动保证数据正确。工程团队仍然需要定义产物版本、数据快照、幂等性和模型发布规则。例如,重试 train_model 时,它应读取同一份固定特征集,而不是悄悄读取已经更新的在线表。
可以这样实践:实现一个最小 DAG 执行器
下面的示例使用 Python 和 PyYAML 读取前面的配置,检查依赖并按拓扑顺序执行命令。它适合用于理解配置驱动编排的核心机制,不包含生产系统需要的分布式锁、持久化状态和远程任务调度。
将 YAML 保存为 training.yaml,然后安装依赖:
python -m venv .venv
. .venv/bin/activate
python -m pip install PyYAML
创建 orchestrator.py:
from __future__ import annotations
import argparse
import subprocess
from collections import deque
from pathlib import Path
import yaml
def execution_order(jobs: dict[str, dict]) -> list[str]:
unknown = {
dep
for job in jobs.values()
for dep in job.get("needs", [])
if dep not in jobs
}
if unknown:
raise ValueError(f"Unknown dependencies: {sorted(unknown)}")
indegree = {name: 0 for name in jobs}
children = {name: [] for name in jobs}
for name, job in jobs.items():
for dependency in job.get("needs", []):
indegree[name] += 1
children[dependency].append(name)
ready = deque(name for name, degree in indegree.items() if degree == 0)
order = []
while ready:
current = ready.popleft()
order.append(current)
for child in children[current]:
indegree[child] -= 1
if indegree[child] == 0:
ready.append(child)
if len(order) != len(jobs):
raise ValueError("The training DAG contains a cycle")
return order
def run_job(name: str, job: dict) -> None:
command = [*job["command"].split()]
for key, value in job.get("args", {}).items():
command.extend([f"--{key.replace('_', '-')}", str(value)])
print(f"[start] {name}: {' '.join(command)}", flush=True)
subprocess.run(command, check=True)
print(f"[done] {name}", flush=True)
def main() -> None:
parser = argparse.ArgumentParser()
parser.add_argument("config", type=Path)
parser.add_argument("--dry-run", action="store_true")
args = parser.parse_args()
config = yaml.safe_load(args.config.read_text(encoding="utf-8"))
jobs = config["jobs"]
order = execution_order(jobs)
print("Execution order:", " -> ".join(order))
if not args.dry_run:
for name in order:
run_job(name, jobs[name])
if __name__ == "__main__":
main()
在还没有准备三个任务脚本时,可以先验证 DAG:
python orchestrator.py training.yaml --dry-run
预期输出类似:
Execution order: prepare_features -> train_model -> evaluate_model
生产实现不应直接用空格拆分任意 shell 命令,也不应无条件执行用户提供的字符串。更稳妥的做法是让配置引用经过注册的任务类型或容器镜像,并对参数、镜像来源及执行身份进行校验。
统一编排之后,治理能力才是关键
把训练脚本迁移到统一框架只是第一步。平台能否长期发挥作用,取决于它是否覆盖开发者真正需要的工程能力:
- 可复现性:记录代码版本、配置版本、数据快照、依赖环境和随机种子。
- 可观测性:为每次运行及每个 DAG 节点生成结构化日志、耗时、资源消耗和错误分类。
- 幂等与重试:明确哪些任务可以重跑,避免重试时覆盖有效模型或重复发布。
- 产物契约:校验上游输出的 schema、路径和版本,阻止错误输入传播到训练阶段。
- 渐进迁移:允许旧 Spark 脚本先作为 DAG 节点运行,再逐步拆分公共逻辑。
- 逃生通道:为无法被通用配置覆盖的训练任务提供受控扩展点,但限制其权限和影响范围。
配置驱动并不意味着所有训练任务都必须完全相同。平台应统一运行生命周期、依赖模型和治理规则,同时保留算法代码、特征逻辑及资源需求的差异。如果配置语言不断增加条件分支,最终变成另一种难以调试的编程语言,就需要把复杂逻辑移回有测试覆盖的任务代码中。
采用时的检查清单
团队评估类似 Training Orchestrator 的方案时,可以从四个问题开始:现有脚本中有哪些重复的基础设施逻辑;训练产物能否被稳定版本化;失败后能否从单个节点恢复;平台是否允许团队在不修改编排器核心代码的情况下接入新任务。
Yelp 的方向说明,训练平台的价值不只在于启动 Spark 作业。更重要的是建立一个统一控制面,让依赖、配置、执行状态和模型产物进入同一套可验证流程。对于脚本数量已经开始失控的团队,优先统一任务契约和运行元数据,通常比立即重写所有训练代码更稳妥。