{T}

数据清洗管道实战

数据清洗是数据工程中最耗时也最关键的环节。 一个设计良好的清洗管道不仅能提升数据质量,还能让数据流转过程可追溯、可复现、可维护。本文聚焦 Python 数据清洗管道的设计原则、验证策略、编排模式与实战场景。

阅读提示

数据清洗管道设计原则

ETL 管道全景

图表渲染中…

核心设计原则

原则说明实践方式
幂等性同一输入多次执行产生相同结果避免随机填充、使用确定性种子
可追溯每条数据可追溯到原始来源添加 _source_ingested_at 元数据列
可复现相同数据 + 相同代码 = 相同输出固定随机种子、版本化清洗规则
渐进式清洗步骤可独立运行和调试每步返回 DataFrame,不修改原始数据
防御性对异常输入优雅降级而非崩溃try/except + 日志记录 + 默认值
可观测管道运行状态可监控统计每步处理行数、耗时、异常数

管道基类设计

python
from __future__ import annotations

import logging
import time
from abc import ABC, abstractmethod
from dataclasses import dataclass, field

import pandas as pd

logger = logging.getLogger(__name__)


@dataclass
class PipelineStats:
    """管道运行统计"""
    step_name: str
    input_rows: int = 0
    output_rows: int = 0
    dropped_rows: int = 0
    duration_seconds: float = 0.0
    errors: list[str] = field(default_factory=list)

    @property
    def drop_rate(self) -> float:
        return self.dropped_rows / self.input_rows if self.input_rows else 0.0


class PipelineStep(ABC):
    """管道步骤基类:每个清洗步骤继承此类"""

    def __init__(self, name: str) -> None:
        self.name = name

    @abstractmethod
    def process(self, df: pd.DataFrame) -> pd.DataFrame:
        """处理数据,返回清洗后的 DataFrame"""
        ...

    def run(self, df: pd.DataFrame) -> tuple[pd.DataFrame, PipelineStats]:
        """执行步骤并收集统计信息"""
        stats = PipelineStats(step_name=self.name, input_rows=len(df))
        start = time.perf_counter()

        try:
            result = self.process(df)
            stats.output_rows = len(result)
            stats.dropped_rows = stats.input_rows - stats.output_rows
        except Exception as e:
            stats.errors.append(str(e))
            logger.error(f"步骤 [{self.name}] 执行失败: {e}")
            result = df  # 失败时返回原始数据(防御性设计)

        stats.duration_seconds = round(time.perf_counter() - start, 3)
        logger.info(
            f"步骤 [{self.name}]: {stats.input_rows} → {stats.output_rows} 行 "
            f"(丢弃 {stats.dropped_rows}, 耗时 {stats.duration_seconds}s)"
        )
        return result, stats


class DataPipeline:
    """数据清洗管道:组合多个步骤顺序执行"""

    def __init__(self, name: str) -> None:
        self.name = name
        self.steps: list[PipelineStep] = []
        self.stats: list[PipelineStats] = []

    def add_step(self, step: PipelineStep) -> DataPipeline:
        """添加步骤(支持链式调用)"""
        self.steps.append(step)
        return self

    def run(self, df: pd.DataFrame) -> pd.DataFrame:
        """顺序执行所有步骤"""
        logger.info(f"管道 [{self.name}] 启动,输入 {len(df)} 行")
        self.stats = []

        for step in self.steps:
            df, step_stats = step.run(df)
            self.stats.append(step_stats)

        logger.info(f"管道 [{self.name}] 完成,输出 {len(df)} 行")
        return df

    def report(self) -> pd.DataFrame:
        """生成管道运行报告"""
        return pd.DataFrame([
            {
                "步骤": s.step_name,
                "输入行数": s.input_rows,
                "输出行数": s.output_rows,
                "丢弃行数": s.dropped_rows,
                "丢弃率": f"{s.drop_rate:.2%}",
                "耗时(秒)": s.duration_seconds,
                "错误": "; ".join(s.errors) if s.errors else "",
            }
            for s in self.stats
        ])
管道基类设计要点
  • 每个步骤继承 PipelineStep,只实现 process() 方法
  • run() 方法自动收集统计信息,无需手动埋点
  • 失败时返回原始数据而非抛异常,保证管道不中断
  • DataPipeline.report() 可生成完整的运行报告

数据验证策略

数据验证是清洗管道的第一道防线。在数据进入转换流程之前,先验证其结构和内容是否符合预期。

验证层次

图表渲染中…

Schema 验证

Schema 验证确保数据包含预期的列名、数据类型和约束条件。

python
from dataclasses import dataclass
from typing import Any

import pandas as pd


@dataclass
class ColumnSchema:
    """列级 Schema 定义"""
    name: str
    dtype: str                # 期望的数据类型
    nullable: bool = True     # 是否允许空值
    unique: bool = False      # 是否要求唯一
    min_value: float | None = None
    max_value: float | None = None
    allowed_values: list[str] | None = None
    regex_pattern: str | None = None


class SchemaValidator:
    """Schema 验证器"""

    def __init__(self, schema: list[ColumnSchema]) -> None:
        self.schema = schema
        self.errors: list[dict[str, Any]] = []

    def validate(self, df: pd.DataFrame) -> tuple[pd.DataFrame, list[dict[str, Any]]]:
        """验证 DataFrame 是否符合 Schema"""
        self.errors = []
        schema_dict = {s.name: s for s in self.schema}

        # 1. 检查必需列是否存在
        missing_cols = set(schema_dict.keys()) - set(df.columns)
        for col in missing_cols:
            self.errors.append({
                "列": col, "规则": "列存在性", "问题": f"缺少必需列: {col}",
            })

        # 2. 逐列验证
        for col_name, col_schema in schema_dict.items():
            if col_name not in df.columns:
                continue
            series = df[col_name]

            # 类型检查
            if not self._check_dtype(series, col_schema.dtype):
                self.errors.append({
                    "列": col_name, "规则": "数据类型",
                    "问题": f"期望 {col_schema.dtype},实际 {series.dtype}",
                })

            # 空值检查
            if not col_schema.nullable and series.isnull().any():
                null_count = series.isnull().sum()
                self.errors.append({
                    "列": col_name, "规则": "非空约束",
                    "问题": f"发现 {null_count} 个空值",
                })

            # 唯一性检查
            if col_schema.unique and series.duplicated().any():
                dup_count = series.duplicated().sum()
                self.errors.append({
                    "列": col_name, "规则": "唯一性",
                    "问题": f"发现 {dup_count} 个重复值",
                })

            # 范围检查
            if col_schema.min_value is not None:
                numeric = pd.to_numeric(series, errors='coerce')
                below = (numeric < col_schema.min_value).sum()
                if below > 0:
                    self.errors.append({
                        "列": col_name, "规则": "最小值",
                        "问题": f"{below} 个值低于最小值 {col_schema.min_value}",
                    })

            if col_schema.max_value is not None:
                numeric = pd.to_numeric(series, errors='coerce')
                above = (numeric > col_schema.max_value).sum()
                if above > 0:
                    self.errors.append({
                        "列": col_name, "规则": "最大值",
                        "问题": f"{above} 个值超过最大值 {col_schema.max_value}",
                    })

            # 枚举值检查
            if col_schema.allowed_values is not None:
                invalid = ~series.isin(col_schema.allowed_values) & series.notna()
                if invalid.any():
                    bad_values = series[invalid].unique().tolist()[:5]
                    self.errors.append({
                        "列": col_name, "规则": "枚举值",
                        "问题": f"非法值: {bad_values}",
                    })

        return df, self.errors

    def _check_dtype(self, series: pd.Series, expected: str) -> bool:
        """检查数据类型是否匹配"""
        type_map = {
            "int": ["int64", "int32", "Int64", "Int32"],
            "float": ["float64", "float32"],
            "str": ["object", "string"],
            "datetime": ["datetime64[ns]", "datetime64[ns, UTC]"],
            "bool": ["bool", "boolean"],
        }
        allowed = type_map.get(expected, [expected])
        return str(series.dtype) in allowed


