用规格驱动组合,让数据流水线少复制、多复用

2026-07-10 20 预计阅读时间: 1 分钟
来源: aws.amazon.com AI 摘要 Original link

Disclaimer: This article is an AI-assisted summary. Read it together with the original source when precision matters. The summary may omit context, version differences, or edge cases and is not official documentation.

预计阅读时间:8 分钟

很多数据流水线一开始只是几段脚本:读表、清洗、聚合、写出。问题通常不是第一条流水线,而是第十条、第五十条。团队开始复制转换逻辑,改几个字段名,换一个过滤条件,然后小变更在多个脚本里连锁扩散。规格驱动组合的核心思路是:把“做什么”写成可追踪的规格,把“怎么做”沉淀成可复用的组件,再由组合层把它们装配成具体工作流。

脚本变多后,真正失控的是意图

传统脚本式流水线的问题不只是代码重复。更麻烦的是,没人能快速回答这些问题:

  • 这条流水线到底做了哪些转换?
  • 某个字段在哪些工作流里被重命名、过滤或聚合?
  • 一个业务规则变化后,需要改哪些脚本?
  • 两条看起来相似的流水线,差异到底在哪里?

复制代码会让转换逻辑分散在多个文件中。每个脚本都携带一份“局部真相”,时间久了,团队会开始害怕修改。规格驱动组合试图把这些局部真相拉回到一个更清晰的层次:用声明式规格描述输入、步骤、参数和输出,让流水线更像配置出来的产品,而不是一次性脚本。

规格不是配置杂物箱,而是工作流契约

“规格驱动”容易被误解成“把所有东西都塞进 YAML”。这不是重点。重点是规格要表达稳定的工作流契约:

  • 输入数据源是什么;
  • 需要执行哪些转换;
  • 每个转换的参数是什么;
  • 输出写到哪里;
  • 哪些规则应该可审计、可比较、可复用。

转换函数仍然应该写在代码里,并接受明确参数。规格只负责选择和组合这些函数。这样做的好处是,转换逻辑可以测试,工作流差异可以通过规格比较,新增流水线也不必复制已有脚本。

可以把它理解成两层:组件层提供 filter_rowsrename_columnsaggregate 这类积木;规格层声明“这条流水线要用哪些积木,以什么参数拼起来”。

可以这样实践:用 YAML 组合一条 Pandas 流水线

下面是一个最小可运行示例,用 YAML 描述数据工作流,用 Python 解释执行。它不是某个特定产品 API,而是演示规格驱动组合的实现方式。运行前需要安装 pandaspyyaml

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 或某个框架本身,而在于它强迫团队把数据工作流的意图显式化。脚本仍然存在,但不再是唯一事实来源。对于规模持续增长的数据团队,这个边界会越来越重要。


相关推荐