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 策略与数据验证的完整方法论。