# 定义 Schema
user_schema = [
    ColumnSchema("user_id", dtype="int", nullable=False, unique=True),
    ColumnSchema("username", dtype="str", nullable=False, regex_pattern=r"^[a-zA-Z0-9_]{3,20}$"),
    ColumnSchema("age", dtype="int", nullable=True, min_value=0, max_value=150),
    ColumnSchema("email", dtype="str", nullable=False),
    ColumnSchema("status", dtype="str", nullable=False, allowed_values=["active", "inactive", "suspended"]),
    ColumnSchema("created_at", dtype="datetime", nullable=False),
]

# 使用
validator = SchemaValidator(user_schema)
df, errors = validator.validate(df)
if errors:
    for err in errors:
        print(f"  验证失败 [{err['列']}] {err['规则']}: {err['问题']}")

使用 Pandera 进行声明式验证

Pandera 是一个专门的数据验证库,支持声明式 Schema 定义和自动类型强制转换。

python
import pandera as pa
from pandera import Column, Check, Index

# 定义 Schema
sales_schema = pa.DataFrameSchema(
    columns={
        "order_id": Column(int, Check(lambda x: x > 0, error="订单ID必须为正整数"), nullable=False),
        "product": Column(str, Check.isin(["A", "B", "C", "D"]), nullable=False),
        "quantity": Column(int, Check.in_range(1, 1000), nullable=False),
        "price": Column(float, Check.greater_than(0), nullable=False),
        "order_date": Column("datetime64[ns]", nullable=False),
        "discount": Column(float, Check.in_range(0.0, 1.0), nullable=True),
    },
    checks=[
        # 行级验证:折扣价不能为负
        Check(
            lambda df: (df["price"] * (1 - df["discount"].fillna(0))) >= 0,
            error="折扣后价格不能为负",
        ),
    ],
    strict=True,       # 禁止出现未定义的列
    coerce=True,       # 自动类型转换
)

# 验证
validated_df = sales_schema.validate(df)

验证策略对比

策略工具优点缺点适用场景
手动验证自定义函数完全灵活代码量大、易遗漏简单场景、一次性验证
Schema 验证自定义 Schema 类可复用、可扩展需要维护 Schema 定义中等复杂度项目
Panderapandera 库声明式、自动类型转换、统计信息学习成本、依赖第三方库正式项目、团队协作
Great Expectationsgreat_expectations丰富的验证规则、数据文档、集成广泛重量级、配置复杂企业级数据平台
Pydanticpydantic与 FastAPI 集成好、类型安全逐行验证性能差API 数据验证

常见数据质量问题与修复

问题全景

图表渲染中…

缺失值处理

python
import pandas as pd
import numpy as np


class MissingValueHandler:
    """缺失值处理器"""

    @staticmethod
    def diagnose(df: pd.DataFrame) -> pd.DataFrame:
        """诊断缺失值分布"""
        total = len(df)
        missing = df.isnull().sum()
        report = pd.DataFrame({
            "列名": df.columns,
            "缺失数": missing.values,
            "缺失率": (missing.values / total * 100).round(2),
            "非空数": total - missing.values,
            "数据类型": df.dtypes.values,
        })
        return report[report["缺失数"] > 0].sort_values("缺失率", ascending=False)

    @staticmethod
    def fill_numeric(df: pd.DataFrame, col: str, strategy: str = "median") -> pd.DataFrame:
        """数值列缺失值填充"""
        df = df.copy()
        series = df[col]

        strategies = {
            "mean": series.mean(),
            "median": series.median(),
            "mode": series.mode()[0],
            "zero": 0,
            "min": series.min(),
            "max": series.max(),
        }

        if strategy not in strategies:
            raise ValueError(f"未知策略: {strategy},可选: {list(strategies.keys())}")

        fill_value = strategies[strategy]
        missing_count = series.isnull().sum()
        df[col] = series.fillna(fill_value)
        print(f"  [{col}] 填充 {missing_count} 个缺失值,策略: {strategy},填充值: {fill_value:.4f}")
        return df

    @staticmethod
    def fill_categorical(df: pd.DataFrame, col: str, fill_value: str = "未知") -> pd.DataFrame:
        """分类列缺失值填充"""
        df = df.copy()
        missing_count = df[col].isnull().sum()
        df[col] = df[col].fillna(fill_value)
        print(f"  [{col}] 填充 {missing_count} 个缺失值,填充值: '{fill_value}'")
        return df

    @staticmethod
    def fill_with_indicator(df: pd.DataFrame, col: str, fill_value: float = 0) -> pd.DataFrame:
        """填充并添加缺失指示列(保留缺失信息)"""
        df = df.copy()
        df[f"_is_{col}_missing"] = df[col].isnull().astype(int)
        df[col] = df[col].fillna(fill_value)
        return df

    @staticmethod
    def interpolate_time(df: pd.DataFrame, col: str, method: str = "linear") -> pd.DataFrame:
        """时间序列插值填充"""
        df = df.copy()
        missing_count = df[col].isnull().sum()
        df[col] = df[col].interpolate(method=method)
        print(f"  [{col}] 插值填充 {missing_count} 个缺失值,方法: {method}")
        return df

    @staticmethod
    def drop_by_threshold(df: pd.DataFrame, threshold: float = 0.5) -> pd.DataFrame:
        """删除缺失率超过阈值的列"""
        missing_rate = df.isnull().sum() / len(df)
        drop_cols = missing_rate[missing_rate > threshold].index.tolist()
        if drop_cols:
            print(f"  删除缺失率 > {threshold:.0%} 的列: {drop_cols}")
            df = df.drop(columns=drop_cols)
        return df
缺失值填充策略选择
  • 删除法:缺失率 < 5% 且非关键列时可用
  • 均值/中位数:数值列,分布偏斜时选中位数
  • 众数:分类列
  • 插值法:时间序列数据
  • 指示列法:缺失本身可能有业务含义(如用户未填写某字段暗示某种行为)
  • 不要统一用 0 填充,0 在不同列有不同含义

异常值处理

python
import pandas as pd
import numpy as np


