4 minutes
数据转换(Transform):清洗与加工
Transform:ETL 的核心环节
如果说抽取是"采购食材",加载是"摆盘上桌",那么转换就是整个 ETL 流程中的"烹饪"环节。这也是数据工程师投入精力最多、代码量最大的部分。
原始数据通常充满了各种问题:字段缺失、格式不统一、存在重复记录、数值异常等等。转换环节的目标就是把这些"脏数据"变成干净、一致、可分析的高质量数据。毫不夸张地说,一个 ETL 流程的质量高低,几乎完全取决于转换环节做得怎么样。
常见的数据质量问题
在实际项目中,我们经常遇到以下几类数据质量问题:
缺失值
数据集中某些字段为空,这是最常见的问题。缺失的原因可能是源系统没有采集到、传输过程中丢失、或者该字段本身就不适用。
重复记录
同一行数据在数据集中出现了多次,通常是因为源系统没有正确的唯一约束、多次导入未去重、或者数据合并时出现了重复。
异常值
数据值明显超出正常范围。例如,年龄字段出现 200 岁、订单金额为负数、温度字段出现 9999 等。异常值可能是录入错误,也可能是正常但罕见的业务场景。
格式不一致
同一类型的数据使用了不同的表示方式。例如日期字段有的用"2026-04-05",有的用"2026/04/05",还有的用"05/04/2026";电话号码有的带国家区号,有的不带。
| 问题类型 | 示例 | 原因 | 影响 |
|---|---|---|---|
| 缺失值 | email 列为 NULL |
未采集、未必填 | 统计分析偏差 |
| 重复记录 | 同一订单出现 2 次 | 重复导入 | 汇总数据翻倍 |
| 异常值 | 年龄=500 | 录入错误 | 平均值失真 |
| 格式不一致 | 日期格式混杂 | 多源合并 | 解析失败 |
| 逻辑矛盾 | 下单日期 > 发货日期 | 业务约束违反 | 数据不可信 |
使用 Pandas 进行数据清洗
Pandas 是 Python 生态中最强大的数据处理库,也是 ETL 转换环节的首选工具。
缺失值处理
import pandas as pd
import numpy as np
# 加载原始数据
df = pd.read_csv("raw_customers.csv")
print("原始数据概览:")
print(df.info())
print(f"缺失值统计:\n{df.isnull().sum()}")
# 策略一:删除缺失值较多的列
threshold = len(df) * 0.5 # 如果缺失超过 50% 则删除该列
df_cleaned = df.dropna(thresh=threshold, axis=1)
# 策略二:用平均值/中位数/众数填充
df_cleaned["age"].fillna(df["age"].median(), inplace=True)
df_cleaned["income"].fillna(df["income"].mean(), inplace=True)
df_cleaned["gender"].fillna(df["gender"].mode()[0], inplace=True)
# 策略三:前向填充(适用于时间序列数据)
df_cleaned["last_login"].fillna(method="ffill", inplace=True)
# 策略四:用业务规则填充
# 如果 email 缺失但 user_id 存在,可以标记为"未知"
df_cleaned["email"].fillna("unknown@placeholder.com", inplace=True)
重复数据处理
# 查找重复行
duplicates = df.duplicated()
print(f"完全重复的行数: {duplicates.sum()}")
# 基于特定列查找重复
id_duplicates = df.duplicated(subset=["customer_id"], keep="first")
print(f"customer_id 重复的行数: {id_duplicates.sum()}")
# 查看重复记录
print(df[df.duplicated(subset=["customer_id"], keep=False)].head(10))
# 删除重复,保留第一条
df_deduped = df.drop_duplicates(subset=["customer_id"], keep="first")
# 更精细的去重:保留时间戳最新的那条
df_deduped = (
df.sort_values("updated_at", ascending=False)
.drop_duplicates(subset=["customer_id"], keep="first")
)
异常值检测与处理
# 方法一:基于 Z-Score 检测
from scipy import stats
z_scores = np.abs(stats.zscore(df["amount"]))
df_no_outliers = df[z_scores < 3] # 保留 Z-Score 在 3 以内的数据
print(f"Z-Score 方法移除了 {len(df) - len(df_no_outliers)} 条异常值")
# 方法二:基于 IQR(四分位距)检测
Q1 = df["amount"].quantile(0.25)
Q3 = df["amount"].quantile(0.75)
IQR = Q3 - Q1
lower_bound = Q1 - 1.5 * IQR
upper_bound = Q3 + 1.5 * IQR
df_iqr_filtered = df[(df["amount"] >= lower_bound) & (df["amount"] <= upper_bound)]
print(f"IQR 方法移除了 {len(df) - len(df_iqr_filtered)} 条异常值")
# 方法三:基于业务规则的过滤
df_valid = df[
(df["age"].between(0, 120)) &
(df["amount"] >= 0) &
(df["quantity"] > 0)
]
数据类型转换
将数据从源系统的类型映射到目标系统的类型,是转换环节的基础任务。
# 字符串转日期
df["order_date"] = pd.to_datetime(df["order_date"], format="%Y-%m-%d")
# 字符串转数值(处理千分位分隔符和货币符号)
df["price"] = (
df["price"]
.str.replace("$", "", regex=False)
.str.replace(",", "", regex=False)
.astype(float)
)
# 类别编码:将文本标签转换为数字
from sklearn.preprocessing import LabelEncoder
le = LabelEncoder()
df["category_code"] = le.fit_transform(df["category_name"])
# 或者使用 Pandas 的 categorical 类型(适合有序类别)
df["tier"] = pd.Categorical(
df["tier"],
categories=["bronze", "silver", "gold", "platinum"],
ordered=True
)
# 布尔值转换
df["is_active"] = df["status"].map({"active": True, "inactive": False})
数据聚合与关联
SQL 风格的 GROUP BY 操作
# 按类别计算销售汇总
sales_summary = df.groupby("category").agg(
总销售额=("amount", "sum"),
订单数=("order_id", "count"),
平均金额=("amount", "mean"),
最大金额=("amount", "max"),
客户数=("customer_id", "nunique")
).reset_index()
print(sales_summary)
# 多层级分组
monthly_category_sales = df.groupby(
[df["order_date"].dt.to_period("M"), "category"]
)["amount"].sum().reset_index()
多表关联
# 读取多张表
orders = pd.read_csv("orders.csv")
customers = pd.read_csv("customers.csv")
order_items = pd.read_csv("order_items.csv")
# 链式 JOIN 构建宽表
fact_orders = (
orders
.merge(customers, on="customer_id", how="left")
.merge(order_items, on="order_id", how="left")
)
print(f"关联后的宽表: {fact_orders.shape[0]} 行, {fact_orders.shape[1]} 列")
使用 SQL 进行数据转换
很多时候,直接在数据库中执行 SQL 转换比把数据拉到应用层处理更高效。SQL 的集合操作能力非常强大,而且数据库引擎会对查询进行优化。
-- 数据清洗:处理缺失值和异常值
CREATE OR REPLACE VIEW clean_orders AS
SELECT
order_id,
customer_id,
COALESCE(order_amount, 0) AS order_amount, -- NULL 替换为 0
CASE
WHEN order_amount < 0 THEN 0 -- 负值归零
WHEN order_amount > 100000 THEN NULL -- 异常值标记
ELSE order_amount
END AS cleaned_amount,
TO_DATE(order_date, 'YYYY-MM-DD') AS order_date, -- 统一日期格式
UPPER(TRIM(customer_email)) AS customer_email -- 统一大小写、去空格
FROM raw_orders
WHERE order_id IS NOT NULL; -- 剔除无主键的记录
-- 数据聚合:月度销售统计
SELECT
DATE_TRUNC('month', order_date) AS month,
customer_tier,
COUNT(DISTINCT order_id) AS order_count,
SUM(cleaned_amount) AS total_revenue,
AVG(cleaned_amount) AS avg_order_value
FROM clean_orders o
JOIN dim_customers c ON o.customer_id = c.customer_id
WHERE cleaned_amount IS NOT NULL
GROUP BY 1, 2
ORDER BY 1, 2;
-- 窗口函数:计算累计销售额和排名
SELECT
order_id,
order_date,
cleaned_amount,
SUM(cleaned_amount) OVER (
PARTITION BY DATE_TRUNC('month', order_date)
ORDER BY order_date
) AS monthly_running_total,
RANK() OVER (
PARTITION BY DATE_TRUNC('month', order_date)
ORDER BY cleaned_amount DESC
) AS monthly_rank
FROM clean_orders;
慢变化维度处理(SCD)
在数据仓库中,维度表(如客户表、产品表)的属性会随时间变化。如何处理这些变化就是 Slowly Changing Dimensions(SCD)要解决的问题。
SCD Type 1:直接覆盖
直接更新属性值,不保留历史记录。适用于纠正错误或不太重要的属性。
-- Type 1:直接在维度表中更新
UPDATE dim_customer
SET email = 'newemail@example.com',
updated_at = CURRENT_TIMESTAMP
WHERE customer_id = 12345;
SCD Type 2:保留历史
当属性变化时,新增一行记录,用生效日期和过期日期来标记版本。这是最常用的模式,能够完整追溯历史。
def apply_scd_type2(
existing_dim: pd.DataFrame,
new_records: pd.DataFrame,
business_key: str,
tracked_columns: list
) -> pd.DataFrame:
"""
实现 SCD Type 2 逻辑
- existing_dim: 当前维度表快照
- new_records: 从源系统抽取的最新数据
- tracked_columns: 需要跟踪变化的列
"""
current_date = pd.Timestamp.now()
for _, new_row in new_records.iterrows():
# 查找已有的业务实体
existing = existing_dim[
existing_dim[business_key] == new_row[business_key]
]
if existing.empty:
# 新增记录
new_row["valid_from"] = current_date
new_row["valid_to"] = pd.NaT
new_row["is_current"] = True
existing_dim = pd.concat([existing_dim, pd.DataFrame([new_row])], ignore_index=True)
else:
# 检查属性是否变化
current_version = existing[existing["is_current"] == True]
if not current_version.empty:
has_changed = any(
current_version[col].iloc[0] != new_row[col]
for col in tracked_columns
)
if has_changed:
# 关闭旧版本
existing_dim.loc[current_version.index, "valid_to"] = current_date
existing_dim.loc[current_version.index, "is_current"] = False
# 创建新版本
new_row["valid_from"] = current_date
new_row["valid_to"] = pd.NaT
new_row["is_current"] = True
existing_dim = pd.concat(
[existing_dim, pd.DataFrame([new_row])], ignore_index=True
)
return existing_dim
SCD Type 3:保留有限历史
在维度表中添加额外的列来保存前后值变化。适用于只需要跟踪最近一次或少数几次变化的场景。
-- Type 3:添加"上一个版本"列
ALTER TABLE dim_customer
ADD COLUMN previous_email VARCHAR(255),
ADD COLUMN email_changed_at TIMESTAMP;
-- 更新时,将当前值移到"上一个"字段
UPDATE dim_customer
SET previous_email = email,
email = 'newemail@example.com',
email_changed_at = CURRENT_TIMESTAMP
WHERE customer_id = 12345;
数据验证与质量检查
在转换环节的最后,应该加入数据质量检查,确保输出数据的可靠性。
def validate_transform_output(df: pd.DataFrame) -> dict:
"""对转换后的数据进行质量检查"""
checks = {
"行数变动": None,
"完整性检查": {},
"唯一性检查": {},
"范围检查": {}
}
# 1. 行数检查
checks["行数变动"] = f"{len(df)} 行"
# 2. 完整性检查:关键字段不应有缺失
for col in ["order_id", "customer_id", "amount"]:
null_count = df[col].isnull().sum()
null_pct = null_count / len(df) * 100
checks["完整性检查"][col] = f"缺失 {null_count} ({null_pct:.2f}%)"
# 3. 唯一性检查
id_dupes = df["order_id"].duplicated().sum()
checks["唯一性检查"]["order_id"] = f"重复 {id_dupes} 条"
# 4. 范围检查
for col, (low, high) in [("amount", (0, 1_000_000)), ("quantity", (0, 1000))]:
out_of_range = ((df[col] < low) | (df[col] > high)).sum()
checks["范围检查"][col] = f"越界 {out_of_range} 条"
return checks
# 执行检查
quality_report = validate_transform_output(transformed_df)
for check_name, result in quality_report.items():
print(f"{check_name}: {result}")
小结
数据转换是 ETL 流程的核心环节,它决定了最终数据质量的上限。本文从常见的数据质量问题出发,介绍了使用 Pandas 进行缺失值处理、重复数据去重、异常值检测的具体方法。我们还讨论了数据类型转换、聚合关联、SQL 转换技巧以及慢变化维度的处理策略。最后,数据验证环节帮助我们确保转换结果的可信度。掌握这些技术,你就能把脏乱的原始数据变成干净、可用的分析数据。
下一篇文章将探讨 ETL 的最后一个环节——数据加载,看看如何高效地将数据写入目标系统。
Summary: 数据清洗、格式转换、SCD 策略与数据验证的完整方法论。