数据平台接入一个新数据集,真正耗时的往往不是复制文件,而是反复确认字段、分类级别、处理步骤、质量规则和责任人,再把这些要求写进一套专用代码。AWS 介绍的规格驱动组合(specification-driven composition)换了一个切入点:用声明式规格表达意图,把实际处理交给可复用能力,并在执行前完成验证。AWS 报告称,这种方式可以把数据集接入周期从数周缩短到数天,同时保留追踪、版本控制、数据分类和治理能力。
从“每个数据集一套流程”转向“意图加能力目录”
传统数据工作流常把业务要求和实现细节揉在一起。比如,为订单数据编写一条流水线时,字段映射、脱敏、格式转换、质量检查和落盘路径都可能直接固化在 DAG 或脚本中。下一个数据集即使只改了几个字段,也要复制并修改整套流程。
规格驱动方式可以拆成三层:
- 声明式规格:描述数据是什么、要达到什么状态、适用哪些治理要求。
- 可复用处理能力:实现标准化、校验、分类、脱敏、发布等动作。
- 组合与验证引擎:读取规格,检查其完整性和兼容性,再组装并执行工作流。
关键变化是:规格只说明“需要什么”,处理组件负责“怎么做”。因此,接入新数据集更像提交一份经过校验的配置,而不是开发一条全新流水线。
一份规格应当表达什么
下面是一个可以直接改造的 YAML 示例。它不是 AWS 公布 API 的逐项复刻,而是一种可实践的最小设计:
apiVersion: data.example.com/v1
kind: DatasetWorkflow
metadata:
name: customer-orders
version: 1.2.0
owner: commerce-data
spec:
source:
format: csv
uri: s3://incoming-data/customer-orders/
schema:
required:
- order_id
- customer_email
- amount
- created_at
classification:
default: internal
fields:
customer_email: pii
steps:
- capability: validate-schema
- capability: mask-fields
inputs:
fields:
- customer_email
- capability: check-quality
inputs:
rules:
- column: amount
operator: greater_than_or_equal
value: 0
- capability: publish-parquet
inputs:
destination: s3://curated-data/customer-orders/
governance:
retentionDays: 365
approvalRequired: true
这份规格同时承载了四类信息:输入位置与格式、结构契约、字段分类和处理意图。capability 名称不应对应任意脚本,而应来自经过登记、测试和授权的能力目录。这样,平台才能知道哪些组件可以组合、需要哪些参数,以及它们会产生什么输出。
规格还应有明确版本。版本变化不仅用于回滚,也能回答治理场景中的实际问题:某批数据由哪版规则处理、当时哪些字段被判定为敏感信息、发布前通过了哪些检查。
执行前验证比运行时失败更有价值
声明式方案并不意味着“YAML 写对了就能跑”。真正降低接入成本的是执行前验证:在占用计算资源、读取生产数据之前,尽早发现缺失字段、未知能力、参数错误和治理冲突。
下面的 Python 脚本可以作为本地或 CI 中的最小验证器。运行前安装 PyYAML,并把上面的规格保存为 workflow.yaml:
python -m pip install PyYAML
python validate_workflow.py workflow.yaml
#!/usr/bin/env python3
import sys
from pathlib import Path
import yaml
KNOWN_CAPABILITIES = {
"validate-schema",
"mask-fields",
"check-quality",
"publish-parquet",
}
def validate(document: dict) -> list[str]:
errors = []
metadata = document.get("metadata", {})
spec = document.get("spec", {})
for field in ("name", "version", "owner"):
if not metadata.get(field):
errors.append(f"metadata.{field} is required")
steps = spec.get("steps", [])
if not steps:
errors.append("spec.steps must contain at least one capability")
for index, step in enumerate(steps):
capability = step.get("capability")
if capability not in KNOWN_CAPABILITIES:
errors.append(
f"spec.steps[{index}].capability is unknown: {capability!r}"
)
classified_fields = set(
spec.get("classification", {}).get("fields", {}).keys()
)
masked_fields = {
field
for step in steps
if step.get("capability") == "mask-fields"
for field in step.get("inputs", {}).get("fields", [])
}
missing_masks = classified_fields - masked_fields
if missing_masks:
errors.append(
"classified fields are not masked: " + ", ".join(sorted(missing_masks))
)
return errors
def main() -> int:
path = Path(sys.argv[1])
document = yaml.safe_load(path.read_text(encoding="utf-8"))
errors = validate(document)
if errors:
print("Workflow specification is invalid:")
for error in errors:
print(f"- {error}")
return 1
print("Workflow specification is valid")
return 0
if __name__ == "__main__":
raise SystemExit(main())
这个示例只验证了基础约束,但已经展示了核心模式:能力必须来自白名单,元数据必须完整,被分类的敏感字段必须进入脱敏步骤。在生产环境中,还可以加入 JSON Schema 校验、能力输入输出类型检查、数据驻留策略、审批状态和成本上限。
可追踪性不能只停留在 Git 历史
把规格放进 Git 是一个合理起点,但治理链路还需要连接规格版本、执行记录和产出数据。每次执行至少应记录:
- 规格名称、版本和内容摘要;
- 被解析出的能力及其组件版本;
- 输入数据标识和输出数据标识;
- 验证结果、审批结果和执行状态;
- 分类规则与质量规则的判定结果。
这样才能从一张下游表反查它由哪份规格生成,也能在规则变化后判断哪些历史数据需要重跑。单纯保存最新 YAML 无法提供这种证据链。
能力目录同样需要治理。组件一旦被多个数据集复用,它的升级就有更大影响面。建议为能力定义明确的输入输出契约、语义化版本和兼容策略;涉及删除字段、改变脱敏算法或调整质量阈值含义时,应发布新主版本,而不是静默修改。
落地时从边界清晰的流程开始
规格驱动组合适合规则可结构化、处理动作可复用的数据接入场景,但不会自动消除复杂度。它只是把复杂度从散落的专用代码集中到规格模型、能力目录和验证系统中。如果能力边界模糊,最终可能得到一批名称统一、行为各异的脚本。
实际采用时,可以按这份清单推进:
- 先选择结构稳定、合规要求明确的一个数据域,不要一开始覆盖所有流水线。
- 从 5 到 10 个高频能力入手,例如结构校验、字段分类、脱敏、质量检查和标准格式发布。
- 为规格建立机器可执行的 schema,并把验证放进拉取请求和部署入口。
- 让每次执行绑定不可变的规格版本、能力版本和数据标识。
- 同时衡量接入周期、失败前移比例、能力复用率和变更影响面。
AWS 所描述的价值不只是“用配置代替代码”,而是建立一份可验证、可版本化、可审计的数据处理契约。做到这一点后,数据集接入速度和治理要求才不必互相牺牲。