class OutlierHandler:
    """异常值处理器"""

    @staticmethod
    def detect_iqr(df: pd.DataFrame, col: str, factor: float = 1.5) -> pd.Series:
        """IQR 法检测异常值"""
        q1 = df[col].quantile(0.25)
        q3 = df[col].quantile(0.75)
        iqr = q3 - q1
        lower = q1 - factor * iqr
        upper = q3 + factor * iqr
        return (df[col] < lower) | (df[col] > upper)

    @staticmethod
    def detect_zscore(df: pd.DataFrame, col: str, threshold: float = 3.0) -> pd.Series:
        """Z-Score 法检测异常值"""
        z_scores = np.abs((df[col] - df[col].mean()) / df[col].std())
        return z_scores > threshold

    @staticmethod
    def clip_iqr(df: pd.DataFrame, col: str, factor: float = 1.5) -> pd.DataFrame:
        """IQR 截断法:将异常值截断到边界"""
        df = df.copy()
        q1 = df[col].quantile(0.25)
        q3 = df[col].quantile(0.75)
        iqr = q3 - q1
        lower = q1 - factor * iqr
        upper = q3 + factor * iqr
        outlier_count = ((df[col] < lower) | (df[col] > upper)).sum()
        df[col] = df[col].clip(lower, upper)
        print(f"  [{col}] IQR 截断: {outlier_count} 个异常值,范围 [{lower:.2f}, {upper:.2f}]")
        return df

    @staticmethod
    def mark_outliers(df: pd.DataFrame, col: str, factor: float = 1.5) -> pd.DataFrame:
        """标记异常值(不修改原值,添加标记列)"""
        df = df.copy()
        is_outlier = OutlierHandler.detect_iqr(df, col, factor)
        df[f"_is_{col}_outlier"] = is_outlier.astype(int)
        print(f"  [{col}] 标记 {is_outlier.sum()} 个异常值")
        return df

    @staticmethod
    def replace_with_percentile(
        df: pd.DataFrame, col: str, lower_pct: float = 0.01, upper_pct: float = 0.99
    ) -> pd.DataFrame:
        """百分位替换法"""
        df = df.copy()
        lower = df[col].quantile(lower_pct)
        upper = df[col].quantile(upper_pct)
        outlier_count = ((df[col] < lower) | (df[col] > upper)).sum()
        df[col] = df[col].clip(lower, upper)
        print(f"  [{col}] 百分位替换: {outlier_count} 个异常值,范围 [{lower:.2f}, {upper:.2f}]")
        return df

异常值检测方法对比

方法原理优点缺点适用分布
IQR 法Q1-1.5IQR ~ Q3+1.5IQR稳健、不受极端值影响对非对称分布效果差近似正态分布
Z-Score偏离均值超过 3 个标准差直观、标准化对极端值敏感(均值被拉偏)正态分布
百分位法1% ~ 99% 分位数简单、可控可能误删正常极值任意分布
业务规则领域知识定义边界最准确需要领域专家特定业务场景
孤立森林无监督异常检测多维异常检测需要训练、参数调优高维数据

重复值处理

python
import pandas as pd


class DuplicateHandler:
    """重复值处理器"""

    @staticmethod
    def diagnose(df: pd.DataFrame, subset: list[str] | None = None) -> dict:
        """诊断重复值"""
        total = len(df)
        exact_dupes = df.duplicated().sum()
        key_dupes = 0
        if subset:
            key_dupes = df.duplicated(subset=subset).sum()

        return {
            "总行数": total,
            "完全重复行": exact_dupes,
            "完全重复率": f"{exact_dupes / total:.2%}",
            "业务键重复行": key_dupes,
            "业务键重复率": f"{key_dupes / total:.2%}" if subset else "N/A",
        }

    @staticmethod
    def remove_exact(df: pd.DataFrame, keep: str = "first") -> pd.DataFrame:
        """删除完全重复行"""
        before = len(df)
        df = df.drop_duplicates(keep=keep).reset_index(drop=True)
        print(f"  删除完全重复: {before} → {len(df)} 行 (移除 {before - len(df)} 行)")
        return df

    @staticmethod
    def remove_by_key(
        df: pd.DataFrame, key_cols: list[str], keep: str = "last"
    ) -> pd.DataFrame:
        """按业务键去重(保留最新或最早记录)"""
        before = len(df)
        df = df.drop_duplicates(subset=key_cols, keep=keep).reset_index(drop=True)
        print(f"  按键 {key_cols} 去重: {before} → {len(df)} 行")
        return df

    @staticmethod
    def merge_duplicates(
        df: pd.DataFrame, key_cols: list[str], agg_strategy: dict[str, str]
    ) -> pd.DataFrame:
        """合并重复记录(聚合冲突字段)"""
        before = len(df)
        # key_cols 以外的列按策略聚合
        result = df.groupby(key_cols, as_index=False).agg(agg_strategy)
        print(f"  合并重复: {before} → {len(result)} 行")
        return result


# 使用示例:合并重复用户记录
# agg_strategy = {"email": "first", "age": "max", "updated_at": "max", "score": "mean"}
# df = DuplicateHandler.merge_duplicates(df, key_cols=["user_id"], agg_strategy=agg_strategy)

格式不一致处理

python
import pandas as pd
import re


class FormatNormalizer:
    """格式标准化处理器"""

    @staticmethod
    def normalize_date(df: pd.DataFrame, col: str, target_format: str = "%Y-%m-%d") -> pd.DataFrame:
        """日期格式标准化"""
        df = df.copy()
        df[col] = pd.to_datetime(df[col], errors="coerce", infer_datetime_format=True)
        invalid = df[col].isnull()
        if invalid.any():
            print(f"  [{col}] {invalid.sum()} 个日期无法解析,设为 NaT")
        return df

    @staticmethod
    def normalize_phone(df: pd.DataFrame, col: str) -> pd.DataFrame:
        """电话号码标准化(保留数字,添加国家代码)"""
        df = df.copy()
        original = df[col].copy()

        def clean_phone(phone: str) -> str:
            if pd.isna(phone):
                return phone
            digits = re.sub(r"[^\d+]", "", str(phone))
            # 中国手机号:11位,添加 +86
            if re.match(r"^1[3-9]\d{9}$", digits):
                return f"+86{digits}"
            # 已有国家代码
            if digits.startswith("+"):
                return digits
            return digits

        df[col] = df[col].apply(clean_phone)
        changed = (df[col] != original).sum()
        print(f"  [{col}] 标准化 {changed} 个电话号码")
        return df

    @staticmethod
    def normalize_string(
        df: pd.DataFrame, col: str, strip: bool = True, lower: bool = False,
        remove_extra_spaces: bool = True,
    ) -> pd.DataFrame:
        """字符串标准化"""
        df = df.copy()
        series = df[col].astype(str)

        if strip:
            series = series.str.strip()
        if lower:
            series = series.str.lower()
        if remove_extra_spaces:
            series = series.str.replace(r"\s+", " ", regex=True)

        df[col] = series.replace("nan", pd.NA)
        return df

    @staticmethod
    def normalize_column_names(df: pd.DataFrame) -> pd.DataFrame:
        """列名标准化:小写 + 下划线"""
        df = df.copy()
        original = df.columns.tolist()
        df.columns = (
            df.columns.str.strip()
            .str.lower()
            .str.replace(r"[^\w]", "_", regex=True)
            .str.replace(r"_+", "_", regex=True)
            .str.strip("_")
        )
        changed = [f"{o} → {n}" for o, n in zip(original, df.columns) if o != n]
        if changed:
            print(f"  列名标准化: {changed[:5]}{'...' if len(changed) > 5 else ''}")
        return df

编码问题处理

python
import chardet
import pandas as pd


