后端开发
数据 Transform:ETL 管道设计实践指南
{}
ETL -- 数据工程的核心范式
ETL(Extract-Transform-Load)是数据仓库建设的基石。Transform 是其中最复杂的环节,负责将原始数据清洗、转换、映射为目标系统的标准格式。
Transform 的三大职责
1. 数据清洗 Data Cleaning
处理缺失值、异常值、重复数据:
# Python 数据清洗示例
import pandas as pd
df = pd.read_csv("raw_data.csv")
# 处理缺失值
df["age"].fillna(df["age"].median(), inplace=True)
df.dropna(subset=["user_id"], inplace=True)
# 去重
df.drop_duplicates(subset=["user_id", "date"], keep="last", inplace=True)
# 异常值过滤
df = df[(df["amount"] > 0) & (df["amount"] < 100000)]
2. 类型映射 Type Mapping
源系统与目标系统的数据类型往往不一致,需要做 Transform 映射:
// JSON 源数据 -> 数据库目标类型
const typeMap = {
"string": "VARCHAR(255)",
"integer": "BIGINT",
"float": "DOUBLE",
"boolean": "TINYINT(1)",
"datetime": "TIMESTAMP"
};
function transformSchema(sourceSchema) {
return sourceSchema.map(field => ({
...field,
target_type: typeMap[field.type] || "TEXT"
}));
}
3. Schema 演化 Schema Evolution
源系统新增字段时,Transform 层需要向后兼容。常用策略包括:
- 前向兼容:新 schema 读取旧数据时,新增字段用默认值填充
- 后向兼容:旧 schema 读取新数据时,忽略未知字段
- 全兼容:使用 Avro/Protobuf 等支持 schema 注册的格式
Transform 管道架构
// 现代 ETL 管道设计
class TransformPipeline {
constructor() {
this.steps = [];
}
addStep(name, fn) {
this.steps.push({ name, fn });
return this;
}
async run(data) {
let result = data;
for (const step of this.steps) {
console.log(`[Transform] ${step.name}...`);
result = await step.fn(result);
}
return result;
}
}
const pipeline = new TransformPipeline()
.addStep("clean_nulls", cleanNulls)
.addStep("validate_schema", validateSchema)
.addStep("map_types", mapTypes)
.addStep("enrich", enrichData)
.addStep("deduplicate", deduplicate);
const result = await pipeline.run(rawData);
性能优化策略
- 批量处理:避免逐行 Transform,使用批量操作减少 IO 次数
- 列式存储:Parquet/ORC 格式只读取需要的列,减少 IO
- 增量 Transform:只处理变更数据(CDC),而非全量重算
- 内存计算:Spark/Flink 利用内存缓存中间结果
现代 ELT 趋势
随着云数据仓库(Snowflake, BigQuery)的兴起,传统 ETL 正向 ELT 演进:先 Load 原始数据到仓库,再在仓库内做 Transform。这利用了仓库的弹性算力,简化了管道架构。
数据 Transform 的核心不是代码,而是对数据本身的理解。最好的管道设计,是对业务逻辑的精确映射。