前言
在上一篇文章《ETL入门:从数据集成到数据仓库》中,我们讨论了 ETL 的基本概念和架构。但在实际生产环境中,ETL 最让人头疼的从来不是数据怎么搬,而是搬过来的数据能不能用。
做过数据工程的朋友都懂:80% 的时间不是在写 ETL 逻辑,而是在修数据。字段为空、格式乱码、同一客户出现三次、金额字段里混着中文备注……这些「脏数据」才是 ETL 工程师的日常。我曾经在一个电商项目中遇到过一个经典案例:上游 CRM 系统导出的用户表中,「手机号」一列混入了微信号、QQ 号甚至"无"这样的文本,而金额字段里则夹杂着货币符号和换行符。那一次的数据修复花掉了整整两个迭代周期。
本文将深入 ETL 的核心环节——数据清洗与转换,从实战出发,覆盖质量分类、框架设计、增量抽取、Python 管道实现、质量监控和调度告警,希望能为正在搭建或优化 ETL 管道的你提供一些可落地的思路。
在开始之前,我们先明确一个核心理念:数据清洗不是一次性活动,而是一个持续迭代的过程。每一次接入新数据源、每一次业务逻辑变更,都可能引入新的质量问题。因此,我们不仅要学会「清洗」,更要学会「设计可复用的清洗体系」。
一、数据质量问题的分类
要清洗数据,首先得知道「脏」在哪。经过多年实践,我把常见的数据质量问题归纳为四大类。需要注意的是,实际业务中这些问题往往不是孤立出现的——一个字段可能同时存在缺失和格式不一致的问题,而不同来源的数据合并时,重复和不一致更是家常便饭。因此,分类是为了更好地识别,但设计清洗策略时一定要考虑组合场景。
1. 缺失值(Missing Values)
最普遍的问题,几乎每个 ETL 工程师每天都要面对。表现形式多样:
- **完全为空**:某条记录的 email 字段为 NULL
- **占位符填充**:用 `"N/A"`、`"-"`、`"unknown"` 填充实际为空的数据
- **默认值侵入**:系统用 `"0"` 或 `"1900-01-01"` 表示「无值」
为什么缺失值会如此普遍?一方面,源系统的业务逻辑可能允许某些字段为空(比如用户不强制填写手机号);另一方面,数据在传输过程中可能因网络中断、序列化异常等导致部分字段丢失。在金融场景中,缺失值可能导致风控模型评分异常;在报表场景中,缺失值可能导致聚合结果偏差。
处理缺失值需要区分场景:对于必填字段(如订单 ID、用户 ID),应该直接拒绝写入并告警;对于可选字段(如备注、昵称),可以填充默认值或直接保留为空。
# 缺失值的多种表现
missing_patterns = {
"email": [None, "", "N/A", "NULL", "unknown@unknown.com"],
"age": [None, 0, -1, 999],
"date": [None, "", "1900-01-01", "1970-01-01"]
}
2. 异常值(Outliers / Anomalies)
在数值字段中出现不符合业务逻辑的值:
- 年龄:200 岁、-5 岁
- 金额:-999999.00 或 999999999.99(可能是测试数据)
- 时间:未来时间、1900 年的时间
异常值的检测方法有很多,除了经典的 IQR(四分位距)方法外,还可以结合业务规则做硬边界校验。比如金额字段,虽然 IQR 可能认为 98000 元是异常值(因为大多数订单只有几十到几百元),但从业务角度它可能是一个合法的大额订单。因此我通常建议统计方法 + 业务规则双管齐下:先用统计方法圈定可疑值,再用业务规则做最终判定。
import pandas as pd
import numpy as np
def detect_outliers_iqr(df: pd.DataFrame, column: str) -> pd.Series:
"""基于 IQR 方法检测异常值"""
Q1 = df[column].quantile(0.25)
Q3 = df[column].quantile(0.75)
IQR = Q3 - Q1
lower = Q1 - 1.5 * IQR
upper = Q3 + 1.5 * IQR
return (df[column] < lower) | (df[column] > upper)
def detect_outliers_business(df: pd.DataFrame, column: str,
min_val: float, max_val: float) -> pd.Series:
"""基于业务规则的异常值检测"""
return (df[column] < min_val) | (df[column] > max_val)
3. 重复数据(Duplicates)
记录层面和字段层面的重复,在实际生产中非常常见,尤其是在多源数据合并的场景下。比如你从三个不同的渠道(官网、小程序、线下 POS)接收订单数据,同一位客户可能会在不同系统中产生多条看似独立但实际关联的记录。
重复数据可以分为三个层次:
- **完全重复**:整行完全一致,通常由数据重传或网络重试导致
- **部分重复**:关键字段一致(如同一身份证号、同一手机号),但其他字段有差异,需要根据业务规则决定合并策略
- **近似重复**:姓名「张三」和「张 三」、「张san」,需要模糊匹配甚至人工介入
处理重复数据时有一个容易踩的坑:不是所有重复都要删除。有些场景下,重复数据本身携带了有价值的信息(比如同一用户在不同时间点的状态变更),贸然去重反而会丢失历史轨迹。因此我建议区分「业务主键去重」和「记录级去重」,前者严格按业务键去重,后者只在整行完全一致时才去除。
# 多层去重策略
def deduplicate(df: pd.DataFrame, keys: list[str]) -> pd.DataFrame:
# 第一步:完全重复去除
df = df.drop_duplicates()
# 第二步:基于业务主键去重,保留最新记录
df = df.sort_values("updated_at").drop_duplicates(subset=keys, keep="last")
return df
4. 不一致数据(Inconsistent Data)
同一含义的数据在不同来源中格式迥异,这是多源数据集成时最大的痛点之一。我之前处理过一个跨国电商项目,同一笔订单的金额在不同系统中分别以 USD、CNY、EUR 计价,而且数字格式也不同——美国系统用 1,234.56,欧洲系统用 `1.234,56`,如果不做标准化,聚合后的报表完全不可用。
常见的不一致场景包括:
- **日期格式**:`2024-01-15` vs `01/15/2024` vs `2024年1月15日` vs `15-Jan-2024`
- **性别编码**:`M/F` vs `1/0` vs `男/女` vs `Male/Female`
- **货币单位**:`USD 100` vs `$100` vs `100美元` vs `100`
- **地址格式**:不同国家、不同系统的地址字段结构完全不一致
格式不一致的处理思路是「先识别、后转换」。识别阶段需要通过正则或元数据判断当前格式,转换阶段则统一映射到目标格式。需要特别注意的是,日期解析是出错率最高的环节——03/04/2024 到底是 3 月 4 日还是 4 月 3 日?如果不明确来源,不要猜测,直接告警。
date_formats = {
"source_a": "%Y-%m-%d",
"source_b": "%m/%d/%Y",
"source_c": "%Y年%m月%d日"
}
def normalize_date(value: str, source: str) -> str:
fmt = date_formats.get(source)
if fmt:
dt = datetime.strptime(value.strip(), fmt)
return dt.strftime("%Y-%m-%d") # 统一输出
raise ValueError(f"Unknown date format for source: {source}")
二、数据清洗策略与框架设计
2.1 清洗策略矩阵
不要试图在一个函数里解决所有脏数据。我推荐按严重程度和**处理方式**两个维度分类。这个矩阵的核心思想是:不是所有脏数据都值得修复,也不是所有错误都要立即告警。我们要把有限的精力花在影响最大的问题上。
| 严重程度 | 处理方式 | 典型场景 | 响应方式 |
|---------|---------|---------|---------|
| 致命 | 拒绝写入(Reject) | 缺少必填字段、主键冲突 | 立即告警 + 阻断流程 |
| 高 | 修复 + 告警(Repair & Alert) | 格式不一致、明显异常值 | 自动修复 + 通知负责人 |
| 中 | 自动修复(Auto-fix) | 空白填充、类型转换 | 静默修复 + 日志记录 |
| 低 | 标记 + 跳过(Tag & Skip) | 备注字段中的特殊字符 | 仅日志记录 |
举个例子:订单表中的 order_id 为空是致命问题,因为没有订单 ID 的记录在目标系统中没有任何意义,直接拒绝;而用户填写的手机号多了几个空格,属于中等问题,自动 trim 即可,无需打扰任何人。
2.2 Cleaner 模式架构
在实际项目中,我倾向将每个清洗规则封装为独立的 Cleaner,形成管道链:
from abc import ABC, abstractmethod
from dataclasses import dataclass, field
from typing import Optional, Any
import logging
logger = logging.getLogger(__name__)
@dataclass
class ValidationResult:
row_id: Optional[str]
field: str
severity: str # fatal / error / warning / info
message: str
raw_value: Any
class BaseCleaner(ABC):
"""清洗器基类——每个清洗规则是一个独立单元"""
def __init__(self, name: str, severity: str = "error"):
self.name = name
self.severity = severity
@abstractmethod
def validate(self, df: pd.DataFrame) -> list[ValidationResult]:
...
@abstractmethod
def fix(self, df: pd.DataFrame) -> pd.DataFrame:
...
def run(self, df: pd.DataFrame) -> tuple[pd.DataFrame, list[ValidationResult]]:
"""执行校验 + 修复"""
issues = self.validate(df)
df_clean = self.fix(df)
for issue in issues:
log_func = logger.warning if issue.severity in ("error",) else logger.info
log_func(f"[{issue.severity.upper()}] {issue.field}: {issue.message}")
return df_clean, issues
2.3 管道路由
清洗管道不应该是「全部套用」,而是根据数据源动态组装:
class CleaningPipeline:
def __init__(self):
self._cleaners: list[BaseCleaner] = []
self._rejects: list[ValidationResult] = []
def add_cleaner(self, cleaner: BaseCleaner):
self._cleaners.append(cleaner)
def run(self, df: pd.DataFrame, reject_on_fatal: bool = True) -> pd.DataFrame:
all_issues = []
for cleaner in self._cleaners:
df, issues = cleaner.run(df)
all_issues.extend(issues)
# 致命错误记录拒绝
if reject_on_fatal:
fatal_issues = [i for i in all_issues if i.severity == "fatal"]
if fatal_issues:
self._rejects.extend(fatal_issues)
fatal_ids = {i.row_id for i in fatal_issues if i.row_id}
df = df[~df.index.isin(fatal_ids)]
logger.warning(f"Rejected {len(fatal_ids)} records due to fatal errors")
return df
这种设计的优势是:每个 Cleaner 独立可测试,管道可插拔可复用,严重级别可动态调整。
三、常见数据转换操作
清洗之后就是转换。如果说清洗是做「减法」(去除脏数据),那转换就是在做「映射」和「加法」——把数据从一种形态变成另一种形态,并赋予新的业务含义。
转换的本质是:将数据从源端的业务语义映射到目标端的标准语义。比如源端叫 created,目标端叫 `created_at`;源端存的是时间戳数字,目标端要的是 ISO 日期字符串。这些差异都必须在转换层统一处理。
下面介绍三种最常见的数据转换操作,几乎每个 ETL 项目都会用到。
3.1 类型转换与格式标准化
上游系统可能是 MySQL(严格类型),也可能是 CSV(全是字符串),还可能是 MongoDB(灵活 schema)。作为 ETL 工程师,你需要主动做类型断言:
def safe_type_cast(df: pd.DataFrame, mapping: dict[str, str]) -> pd.DataFrame:
"""安全的类型转换,失败时置空并告警"""
for col, target_type in mapping.items():
if col not in df.columns:
continue
try:
if target_type == "int":
df[col] = pd.to_numeric(df[col], errors="coerce").astype("Int64")
elif target_type == "float":
df[col] = pd.to_numeric(df[col], errors="coerce")
elif target_type == "datetime":
df[col] = pd.to_datetime(df[col], errors="coerce")
elif target_type == "str":
df[col] = df[col].astype(str)
except Exception as e:
logger.error(f"Type cast failed for {col}: {e}")
return df
# 使用示例
type_mapping = {
"order_id": "int",
"amount": "float",
"created_at": "datetime",
"user_name": "str"
}
df = safe_type_cast(df, type_mapping)
3.2 字段映射(Field Mapping)
来源于不同系统的字段名和字段含义往往不同,需要建立映射层:
FIELD_MAPPING = {
"src_order_id": "order_id",
"src_user_id": "user_id",
"order_amount": "amount",
"order_date": "created_at",
}
def apply_field_mapping(df: pd.DataFrame, mapping: dict[str, str]) -> pd.DataFrame:
return df.rename(columns=mapping)
更复杂的场景还需要做值映射(Value Mapping):
VALUE_MAPPING = {
"gender": {
"M": "male", "F": "female",
"1": "male", "0": "female",
"男": "male", "女": "female"
},
"status": {
"1": "active", "2": "inactive", "3": "deleted",
"Y": "active", "N": "inactive"
}
}
def apply_value_mapping(df: pd.DataFrame, mapping: dict) -> pd.DataFrame:
for col, mapper in mapping.items():
if col in df.columns:
df[col] = df[col].map(mapper).fillna(df[col])
return df
3.3 聚合计算与衍生字段
在转换层做轻度聚合可以大幅降低下游分析的计算压力:
def add_derived_fields(df: pd.DataFrame) -> pd.DataFrame:
"""添加业务常用的衍生字段"""
# 订单金额分类
df["amount_category"] = pd.cut(
df["amount"],
bins=[0, 100, 1000, 10000, float("inf")],
labels=["小额", "中额", "大额", "超大额"]
)
# 时间维度拆解
if "created_at" in df.columns:
df["order_year"] = df["created_at"].dt.year
df["order_month"] = df["created_at"].dt.month
df["order_weekday"] = df["created_at"].dt.weekday
df["is_weekend"] = df["order_weekday"].isin([5, 6])
# 用户首单标记(假设数据已按 user_id 排序)
df["is_first_order"] = ~df["user_id"].duplicated()
return df
四、增量抽取方案设计
全量抽取在小数据量时没有问题,但当表达到亿级时,每次全量扫描就是灾难。增量抽取才是生产环境的标配。
4.1 基于时间戳(Timestamp-based)
最常见的方式,在源表上存在 updated_at 或 `modified_at` 字段,每次抽取筛选增量时间窗口:
def incremental_extract_timestamp(
conn, table: str, last_run: datetime, batch_size: int = 10000
) -> pd.DataFrame:
"""基于时间戳的增量抽取"""
query = f"""
SELECT * FROM {table}
WHERE updated_at >= %(last_run)s
ORDER BY updated_at ASC
LIMIT %(batch_size)s
"""
df = pd.read_sql(query, conn, params={
"last_run": last_run,
"batch_size": batch_size
})
return df
优点:实现简单,几乎零侵入。
缺点:严重依赖业务代码正确维护时间字段;无法捕获物理删除记录。
4.2 基于 CDC(Change Data Capture)
使用数据库原生变更日志机制:
- **MySQL**: 解析 Binlog(通过 Canal / Maxwell / Debezium)
- **PostgreSQL**: 逻辑复制(Logical Replication / pgoutput)
- **SQL Server**: Change Tracking / CDC 表
# 以 Debezium + Kafka 为例的伪代码示意
def consume_cdc_events(topic: str, bootstrap_servers: str):
"""消费 CDC 事件流"""
consumer = KafkaConsumer(
topic,
bootstrap_servers=bootstrap_servers,
value_deserializer=lambda m: json.loads(m.decode("utf-8")),
auto_offset_reset="earliest",
enable_auto_commit=False,
)
for message in consumer:
event = message.value
if event["op"] == "c": # Create
upsert_to_target(event["after"])
elif event["op"] == "u": # Update
upsert_to_target(event["after"])
elif event["op"] == "d": # Delete
soft_delete_in_target(event["before"]["id"])
consumer.commit()
优点:实时性强,能捕获删除,对业务表无侵入。
缺点:架构复杂度大幅增加,需要维护 Kafka / 数据同步组件。
4.3 基于日志(Log-based)
适用于没有时间戳也无法开启 CDC 的遗留系统。通过业务日志或 API 的变更推送来捕获数据:
def incremental_from_api_logs(log_file: str, parser: callable) -> Generator:
"""解析业务操作日志,还原数据变更"""
with open(log_file, "r") as f:
for line in tail(f): # tail -f 类似效果
event = parser(line)
if event["action"] in ("create", "update"):
yield fetch_full_record(event["record_id"])
4.4 三种方案对比
| 方案 | 实时性 | 捕获删除 | 侵入性 | 维护成本 |
|------|--------|----------|--------|----------|
| 时间戳 | 分钟级 | ❌ | 低 | 低 |
| CDC | 秒级 | ✅ | 无 | 高 |
| 日志 | 分钟级 | 取决于日志 | 中 | 中 |
**建议**:初创团队用时间戳方案快速上线;团队成熟后逐步引入 CDC 覆盖核心链路。
在实际项目中,我通常采用「混合策略」:核心交易表用 CDC 保证实时性和数据完整性,一般业务表用时间戳降低架构复杂度,遗留系统则通过日志补全。没有银弹,只有最适合当前阶段的方案。
五、完整管道实战:Python 数据清洗管道
下面我们用 Python 实现一个完整的 ETL 清洗管道,包含数据校验、异常处理、日志记录和断点恢复。
#!/usr/bin/env python3
"""
etl_pipeline.py —— 通用 ETL 数据清洗管道
支持:多层校验、异常隔离、日志链路追踪、断点续传
"""
import hashlib
import json
import logging
import sys
import time
from datetime import datetime, timedelta
from pathlib import Path
from typing import Optional
import pandas as pd
import yaml
# ── 日志配置 ────────────────────────────────────────────
def setup_logger(name: str, log_dir: str = "logs") -> logging.Logger:
Path(log_dir).mkdir(parents=True, exist_ok=True)
logger = logging.getLogger(name)
logger.setLevel(logging.DEBUG)
fmt = logging.Formatter(
"[%(asctime)s] %(levelname)-8s [%(name)s] %(message)s",
datefmt="%Y-%m-%d %H:%M:%S"
)
# 文件日志(全量)
fh = logging.FileHandler(f"{log_dir}/{name}_{datetime.now():%Y%m%d}.log")
fh.setLevel(logging.DEBUG)
fh.setFormatter(fmt)
# 控制台日志(INFO 以上)
ch = logging.StreamHandler()
ch.setLevel(logging.INFO)
ch.setFormatter(fmt)
logger.addHandler(fh)
logger.addHandler(ch)
return logger
logger = setup_logger("etl_pipeline")
# ── 数据校验器 ──────────────────────────────────────────
class DataValidator:
"""数据校验规则集合"""
@staticmethod
def check_not_null(df: pd.DataFrame, columns: list[str]) -> pd.DataFrame:
"""必填字段校验——空值所在行标记为 reject"""
mask = df[columns].isna().any(axis=1)
df.loc[mask, "_reject_reason"] = df.loc[mask, "_reject_reason"].fillna("") + \
f"missing_required:{','.join(columns)};"
return df
@staticmethod
def check_unique(df: pd.DataFrame, keys: list[str]) -> pd.DataFrame:
"""唯一性校验"""
dup_mask = df.duplicated(subset=keys, keep=False)
df.loc[dup_mask, "_reject_reason"] = df.loc[dup_mask, "_reject_reason"].fillna("") + \
f"duplicate_key:{','.join(keys)};"
return df
@staticmethod
def check_value_range(
df: pd.DataFrame, column: str, min_val: float = None, max_val: float = None
) -> pd.DataFrame:
"""值域校验"""
if min_val is not None:
mask = df[column] < min_val
df.loc[mask, "_reject_reason"] = df.loc[mask, "_reject_reason"].fillna("") + \
f"below_min:{column}={min_val};"
if max_val is not None:
mask = df[column] > max_val
df.loc[mask, "_reject_reason"] = df.loc[mask, "_reject_reason"].fillna("") + \
f"above_max:{column}={max_val};"
return df
# ── 清洗管道主类 ────────────────────────────────────────
class ETLPipeline:
"""
可配置的 ETL 数据清洗管道
用法:
pipeline = ETLPipeline(config_path="pipeline_config.yaml")
pipeline.run("orders_20240617.csv")
"""
def __init__(self, config_path: str):
with open(config_path, "r") as f:
self.config = yaml.safe_load(f)
self.logger = logging.getLogger(f"etl_pipeline.{self.config.get('name', 'default')}")
self._stats = {"total": 0, "passed": 0, "rejected": 0, "duration": 0.0}
def _load_data(self, path: str, **kwargs) -> pd.DataFrame:
"""多格式数据加载"""
ext = Path(path).suffix.lower()
loaders = {
".csv": pd.read_csv,
".json": pd.read_json,
".parquet": pd.read_parquet,
".xlsx": pd.read_excel,
}
loader = loaders.get(ext)
if not loader:
raise ValueError(f"Unsupported file format: {ext}")
self.logger.info(f"Loading data from {path}")
df = loader(path, **kwargs)
df["_reject_reason"] = ""
self._stats["total"] = len(df)
self.logger.info(f"Loaded {len(df)} records")
return df
def _validate(self, df: pd.DataFrame) -> pd.DataFrame:
"""执行校验规则"""
rules = self.config.get("validation", {})
validator = DataValidator()
for rule in rules.get("not_null", []):
df = validator.check_not_null(df, rule["columns"])
self.logger.debug(f"Not-null check: {rule['columns']}")
for rule in rules.get("unique", []):
df = validator.check_unique(df, rule["keys"])
self.logger.debug(f"Unique check: {rule['keys']}")
for rule in rules.get("range", []):
df = validator.check_value_range(
df, rule["column"], rule.get("min"), rule.get("max")
)
self.logger.debug(f"Range check: {rule['column']}")
return df
def _transform(self, df: pd.DataFrame) -> pd.DataFrame:
"""执行转换逻辑"""
transforms = self.config.get("transform", {})
# 类型转换
for col, dtype in transforms.get("type_cast", {}).items():
if col in df.columns:
try:
df[col] = df[col].astype(dtype)
except Exception as e:
self.logger.warning(f"Type cast failed for {col}: {e}")
# 字段映射
rename_map = transforms.get("field_mapping", {})
if rename_map:
df = df.rename(columns=rename_map)
# 值映射
for col, mapping in transforms.get("value_mapping", {}).items():
if col in df.columns:
df[col] = df[col].map(mapping).fillna(df[col])
self.logger.info("Transform completed")
return df
def _split_rejects(self, df: pd.DataFrame) -> tuple[pd.DataFrame, pd.DataFrame]:
"""将有拒绝原因的记录分离出来"""
reject_mask = df["_reject_reason"].str.len() > 0
rejects = df[reject_mask].copy()
passed = df[~reject_mask].copy()
passed = passed.drop(columns=["_reject_reason"])
rejects = rejects.drop(columns=["_reject_reason"])
self._stats["rejected"] = len(rejects)
self._stats["passed"] = len(passed)
return passed, rejects
def run(self, source: str, **kwargs) -> dict:
"""执行完整的 ETL 管道"""
start_time = time.time()
self.logger.info(f"Pipeline started — source={source}")
try:
# 1. 数据加载
df = self._load_data(source, **kwargs)
# 2. 数据校验
df = self._validate(df)
# 3. 数据转换
df = self._transform(df)
# 4. 分离合法/拒绝数据
passed, rejects = self._split_rejects(df)
except Exception as e:
self.logger.exception(f"Pipeline failed: {e}")
raise
# 5. 输出
output = self.config.get("output", {})
if not passed.empty:
passed_path = output.get("path", "output") + "/passed"
Path(passed_path).mkdir(parents=True, exist_ok=True)
passed.to_parquet(f"{passed_path}/{Path(source).stem}.parquet", index=False)
self.logger.info(f"Passed records written: {len(passed)}")
if not rejects.empty:
reject_path = output.get("reject_path", "output") + "/rejects"
Path(reject_path).mkdir(parents=True, exist_ok=True)
rejects.to_csv(f"{reject_path}/{Path(source).stem}_rejects.csv", index=False)
self.logger.warning(f"Rejected records written: {len(rejects)}")
self._stats["duration"] = round(time.time() - start_time, 2)
self.logger.info(
f"Pipeline finished — "
f"total={self._stats['total']}, "
f"passed={self._stats['passed']}, "
f"rejected={self._stats['rejected']}, "
f"duration={self._stats['duration']}s"
)
return self._stats
配套配置文件示例
# pipeline_config.yaml
name: order_etl_pipeline
validation:
not_null:
- columns: ["order_id", "user_id", "amount"]
unique:
- keys: ["order_id"]
range:
- column: "amount"
min: 0
max: 1000000
transform:
type_cast:
order_id: int
amount: float
created_at: datetime64[ns]
field_mapping:
order_id: order_id
user_id: user_id
amount: order_amount
created_at: order_date
value_mapping:
status:
"1": active
"2": inactive
output:
path: ./output/data
reject_path: ./output/rejects
使用示例
# 生产调用
pipeline = ETLPipeline("pipeline_config.yaml")
stats = pipeline.run("raw_orders_20240617.csv")
print(stats)
# {'total': 50000, 'passed': 49832, 'rejected': 168, 'duration': 3.42}
这个管道设计体现了几个关键原则:
1. 配置与代码分离——清洗规则写在 YAML 里,不用改代码就能调整规则
2. 异常隔离——脏数据不会污染目标表,单独归档供分析
3. 可观测性——完整的日志链路,每条记录都能追溯到校验失败的原因
4. 多格式输入输出——支持 CSV / JSON / Parquet / Excel
六、数据质量监控与告警机制
写好了管道只是第一步,更关键的是如何持续保证数据质量。没有监控的 ETL 就像没有仪表盘的飞机——你只知道它在飞,但不知道方向对不对,也不知道有没有零件在掉。
我建议建立三层监控体系,从宏观到微观逐层覆盖:
6.1 表级监控(行数 + 波动率)
每次 ETL 完成后记录核心指标到监控表:
-- 数据质量监控表
CREATE TABLE dw_quality_metrics (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
table_name VARCHAR(128),
batch_id VARCHAR(64),
run_date DATE,
row_count INT,
null_count INT,
duplicate_count INT,
reject_count INT,
duration_seconds DECIMAL(10,2),
status VARCHAR(20), -- SUCCESS / WARNING / FAILED
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
每次跑完 ETL,对比历史数据的行数波动,超出阈值(如 ±20%)则告警:
def check_row_count_anomaly(
table: str, current_count: int, history_days: int = 30, threshold: float = 0.2
) -> bool:
"""检测行数是否异常"""
query = """
SELECT AVG(row_count) as avg_rows, STDDEV(row_count) as std_rows
FROM dw_quality_metrics
WHERE table_name = %(table)s
AND run_date >= NOW() - INTERVAL %(days)s DAY
AND status = 'SUCCESS'
"""
# 执行查询…
avg_rows, std_rows = 485000, 12000
lower = avg_rows * (1 - threshold)
upper = avg_rows * (1 + threshold)
if current_count < lower or current_count > upper:
logger.error(f"Row count anomaly for {table}: {current_count} (expected ~{avg_rows})")
return False
return True
6.2 字段级监控(空值率 + 枚举分布)
表级监控能发现「行数不对」,但无法回答「数据质量有没有下降」。字段级监控就是用来回答这个问题的。
对关键字段记录空值率,出现突变时自动告警。比如,昨天 phone 字段的空值率是 2%,今天突然变成 15%,说明上游接口可能出了问题。同样地,对于枚举字段(如 status、gender),如果某个枚举值的占比出现剧烈波动,也值得关注。这些突变往往是数据源侧逻辑变更的前兆信号。
def profile_columns(df: pd.DataFrame) -> dict:
"""生成字段画像"""
profile = {}
for col in df.columns:
null_pct = df[col].isna().mean()
profile[col] = {
"null_rate": round(null_pct, 4),
"dtype": str(df[col].dtype),
"unique_count": df[col].nunique(),
}
if null_pct > 0.05: # 空值率超过 5% 告警
logger.warning(f"Column {col} null rate: {null_pct:.2%}")
return profile
6.3 调度级告警
通过 Prometheus + AlertManager 或自建告警系统,对以下场景发送通知(企业微信 / 钉钉 / Slack):
# 企业微信机器人告警示例
def send_wechat_alert(webhook: str, title: str, content: str):
payload = {
"msgtype": "markdown",
"markdown": {
"content": f"## ⚠️ ETL 告警\n**{title}**\n{content}\n时间: {datetime.now()}"
}
}
requests.post(webhook, json=payload, timeout=5)
七、调度监控基础:ETL 任务状态管理
最后,简单聊聊 ETL 任务调度侧的基础监控。数据管道写好了,质量监控也配上了,但如果任务本身挂了,一切归零。调度监控要解决的核心问题是:任务有没有按时跑?有没有跑成功?花了多长时间?
无论你用 Airflow、DolphinScheduler 还是自建调度,核心指标都是这四类:
在实际生产中,我见过很多团队只关注任务成功/失败状态,忽略了时长和记录数的监控。结果是:任务每天都能跑成功,但处理时间从 10 分钟悄悄膨胀到了 2 小时,直到某一天在业务高峰期彻底超时崩溃。所以趋势监控比状态监控更重要。
| 指标 | 含义 | 告警阈值 |
|------|------|---------|
| etl_task_duration_seconds | 任务执行耗时 | P99 > 历史均值 2x |
| etl_task_status | 最终状态(0=成功/1=失败) | > 0 |
| etl_records_processed | 处理记录数 | 低于历史均值 50% |
| etl_records_rejected | 拒绝记录数 | 超过 1% 总量 |
很多团队只在任务失败时收到告警,但更隐蔽的问题是「任务成功了但数据是错的」。行数骤降 80% 但任务依然显示成功,这种情况我遇到过不止一次。所以我们一定要把 records_processed 也纳入监控范围,与历史均值做对比。
如果团队资源允许,建议将调度监控与质量监控打通。比如在 Airflow 的 DAG 中,最后一个 Task 不是「数据加载完成」,而是「质量校验通过」。只有质量校验通过了,DAG 才标记为成功。这样可以在一个视图中统览所有 ETL 管道的健康状态。
# 使用 prometheus_client 暴露指标
from prometheus_client import Counter, Histogram, Gauge, start_http_server
etl_duration = Histogram(
"etl_task_duration_seconds", "ETL task duration",
["task_name", "source"],
buckets=(1, 5, 10, 30, 60, 120, 300, 600)
)
etl_status = Gauge(
"etl_task_status", "ETL task status (0=success, 1=fail)",
["task_name"]
)
etl_records_rejected = Counter(
"etl_records_rejected_total", "Total rejected records",
["task_name", "reason"]
)
def monitor_task(task_name: str, source: str):
"""装饰器:自动采集 ETL 监控指标"""
def decorator(func):
@functools.wraps(func)
def wrapper(*args, **kwargs):
start = time.time()
try:
result = func(*args, **kwargs)
etl_status.labels(task_name=task_name).set(0)
return result
except Exception as e:
etl_status.labels(task_name=task_name).set(1)
raise
finally:
etl_duration.labels(
task_name=task_name, source=source
).observe(time.time() - start)
return wrapper
return decorator
将这些指标接入 Grafana,你就可以看到这样的面板:
- 今日 ETL 成功率(红线 = 异常)
- 各任务执行时长趋势
- 被拒绝数据的 TOP10 原因
- 数据行数日环比波动
结语
数据清洗从来不是一次性的工作。随着业务演进,新的数据源不断接入、老的业务逻辑持续变更,「脏数据」的形式也在进化。与其每次被动修数据,不如从一开始就建立体系化的清洗框架和监控机制。
回顾一下本文的核心要点:
1. 数据质量问题分类是清洗的第一步——只有定义了什么是「脏」,才能知道怎么「洗」
2. 清洗框架设计要遵循单一职责原则,每个 Cleaner 只做一件事,管道可以灵活组装
3. 数据转换的本质是语义映射——从源端业务语义到目标端标准语义的桥梁
4. 增量抽取选择要看团队阶段——初创用时间戳,成熟上 CDC
5. 质量监控要覆盖表级、字段级和调度级三层,不能只盯着任务成功/失败
6. Prometheus + Grafana 可以很好地支撑 ETL 可观测性建设
本文分享的这些策略和代码,都是我过去几年在多个项目中踩坑后沉淀下来的经验。它们不一定适用于所有场景,但我希望至少能提供一个思考框架,让你在遇到类似问题时知道从哪里入手。
最后想说一点:在做数据清洗时,保持怀疑——对上游数据永远保持怀疑,对自己写的清洗逻辑也保持怀疑。多一层校验、多一条日志,往往就能在关键时刻救你一命。
如果你有更好的清洗策略或者踩过什么有趣的「脏数据」的坑,欢迎在评论区交流分享。下一篇文章,我们将讨论 ETL 性能优化——当数据量从百万级增长到亿级时,哪些设计需要重新思考。
*本文为学习分享系列文章,代码示例可在 [github.com/your-repo/etl-series](https://github.com) 找到完整项目。*
评论