class EncodingHandler:
    """编码问题处理器"""

    @staticmethod
    def detect_encoding(file_path: str, sample_size: int = 10000) -> dict:
        """检测文件编码"""
        with open(file_path, "rb") as f:
            raw = f.read(sample_size)
        result = chardet.detect(raw)
        return {
            "编码": result["encoding"],
            "置信度": f"{result['confidence']:.2%}",
            "语言": result.get("language", "未知"),
        }

    @staticmethod
    def read_with_fallback(file_path: str, encodings: list[str] | None = None) -> pd.DataFrame:
        """尝试多种编码读取文件"""
        if encodings is None:
            encodings = ["utf-8", "gbk", "gb2312", "gb18030", "latin-1"]

        for encoding in encodings:
            try:
                df = pd.read_csv(file_path, encoding=encoding)
                print(f"  成功读取,编码: {encoding}")
                return df
            except (UnicodeDecodeError, UnicodeError):
                continue

        # 最后尝试忽略错误
        df = pd.read_csv(file_path, encoding="utf-8", errors="replace")
        print("  警告: 使用 errors='replace' 读取,部分字符可能丢失")
        return df

    @staticmethod
    def remove_bom(file_path: str, output_path: str | None = None) -> str:
        """移除 BOM 头"""
        output_path = output_path or file_path
        with open(file_path, "rb") as f:
            content = f.read()

        # UTF-8 BOM
        if content.startswith(b"\xef\xbb\xbf"):
            content = content[3:]
            print("  移除 UTF-8 BOM")
        # UTF-16 LE BOM
        elif content.startswith(b"\xff\xfe"):
            content = content[2:]
            print("  移除 UTF-16 LE BOM")

        with open(output_path, "wb") as f:
            f.write(content)
        return output_path

增量更新策略

全量 vs 增量

图表渲染中…
维度全量更新增量更新
数据量全部数据仅新增/变更部分
执行时间长(与数据总量成正比)短(与变更量成正比)
资源消耗
实现复杂度简单较复杂(需识别变更)
数据一致性强(每次完整重建)需额外保证
适用场景数据量小、变更频繁数据量大、变更较少

Watermark 策略

Watermark 通过时间戳标记已处理数据的位置,下次只处理新数据。

python
import os
import json
from datetime import datetime
from pathlib import Path

import pandas as pd


class WatermarkManager:
    """Watermark 管理器:记录和读取数据处理的进度标记"""

    def __init__(self, watermark_dir: str = ".watermarks") -> None:
        self.watermark_dir = Path(watermark_dir)
        self.watermark_dir.mkdir(exist_ok=True)

    def _path(self, pipeline_name: str) -> Path:
        return self.watermark_dir / f"{pipeline_name}.json"

    def get(self, pipeline_name: str) -> datetime | None:
        """获取上次处理的 watermark"""
        path = self._path(pipeline_name)
        if path.exists():
            data = json.loads(path.read_text())
            return datetime.fromisoformat(data["watermark"])
        return None

    def set(self, pipeline_name: str, watermark: datetime) -> None:
        """设置 watermark"""
        path = self._path(pipeline_name)
        path.write_text(json.dumps({
            "pipeline": pipeline_name,
            "watermark": watermark.isoformat(),
            "updated_at": datetime.now().isoformat(),
        }, ensure_ascii=False))
        print(f"  Watermark 已更新: {pipeline_name} → {watermark}")

    def incremental_read(
        self, pipeline_name: str, df: pd.DataFrame, time_col: str
    ) -> pd.DataFrame:
        """基于 watermark 过滤增量数据"""
        last_watermark = self.get(pipeline_name)

        if last_watermark is None:
            print(f"  首次运行,处理全部数据 ({len(df)} 行)")
            return df

        df[time_col] = pd.to_datetime(df[time_col])
        incremental = df[df[time_col] > last_watermark]
        print(f"  增量数据: {len(incremental)} 行 (watermark: {last_watermark})")
        return incremental


# 使用示例
wm = WatermarkManager()

# 读取增量数据
df_all = pd.read_csv("orders.csv", parse_dates=["updated_at"])
df_incremental = wm.incremental_read("orders_pipeline", df_all, "updated_at")

# 处理完成后更新 watermark
if len(df_incremental) > 0:
    new_watermark = df_incremental["updated_at"].max()
    wm.set("orders_pipeline", new_watermark)

Change Data Capture (CDC)

CDC 通过捕获数据源的变更日志来识别新增、修改和删除的记录。

python
from enum import Enum
from dataclasses import dataclass

import pandas as pd


class ChangeType(Enum):
    INSERT = "INSERT"
    UPDATE = "UPDATE"
    DELETE = "DELETE"


@dataclass
class ChangeRecord:
    change_type: ChangeType
    key: dict
    data: dict
    timestamp: str


class CDCProcessor:
    """CDC 变更处理器"""

    @staticmethod
    def detect_changes(
        source: pd.DataFrame,
        target: pd.DataFrame,
        key_cols: list[str],
        compare_cols: list[str] | None = None,
    ) -> pd.DataFrame:
        """对比源表和目标表,检测变更"""
        compare_cols = compare_cols or [c for c in source.columns if c not in key_cols]

        # 标记来源
        source = source.copy()
        target = target.copy()
        source["_source"] = "source"
        target["_source"] = "target"

        merged = pd.merge(
            source, target, on=key_cols, how="outer", suffixes=("_src", "_tgt"), indicator=True
        )

        changes = []

        for _, row in merged.iterrows():
            key = {k: row[k] for k in key_cols}

            if row["_merge"] == "right_only":
                # 目标有、源没有 → 删除
                changes.append({**key, "_change_type": ChangeType.DELETE.value})
            elif row["_merge"] == "left_only":
                # 源有、目标没有 → 新增
                changes.append({**key, "_change_type": ChangeType.INSERT.value})
            else:
                # 两边都有,检查是否有字段变更
                for col in compare_cols:
                    src_val = row.get(f"{col}_src")
                    tgt_val = row.get(f"{col}_tgt")
                    if src_val != tgt_val and not (pd.isna(src_val) and pd.isna(tgt_val)):
                        changes.append({**key, "_change_type": ChangeType.UPDATE.value})
                        break

        result = pd.DataFrame(changes) if changes else pd.DataFrame()
        if len(result) > 0:
            print(f"  CDC 检测: {len(result)} 条变更")
            print(f"    新增: {(result['_change_type'] == 'INSERT').sum()}")
            print(f"    更新: {(result['_change_type'] == 'UPDATE').sum()}")
            print(f"    删除: {(result['_change_type'] == 'DELETE').sum()}")
        else:
            print("  CDC 检测: 无变更")
        return result

    @staticmethod
    def apply_changes(
        target: pd.DataFrame, changes: pd.DataFrame, source: pd.DataFrame, key_cols: list[str]
    ) -> pd.DataFrame:
        """将变更应用到目标表"""
        if changes.empty:
            return target

        target = target.copy()

        # 处理删除
        deletes = changes[changes["_change_type"] == ChangeType.DELETE.value]
        if len(deletes) > 0:
            delete_keys = deletes[key_cols]
            mask = target[key_cols].apply(tuple, axis=1).isin(
                delete_keys.apply(tuple, axis=1)
            )
            target = target[~mask]

        # 处理更新和新增
        upserts = changes[changes["_change_type"].isin([ChangeType.INSERT.value, ChangeType.UPDATE.value])]
        if len(upserts) > 0:
            upsert_keys = upserts[key_cols]
            # 先删除目标中已有的(更新场景)
            mask = target[key_cols].apply(tuple, axis=1).isin(
                upsert_keys.apply(tuple, axis=1)
            )
            target = target[~mask]
            # 再从源表插入
            source_rows = source[
                source[key_cols].apply(tuple, axis=1).isin(
                    upsert_keys.apply(tuple, axis=1)
                )
            ]
            target = pd.concat([target, source_rows], ignore_index=True)

        return target.reset_index(drop=True)

