很多数据流水线一开始只是几段脚本:读表、清洗、聚合、写出。问题通常不是第一条流水线,而是第十条、第五十条。团队开始复制转换逻辑,改几个字段名,换一个过滤条件,然后小变更在多个脚本里连锁扩散。规格驱动组合的核心思路是:把“做什么”写成可追踪的规格,把“怎么做”沉淀成可复用的组件,再由组合层把它们装配成具体工作流。
脚本变多后,真正失控的是意图
传统脚本式流水线的问题不只是代码重复。更麻烦的是,没人能快速回答这些问题:
- 这条流水线到底做了哪些转换?
- 某个字段在哪些工作流里被重命名、过滤或聚合?
- 一个业务规则变化后,需要改哪些脚本?
- 两条看起来相似的流水线,差异到底在哪里?
复制代码会让转换逻辑分散在多个文件中。每个脚本都携带一份“局部真相”,时间久了,团队会开始害怕修改。规格驱动组合试图把这些局部真相拉回到一个更清晰的层次:用声明式规格描述输入、步骤、参数和输出,让流水线更像配置出来的产品,而不是一次性脚本。
规格不是配置杂物箱,而是工作流契约
“规格驱动”容易被误解成“把所有东西都塞进 YAML”。这不是重点。重点是规格要表达稳定的工作流契约:
- 输入数据源是什么;
- 需要执行哪些转换;
- 每个转换的参数是什么;
- 输出写到哪里;
- 哪些规则应该可审计、可比较、可复用。
转换函数仍然应该写在代码里,并接受明确参数。规格只负责选择和组合这些函数。这样做的好处是,转换逻辑可以测试,工作流差异可以通过规格比较,新增流水线也不必复制已有脚本。
可以把它理解成两层:组件层提供 filter_rows、rename_columns、aggregate 这类积木;规格层声明“这条流水线要用哪些积木,以什么参数拼起来”。
可以这样实践:用 YAML 组合一条 Pandas 流水线
下面是一个最小可运行示例,用 YAML 描述数据工作流,用 Python 解释执行。它不是某个特定产品 API,而是演示规格驱动组合的实现方式。运行前需要安装 pandas 和 pyyaml。
python -m venv .venv
source .venv/bin/activate
pip install pandas pyyaml
创建 pipeline.yaml:
input:
path: orders.csv
steps:
- op: filter_equals
column: status
value: paid
- op: rename_columns
mapping:
user_id: customer_id
amount: revenue
- op: aggregate_sum
group_by: customer_id
column: revenue
output_column: total_revenue
output:
path: customer_revenue.csv
创建 run_pipeline.py:
from __future__ import annotations
import sys
from typing import Any, Callable
import pandas as pd
import yaml
def filter_equals(df: pd.DataFrame, *, column: str, value: Any) -> pd.DataFrame:
return df[df[column] == value].copy()
def rename_columns(df: pd.DataFrame, *, mapping: dict[str, str]) -> pd.DataFrame:
return df.rename(columns=mapping)
def aggregate_sum(
df: pd.DataFrame,
*,
group_by: str,
column: str,
output_column: str,
) -> pd.DataFrame:
return (
df.groupby(group_by, as_index=False)[column]
.sum()
.rename(columns={column: output_column})
)
OPERATIONS: dict[str, Callable[..., pd.DataFrame]] = {
"filter_equals": filter_equals,
"rename_columns": rename_columns,
"aggregate_sum": aggregate_sum,
}
def run(spec_path: str) -> None:
with open(spec_path, "r", encoding="utf-8") as f:
spec = yaml.safe_load(f)
df = pd.read_csv(spec["input"]["path"])
for step in spec["steps"]:
op_name = step["op"]
params = {k: v for k, v in step.items() if k != "op"}
if op_name not in OPERATIONS:
raise ValueError(f"Unknown operation: {op_name}")
df = OPERATIONS[op_name](df, **params)
df.to_csv(spec["output"]["path"], index=False)
print(f"Wrote {len(df)} rows to {spec['output']['path']}")
if __name__ == "__main__":
run(sys.argv[1])
准备一份 orders.csv:
order_id,user_id,status,amount
1,u1,paid,20
2,u1,cancelled,10
3,u2,paid,35
4,u1,paid,15
运行:
python run_pipeline.py pipeline.yaml
cat customer_revenue.csv
预期输出类似:
customer_id,total_revenue
u1,35
u2,35
这个例子很小,但已经体现出关键边界:新增一条类似流水线时,可以新增一份规格文件,而不是复制 run_pipeline.py。如果多个流水线都使用 aggregate_sum,业务规则修正也集中在一个函数里。
组合层需要治理,不然会变成另一种混乱
规格驱动组合不是免费午餐。它把复杂度从脚本复制转移到了组件设计、规格校验和版本管理上。几个工程约束很关键:
- 操作集合要克制。不要为每个临时需求都新增一个高度定制的
op。 - 规格需要 schema 校验。否则错误只会从 Python 脚本变成 YAML 拼写错误。
- 转换函数要有单元测试。规格能减少复制,但不能替代测试。
- 参数要显式。隐式读取全局配置会让工作流重新变得难追踪。
- 规格变更要能 code review。数据工作流的行为变化不应该藏在运行时控制台里。
可以继续给上面的例子加一层 JSON Schema 或 Pydantic 校验。对于生产系统,还要考虑数据血缘、运行日志、失败重试、权限隔离和版本化发布。
适合什么时候采用
当团队只有两三条一次性脚本时,规格驱动组合可能显得重。真正值得引入它的信号通常是:同类转换逻辑开始被复制;业务方频繁要求“差不多但略有不同”的数据集;工程师很难说清某条流水线和另一条的差异;字段或规则变更需要全仓库搜索。
比较稳妥的落地路径是从重复最多的转换开始,把它们提成小而稳定的组件,再为新流水线引入规格文件。不要一口气重写所有历史脚本。先让一两类高频工作流变得可组合、可比较、可审计,再扩大范围。
规格驱动组合的价值不在于 YAML、JSON 或某个框架本身,而在于它强迫团队把数据工作流的意图显式化。脚本仍然存在,但不再是唯一事实来源。对于规模持续增长的数据团队,这个边界会越来越重要。