管道编排模式

线性管道

最简单的编排模式:步骤按顺序执行,上一步的输出是下一步的输入。

图表渲染中…
python
# 线性管道
pipeline = DataPipeline("linear_clean")
pipeline.add_step(ReadCSVStep("data.csv"))
pipeline.add_step(ValidateStep(schema))
pipeline.add_step(MissingValueStep(strategy="median"))
pipeline.add_step(OutlierStep(method="iqr"))
pipeline.add_step(FormatNormalizeStep())
pipeline.add_step(ExportStep("output.parquet"))

result = pipeline.run(initial_df)

分支管道

数据需要按不同规则分别处理,最后再合并。

图表渲染中…
python
class BranchPipeline:
    """分支管道:按条件分发到不同子管道"""

    def __init__(self, name: str) -> None:
        self.name = name
        self.branches: dict[str, DataPipeline] = {}

    def add_branch(self, branch_name: str, pipeline: DataPipeline) -> BranchPipeline:
        self.branches[branch_name] = pipeline
        return self

    def run(self, df: pd.DataFrame, branch_fn) -> pd.DataFrame:
        """
        branch_fn: 接收 DataFrame,返回 {分支名: 子 DataFrame} 的字典
        """
        splits = branch_fn(df)
        results = {}

        for branch_name, sub_df in splits.items():
            if branch_name in self.branches:
                results[branch_name] = self.branches[branch_name].run(sub_df)
            else:
                results[branch_name] = sub_df

        # 合并所有分支结果
        return pd.concat(results.values(), ignore_index=True)


# 使用示例
def split_by_type(df: pd.DataFrame) -> dict[str, pd.DataFrame]:
    numeric_cols = df.select_dtypes(include=[np.number]).columns.tolist()
    other_cols = [c for c in df.columns if c not in numeric_cols]
    return {
        "numeric": df[numeric_cols],
        "other": df[other_cols],
    }

branch_pipeline = BranchPipeline("type_aware_clean")
branch_pipeline.add_branch("numeric", numeric_pipeline)
branch_pipeline.add_branch("other", categorical_pipeline)

result = branch_pipeline.run(df, split_by_type)

合并管道

多个数据源分别清洗后合并。

图表渲染中…
python
class MergePipeline:
    """合并管道:多数据源清洗后合并"""

    def __init__(self, name: str) -> None:
        self.name = name
        self.sources: dict[str, DataPipeline] = {}

    def add_source(self, source_name: str, pipeline: DataPipeline) -> MergePipeline:
        self.sources[source_name] = pipeline
        return self

    def run(
        self,
        source_data: dict[str, pd.DataFrame],
        merge_on: list[str],
        how: str = "outer",
        dedup: bool = True,
    ) -> pd.DataFrame:
        """分别清洗后合并"""
        cleaned = {}
        for name, pipeline in self.sources.items():
            if name in source_data:
                cleaned[name] = pipeline.run(source_data[name])

        # 逐步合并
        dfs = list(cleaned.values())
        if not dfs:
            return pd.DataFrame()

        result = dfs[0]
        for df in dfs[1:]:
            result = pd.merge(result, df, on=merge_on, how=how, suffixes=("", f"_{name}"))

        # 合并后去重
        if dedup:
            before = len(result)
            result = result.drop_duplicates(subset=merge_on, keep="last").reset_index(drop=True)
            print(f"  合并后去重: {before} → {len(result)} 行")

        return result

实战场景

场景一:CSV 数据清洗管道

完整的 CSV 数据清洗管道:读取 → 验证 → 清洗 → 转换 → 导出。

python
from __future__ import annotations

import logging
import time
from pathlib import Path

import numpy as np
import pandas as pd

logging.basicConfig(level=logging.INFO, format="%(message)s")
logger = logging.getLogger(__name__)


# ── 步骤 1:数据读取 ──────────────────────────────────────────

class ReadCSVStep(PipelineStep):
    """CSV 数据读取步骤"""

    def __init__(self, file_path: str, **read_kwargs) -> None:
        super().__init__(name="读取CSV")
        self.file_path = file_path
        self.read_kwargs = read_kwargs

    def process(self, df: pd.DataFrame) -> pd.DataFrame:
        path = Path(self.file_path)
        if not path.exists():
            raise FileNotFoundError(f"文件不存在: {self.file_path}")

        result = pd.read_csv(self.file_path, **self.read_kwargs)
        print(f"  读取 {path.name}: {result.shape[0]} 行 × {result.shape[1]} 列")
        return result


# ── 步骤 2:列名标准化 ────────────────────────────────────────

class NormalizeColumnsStep(PipelineStep):
    """列名标准化步骤"""

    def process(self, df: pd.DataFrame) -> pd.DataFrame:
        original = df.columns.tolist()
        df.columns = (
            df.columns.str.strip()
            .str.lower()
            .str.replace(r"[^\w]", "_", regex=True)
            .str.replace(r"_+", "_", regex=True)
            .str.strip("_")
        )
        changed = sum(1 for o, n in zip(original, df.columns) if o != n)
        print(f"  标准化 {changed} 个列名")
        return df


# ── 步骤 3:数据验证 ──────────────────────────────────────────

class ValidateStep(PipelineStep):
    """数据验证步骤"""

    def __init__(self, required_cols: list[str] | None = None) -> None:
        super().__init__(name="数据验证")
        self.required_cols = required_cols or []

    def process(self, df: pd.DataFrame) -> pd.DataFrame:
        # 检查必需列
        missing = set(self.required_cols) - set(df.columns)
        if missing:
            raise ValueError(f"缺少必需列: {missing}")

        # 检查空 DataFrame
        if df.empty:
            raise ValueError("DataFrame 为空")

        # 添加元数据列
        df["_validated_at"] = pd.Timestamp.now()
        print(f"  验证通过: {len(df)} 行")
        return df


# ── 步骤 4:缺失值处理 ────────────────────────────────────────

class HandleMissingStep(PipelineStep):
    """缺失值处理步骤"""

    def __init__(
        self,
        numeric_strategy: str = "median",
        categorical_fill: str = "未知",
        drop_threshold: float = 0.8,
    ) -> None:
        super().__init__(name="缺失值处理")
        self.numeric_strategy = numeric_strategy
        self.categorical_fill = categorical_fill
        self.drop_threshold = drop_threshold

    def process(self, df: pd.DataFrame) -> pd.DataFrame:
        # 删除缺失率过高的列
        missing_rate = df.isnull().sum() / len(df)
        drop_cols = missing_rate[missing_rate > self.drop_threshold].index.tolist()
        if drop_cols:
            df = df.drop(columns=drop_cols)
            print(f"  删除高缺失列: {drop_cols}")

        # 数值列填充
        numeric_cols = df.select_dtypes(include=[np.number]).columns
        for col in numeric_cols:
            if df[col].isnull().any():
                strategies = {
                    "mean": df[col].mean(),
                    "median": df[col].median(),
                    "zero": 0,
                }
                fill_val = strategies.get(self.numeric_strategy, df[col].median())
                count = df[col].isnull().sum()
                df[col] = df[col].fillna(fill_val)
                print(f"  [{col}] 填充 {count} 个缺失值 ({self.numeric_strategy}={fill_val:.2f})")

        # 分类列填充
        cat_cols = df.select_dtypes(include=["object", "string"]).columns
        for col in cat_cols:
            if df[col].isnull().any():
                count = df[col].isnull().sum()
                df[col] = df[col].fillna(self.categorical_fill)
                print(f"  [{col}] 填充 {count} 个缺失值 ('{self.categorical_fill}')")

        return df


# ── 步骤 5:异常值处理 ────────────────────────────────────────

class HandleOutlierStep(PipelineStep):
    """异常值处理步骤"""

    def __init__(self, method: str = "iqr", factor: float = 1.5) -> None:
        super().__init__(name="异常值处理")
        self.method = method
        self.factor = factor

    def process(self, df: pd.DataFrame) -> pd.DataFrame:
        numeric_cols = df.select_dtypes(include=[np.number]).columns
        for col in numeric_cols:
            if self.method == "iqr":
                q1 = df[col].quantile(0.25)
                q3 = df[col].quantile(0.75)
                iqr = q3 - q1
                lower = q1 - self.factor * iqr
                upper = q3 + self.factor * iqr
                outlier_count = ((df[col] < lower) | (df[col] > upper)).sum()
                if outlier_count > 0:
                    df[col] = df[col].clip(lower, upper)
                    print(f"  [{col}] IQR 截断 {outlier_count} 个异常值 [{lower:.2f}, {upper:.2f}]")
            elif self.method == "zscore":
                z = np.abs((df[col] - df[col].mean()) / df[col].std())
                outlier_count = (z > 3).sum()
                if outlier_count > 0:
                    lower = df[col].quantile(0.01)
                    upper = df[col].quantile(0.99)
                    df[col] = df[col].clip(lower, upper)
                    print(f"  [{col}] Z-Score 截断 {outlier_count} 个异常值")
        return df


# ── 步骤 6:重复值处理 ────────────────────────────────────────

class HandleDuplicateStep(PipelineStep):
    """重复值处理步骤"""

    def __init__(self, key_cols: list[str] | None = None, keep: str = "first") -> None:
        super().__init__(name="重复值处理")
        self.key_cols = key_cols
        self.keep = keep

    def process(self, df: pd.DataFrame) -> pd.DataFrame:
        before = len(df)
        if self.key_cols:
            df = df.drop_duplicates(subset=self.key_cols, keep=self.keep)
        else:
            df = df.drop_duplicates(keep=self.keep)
        removed = before - len(df)
        if removed > 0:
            print(f"  去重: 移除 {removed} 行 (保留: {self.keep})")
        return df.reset_index(drop=True)


# ── 步骤 7:类型转换 ──────────────────────────────────────────

class TypeConvertStep(PipelineStep):
    """类型转换步骤"""

    def __init__(
        self,
        date_cols: list[str] | None = None,
        category_cols: list[str] | None = None,
        int_cols: list[str] | None = None,
    ) -> None:
        super().__init__(name="类型转换")
        self.date_cols = date_cols or []
        self.category_cols = category_cols or []
        self.int_cols = int_cols or []

    def process(self, df: pd.DataFrame) -> pd.DataFrame:
        # 日期列
        for col in self.date_cols:
            if col in df.columns:
                df[col] = pd.to_datetime(df[col], errors="coerce")
                print(f"  [{col}] 转换为 datetime")

        # 分类列
        for col in self.category_cols:
            if col in df.columns:
                df[col] = df[col].astype("category")
                print(f"  [{col}] 转换为 category")

        # 整数列
        for col in self.int_cols:
            if col in df.columns:
                df[col] = pd.to_numeric(df[col], errors="coerce").astype("Int64")
                print(f"  [{col}] 转换为 Int64")

        return df


# ── 步骤 8:数据导出 ──────────────────────────────────────────

class ExportStep(PipelineStep):
    """数据导出步骤"""

    def __init__(self, output_path: str, format: str = "parquet") -> None:
        super().__init__(name="数据导出")
        self.output_path = output_path
        self.format = format

    def process(self, df: pd.DataFrame) -> pd.DataFrame:
        path = Path(self.output_path)
        path.parent.mkdir(parents=True, exist_ok=True)

        if self.format == "parquet":
            df.to_parquet(self.output_path, index=False)
        elif self.format == "csv":
            df.to_csv(self.output_path, index=False)
        elif self.format == "json":
            df.to_json(self.output_path, orient="records", force_ascii=False, indent=2)
        else:
            raise ValueError(f"不支持的格式: {self.format}")

        size_mb = path.stat().st_size / 1024 / 1024
        print(f"  导出: {self.output_path} ({size_mb:.2f} MB, {len(df)} 行)")
        return df


# ── 组装管道 ──────────────────────────────────────────────────

def build_csv_pipeline(
    input_path: str,
    output_path: str,
    required_cols: list[str] | None = None,
    key_cols: list[str] | None = None,
    date_cols: list[str] | None = None,
    category_cols: list[str] | None = None,
) -> DataPipeline:
    """构建 CSV 数据清洗管道"""
    pipeline = DataPipeline("csv_clean_pipeline")
    pipeline.add_step(ReadCSVStep(input_path))
    pipeline.add_step(NormalizeColumnsStep())
    pipeline.add_step(ValidateStep(required_cols=required_cols))
    pipeline.add_step(HandleMissingStep(numeric_strategy="median", categorical_fill="未知"))
    pipeline.add_step(HandleOutlierStep(method="iqr"))
    pipeline.add_step(HandleDuplicateStep(key_cols=key_cols))
    pipeline.add_step(TypeConvertStep(date_cols=date_cols, category_cols=category_cols))
    pipeline.add_step(ExportStep(output_path, format="parquet"))
    return pipeline


# ── 运行示例 ──────────────────────────────────────────────────

if __name__ == "__main__":
    pipeline = build_csv_pipeline(
        input_path="raw_sales.csv",
        output_path="clean_sales.parquet",
        required_cols=["order_id", "product", "amount"],
        key_cols=["order_id"],
        date_cols=["order_date", "created_at"],
        category_cols=["product", "status"],
    )

    # 运行管道
    result = pipeline.run(pd.DataFrame())

    # 查看运行报告
    report = pipeline.report()
    print("\n管道运行报告:")
    print(report.to_string(index=False))

场景二:多数据源合并清洗管道

API + CSV + Database 三种数据源的合并清洗管道。

python
from __future__ import annotations

import json
import sqlite3
from datetime import datetime
from typing import Any

import pandas as pd
import requests

import logging

logger = logging.getLogger(__name__)


# ── 数据源读取器 ──────────────────────────────────────────────

class APISourceReader:
    """API 数据源读取器"""

    def __init__(self, url: str, params: dict | None = None, headers: dict | None = None) -> None:
        self.url = url
        self.params = params or {}
        self.headers = headers or {}

    def read(self) -> pd.DataFrame:
        """从 API 读取数据"""
        try:
            response = requests.get(self.url, params=self.params, headers=self.headers, timeout=30)
            response.raise_for_status()
            data = response.json()

            # 适配常见 API 响应格式
            if isinstance(data, list):
                records = data
            elif isinstance(data, dict):
                # 尝试常见的数据字段名
                for key in ["data", "results", "items", "records"]:
                    if key in data:
                        records = data[key]
                        break
                else:
                    records = [data]
            else:
                records = []

            df = pd.DataFrame(records)
            df["_source"] = "api"
            df["_ingested_at"] = datetime.now().isoformat()
            print(f"  API 读取: {len(df)} 条记录")
            return df

        except requests.RequestException as e:
            logger.error(f"API 请求失败: {e}")
            return pd.DataFrame()


class CSVSourceReader:
    """CSV 数据源读取器"""

    def __init__(self, file_path: str, **read_kwargs) -> None:
        self.file_path = file_path
        self.read_kwargs = read_kwargs

    def read(self) -> pd.DataFrame:
        """从 CSV 读取数据"""
        df = pd.read_csv(self.file_path, **self.read_kwargs)
        df["_source"] = "csv"
        df["_ingested_at"] = datetime.now().isoformat()
        print(f"  CSV 读取: {len(df)} 条记录")
        return df


class DatabaseSourceReader:
    """数据库数据源读取器"""

    def __init__(self, connection_string: str, query: str) -> None:
        self.connection_string = connection_string
        self.query = query

    def read(self) -> pd.DataFrame:
        """从数据库读取数据"""
        conn = sqlite3.connect(self.connection_string)
        try:
            df = pd.read_sql(self.query, conn)
            df["_source"] = "database"
            df["_ingested_at"] = datetime.now().isoformat()
            print(f"  数据库读取: {len(df)} 条记录")
            return df
        finally:
            conn.close()


# ── 多源合并管道 ──────────────────────────────────────────────

class MultiSourcePipeline:
    """多数据源合并清洗管道"""

    def __init__(self, name: str) -> None:
        self.name = name
        self.readers: dict[str, Any] = {}
        self.clean_steps: list[PipelineStep] = []
        self.merge_config: dict[str, Any] = {}

    def add_source(self, name: str, reader: Any) -> MultiSourcePipeline:
        """添加数据源"""
        self.readers[name] = reader
        return self

    def add_clean_step(self, step: PipelineStep) -> MultiSourcePipeline:
        """添加清洗步骤(应用于所有数据源)"""
        self.clean_steps.append(step)
        return self

    def set_merge_config(self, on: list[str], how: str = "outer", dedup: bool = True) -> MultiSourcePipeline:
        """设置合并配置"""
        self.merge_config = {"on": on, "how": how, "dedup": dedup}
        return self

    def run(self) -> pd.DataFrame:
        """执行多源合并管道"""
        logger.info(f"多源管道 [{self.name}] 启动")

        # 1. 读取所有数据源
        source_data = {}
        for name, reader in self.readers.items():
            try:
                source_data[name] = reader.read()
            except Exception as e:
                logger.error(f"数据源 [{name}] 读取失败: {e}")
                source_data[name] = pd.DataFrame()

        # 2. 分别清洗
        cleaned = {}
        for name, df in source_data.items():
            if df.empty:
                continue
            for step in self.clean_steps:
                df, _ = step.run(df)
            cleaned[name] = df

        # 3. 合并
        if not cleaned:
            logger.warning("所有数据源为空")
            return pd.DataFrame()

        dfs = list(cleaned.values())
        result = dfs[0]
        for i, df in enumerate(dfs[1:], 1):
            merge_on = self.merge_config.get("on", [])
            merge_how = self.merge_config.get("how", "outer")
            suffixes = ("", f"_{list(cleaned.keys())[i]}")
            result = pd.merge(result, df, on=merge_on, how=merge_how, suffixes=suffixes)

        # 4. 合并后去重
        if self.merge_config.get("dedup", True) and self.merge_config.get("on"):
            before = len(result)
            result = result.drop_duplicates(subset=self.merge_config["on"], keep="last")
            result = result.reset_index(drop=True)
            print(f"  合并后去重: {before} → {len(result)} 行")

        # 5. 添加元数据
        result["_merged_at"] = datetime.now().isoformat()
        result["_pipeline"] = self.name

        logger.info(f"多源管道 [{self.name}] 完成: {len(result)} 行")
        return result


# ── 运行示例 ──────────────────────────────────────────────────

def build_multi_source_pipeline() -> MultiSourcePipeline:
    """构建多数据源合并管道"""

    pipeline = MultiSourcePipeline("user_data_merge")

    # 添加数据源
    pipeline.add_source("api", APISourceReader(
        url="https://api.example.com/users",
        headers={"Authorization": "Bearer token123"},
    ))
    pipeline.add_source("csv", CSVSourceReader("legacy_users.csv"))
    pipeline.add_source("db", DatabaseSourceReader(
        connection_string="users.db",
        query="SELECT * FROM active_users WHERE status = 'active'",
    ))

    # 添加清洗步骤(应用于所有数据源)
    pipeline.add_clean_step(NormalizeColumnsStep())
    pipeline.add_clean_step(HandleMissingStep(numeric_strategy="median"))
    pipeline.add_clean_step(HandleDuplicateStep(key_cols=["user_id"]))

    # 设置合并配置
    pipeline.set_merge_config(on=["user_id"], how="outer", dedup=True)

    return pipeline


if __name__ == "__main__":
    pipeline = build_multi_source_pipeline()
    result = pipeline.run()

    # 导出
    result.to_parquet("merged_users.parquet", index=False)
    print(f"\n最终输出: {len(result)} 行 × {len(result.columns)} 列")
    print(f"数据来源分布: {result['_source'].value_counts().to_dict()}")

多源合并冲突解决策略

当多个数据源对同一实体提供不同值时,需要冲突解决策略。

python
class ConflictResolver:
    """多源数据冲突解决器"""

    @staticmethod
    def resolve_by_priority(
        df: pd.DataFrame,
        key_col: str,
        conflict_cols: list[str],
        source_priority: list[str],  # 数据源优先级,从高到低
    ) -> pd.DataFrame:
        """按数据源优先级解决冲突"""
        df = df.copy()

        for col in conflict_cols:
            # 找到冲突列的所有版本(可能有 _api, _csv 等后缀)
            versions = [c for c in df.columns if c.startswith(col)]
            if len(versions) <= 1:
                continue

            # 按优先级选择非空值
            result_col = []
            for _, row in df.iterrows():
                resolved = None
                for source in source_priority:
                    col_name = f"{col}_{source}" if f"{col}_{source}" in versions else col
                    if col_name in row and pd.notna(row[col_name]):
                        resolved = row[col_name]
                        break
                result_col.append(resolved)

            df[col] = result_col
            # 删除冲突版本列
            drop_cols = [c for c in versions if c != col]
            df = df.drop(columns=drop_cols)

        return df

    @staticmethod
    def resolve_by_recency(
        df: pd.DataFrame, key_col: str, conflict_cols: list[str], time_col: str = "_ingested_at"
    ) -> pd.DataFrame:
        """按数据新鲜度解决冲突(取最新值)"""
        df = df.copy()
        df[time_col] = pd.to_datetime(df[time_col])

        for col in conflict_cols:
            versions = [c for c in df.columns if c.startswith(col)]
            if len(versions) <= 1:
                continue

            # 选择最新时间戳对应的值
            for idx, row in df.iterrows():
                latest_val = None
                latest_time = pd.Timestamp.min
                for v in versions:
                    if pd.notna(row[v]):
                        latest_val = row[v]
                df.at[idx, col] = latest_val

            drop_cols = [c for c in versions if c != col]
            df = df.drop(columns=drop_cols)

        return df

常见陷阱

陷阱现象原因解决方案
修改原始 DataFrame下游步骤看到意外数据Pandas 默认不复制数据每步开头 df = df.copy()
SettingWithCopyWarning链式赋值警告修改了视图而非副本.loc[].copy()
统一用 0 填充缺失值数据分布被破坏0 在不同列含义不同按列选择填充策略(中位数/众数/插值)
IQR 处理偏态分布误删大量正常值IQR 假设近似正态分布偏态数据用百分位法或对数变换后检测
忽略编码问题中文乱码文件编码与读取编码不一致用 chardet 检测编码,或尝试多种编码
合并产生笛卡尔积行数暴增多对多合并合并前检查键的唯一性,用 validate 参数
增量更新遗漏数据部分数据未被处理Watermark 时间戳精度不够使用数据库事务时间而非业务时间
管道无日志出错无法定位缺少运行统计和日志每步记录输入/输出行数、耗时、异常
类型推断错误数值列被读为字符串CSV 中混合了数字和文本读取时指定 dtype,或用 pd.to_numeric
时区不一致时间数据错乱不同源使用不同时区统一转为 UTC,显示时再转本地时区
去重保留错误记录最新数据被删除keep="first" 保留了旧记录按时间排序后再去重,或用 keep="last"
大文件内存溢出OOM 错误一次性加载全部数据chunksize 分块处理或使用 Polars/Dask

陷阱详解:修改原始 DataFrame

python
# ❌ 反面:直接修改传入的 DataFrame
def clean_data(df: pd.DataFrame) -> pd.DataFrame:
    df["age"] = df["age"].fillna(df["age"].median())  # 修改了原始 df!
    return df

original = pd.DataFrame({"age": [25, None, 30]})
cleaned = clean_data(original)
print(original["age"].isnull().sum())  # 0 — 原始数据被污染了!

# ✅ 正面:复制后修改
def clean_data(df: pd.DataFrame) -> pd.DataFrame:
    df = df.copy()  # 显式复制
    df["age"] = df["age"].fillna(df["age"].median())
    return df

original = pd.DataFrame({"age": [25, None, 30]})
cleaned = clean_data(original)
print(original["age"].isnull().sum())  # 1 — 原始数据完好

陷阱详解:IQR 处理偏态分布

python
# ❌ 反面:对收入数据直接用 IQR
# 收入数据通常右偏,IQR 会把高收入正常值标记为异常
income = pd.Series([5000, 6000, 5500, 5800, 50000, 6200, 5300])
q1, q3 = income.quantile(0.25), income.quantile(0.75)
iqr = q3 - q1
upper = q3 + 1.5 * iqr
print(f"上界: {upper:.0f}")  # 可能只有 ~8000,50000 被截断

# ✅ 正面:对数变换后检测,或用百分位法
log_income = np.log1p(income)
q1, q3 = log_income.quantile(0.25), log_income.quantile(0.75)
iqr = q3 - q1
upper = np.expm1(q3 + 1.5 * iqr)  # 转回原始尺度

# ✅ 正面:直接用百分位法
lower = income.quantile(0.01)
upper = income.quantile(0.99)
cleaned = income.clip(lower, upper)

最佳实践速查表

场景推荐做法避免
管道设计每步只做一件事,返回新 DataFrame在一个函数里做所有清洗
数据复制每步开头 df.copy()直接修改传入的 DataFrame
缺失值按列选择策略(中位数/众数/插值)统一用 0 或删除
异常值先标记再决定处理方式直接删除所有异常值
重复值明确去重键和保留策略盲目 drop_duplicates()
类型转换读取时指定 dtype依赖 Pandas 自动推断
编码chardet 检测 + 多编码回退假定所有文件都是 UTF-8
增量更新Watermark + 事务时间依赖业务时间戳
合并先检查键唯一性再 merge直接 merge 不验证
日志每步记录行数变化和耗时管道静默运行
导出Parquet(更快更小)到处用 CSV
大文件chunksize / Polars / Dask一次性 read_csv
元数据添加 _source_ingested_at丢失数据来源信息
Schema用 Pandera 或自定义 Schema 验证不验证直接处理
管道测试用小数据集单元测试每个步骤只在完整数据上测试

术语表

术语英文定义
ETLExtract, Transform, Load数据提取、转换、加载的管道模式
幂等性Idempotency同一操作多次执行产生相同结果的性质
WatermarkWatermark标记数据处理进度的位置标记,用于增量更新
CDCChange Data Capture变更数据捕获,通过日志识别数据变更
SchemaSchema数据结构的定义,包括列名、类型和约束
缺失值Missing Value数据中的空值(NaN、None、NA)
异常值Outlier显著偏离大多数数据的极端值
IQRInterquartile Range四分位距:Q3 - Q1,用于识别异常值
Z-ScoreZ-Score标准分数,衡量数据点偏离均值的标准差数
笛卡尔积Cartesian Product多对多合并产生的行数爆炸
UpsertUpdate or Insert存在则更新,不存在则插入
SCD2Slowly Changing Dimension Type 2数据仓库中历史版本管理策略
ParquetApache Parquet列式存储格式,比 CSV 更快更小
BOMByte Order Mark文件开头的编码标识字节
数据血缘Data Lineage数据从源头到终点的完整流转路径
管道编排Pipeline Orchestration管道步骤的调度和依赖管理
分支管道Branch Pipeline按条件将数据分发到不同子管道
合并管道Merge Pipeline多个数据源清洗后合并的管道模式

延伸阅读

官方文档

经典参考

  • 《数据仓库工具箱》(Ralph Kimball)— ETL 与维度建模经典
  • 《Designing Data-Intensive Applications》(Martin Kleppmann)— 数据系统设计圣经
  • The Data Engineering Cookbook — 数据工程实践指南

推荐阅读

版本差异(数据科学栈 → 当前版本)

本文编写时当前稳定版升级要点
Python3.8-3.123.143.12+ 起性能显著提升;3.14 PEP 649/750
NumPy1.x/2.02.3.xnp.float_ 等别名移除;NEP 50 类型提升
Pandas1.x/2.x3.0.xCopy-on-Write 默认开启;inplace 移除;字符串 dtype 变化
Matplotlib3.x3.x 稳定版API 兼容,样式更新
Seaborn0.12/0.130.13.xAPI 稳定
scikit-learn1.x1.7.xAPI 稳定,新算法持续加入

本文讲解的数据分析流程(读取→清洗→分析→可视化)与核心 API 在最新版本中成立;升级时重点关注 Pandas 3.0 的 Copy-on-Write 与 NumPy 2.x 的类型变化。