前言

大家好,今天想和大家聊聊 ETL 这个话题。

不管你是后端工程师、数据分析师,还是刚接触大数据领域的新人,在工作或学习中大概率都听过"ETL"这个词。ETL 是数据工程最核心的基础能力,也是构建数据仓库和数据管道的起点。

这篇文章完全面向零基础或刚入门的读者,我会用最通俗的语言来拆解 ETL 的方方面面,并附上一个完整的 Python 示例,希望能帮你建立起对 ETL 的系统认知。


什么是 ETL?核心概念与价值

ETL 的三个字母

ETL 是三个单词的缩写:

| 阶段 | 英文 | 中文 | 做什么 |

|------|------|------|--------|

| E | Extract | 抽取 | 从源系统获取数据 |

| T | Transform | 转换 | 清洗、加工、格式化数据 |

| L | Load | 加载 | 将数据写入目标系统 |

为什么需要 ETL?

想象一下:你的公司在用 MySQL 做业务数据库,销售人员用 Excel 记录客户信息,市场团队用 Google Analytics,客服在用 Zendesk——这些系统的数据散落各方、格式各异

你总不能每次做报表时,都手动打开 5 个系统拷贝粘贴吧?

ETL 要做的事就是:自动把这些"脏乱差"的原始数据,变成"干净整齐"的分析数据。具体来说:

  • **业务价值**:打通数据孤岛,让数据真正驱动决策
  • **技术价值**:减轻源库查询压力,数据统一建模,历史数据归档
  • **工程价值**:自动化可编排,出问题时能追溯和重跑
**一句话总结:ETL 是把"散落四处的数据原油"炼成"可直接使用的汽油"。**

数据仓库基础

要理解 ETL,必须先搞清楚数据往哪存——数据仓库

OLTP vs OLAP

这是两个几乎所有数据库文章都会提到的概念:

| 对比维度 | OLTP(联机事务处理) | OLAP(联机分析处理) |

|---------|-------------------|-------------------|

| 典型系统 | MySQL、PostgreSQL、Oracle | ClickHouse、Doris、Snowflake |

| 主要操作 | 增删改查(CRUD) | 复杂查询、聚合分析 |

| 数据特点 | 行式存储,实时更新 | 列式存储,大批量写入 |

| 用户 | 一线业务人员、终端用户 | 数据分析师、管理层 |

| 典型语句 | UPDATE order SET status = 'paid' WHERE id = 123 | `SELECT region, SUM(amount) FROM orders GROUP BY region` |

**简单记法:OLTP 管"今天卖了多少单",OLAP 管"这个季度每个区域的销售趋势"。**

ETL 通常是从 OLTP 抽数据,经过转换,加载到 OLAP 数据仓库

星型模型与雪花模型

数据仓库的建模方式有很多,初学者最需要掌握的是两种经典模型:

星型模型(Star Schema)


      事实表(Fact Table)
     ┌──────────────────────┐
     │ 订单ID │ 金额 │ 日期 │
     └────┬──────────┬──────┘
          │          │
     ┌────▼──┐  ┌────▼────┐
     │ 用户维度│  │ 产品维度 │
     └────────┘  └─────────┘
  • **中心**:一个事实表(fact table),存可量化的指标(销售额、数量)
  • **四周**:多个维度表(dimension table),存描述性信息(用户姓名、产品分类)
  • **优点**:查询简单、性能好,适合大多数 BI 工具
  • **适用**:初学者首选,大部分业务场景够用

雪花模型(Snowflake Schema)

雪花模型是星型的"规范化"版本——把维度表继续拆分成更细的表。比如"产品维度"拆成"产品表"和"分类表"。

  • **优点**:减少数据冗余,节省存储
  • **缺点**:查询时 JOIN 更多,性能稍差
  • **适用**:对存储敏感、层次结构复杂的场景
**新手建议**:一开始用星型模型就够了,90% 的场景下它是更好的选择。

ETL vs ELT:有什么区别?怎么选?

近年来随着 MPP 数据仓库(如 Snowflake、Redshift、ClickHouse)的流行,ELT 模式越来越常见。

ETL(传统方式)


抽取 → 转换(在中间层) → 加载到数仓
  • 数据先经过 ETL 引擎(如 Kettle、Python 脚本)做清洗转换
  • 转换过程在数仓**外部**完成
  • 适用于:源数据质量差、数仓计算能力弱、需要复杂业务逻辑

ELT(现代方式)


抽取 → 加载到数仓 → 在数仓内部转换(SQL)
  • 数据**原封不动**先加载到数仓临时区
  • 用 SQL/dbt 等工具在数仓内做转换
  • 适用于:数仓计算能力强、数据量超大、团队擅长 SQL

怎么选?

| 场景 | 推荐 |

|------|------|

| 数据量 < 100GB,源系统是传统数据库 | ETL |

| 数据量 > 1TB,使用云数仓(Snowflake 等) | ELT |

| 团队 SQL 能力强 | ELT |

| 需要复杂数据清洗和脱敏 | ETL |

| 实时性要求高、流式处理 | ETL(流式) |

一个具体场景来帮助理解

假设你要从业务数据库中抽取用户行为日志,每天大约 500 万条记录,目标是生成一份"用户活跃度报表":

  • **用 ETL**:在 Python/Spark 中过滤掉爬虫数据、解析 User-Agent、计算停留时长,再把结果写入数仓。
  • **用 ELT**:先把 500 万条原始日志全部 load 到 ClickHouse,然后用 SQL 一句 `SELECT ... GROUP BY ...` 完成转换。

前者把"重活"放在 ETL 引擎中做,后者交给数仓。如果你的数仓是 ClickHouse 或 Snowflake 这种高性能 MPP 引擎,ELT 通常更高效。

实际上,很多公司的数据管道是**两者混用**的:简单清洗用 ELT,复杂业务逻辑在中间层用 ETL 处理。

常用 ETL 工具对比

工欲善其事,必先利其器。给新手列一个常用工具的速览:

| 工具 | 类型 | 语言/UI | 适合场景 | 学习曲线 |

|------|------|---------|---------|---------|

| Apache NiFi | 流式/批量 | Web UI / Java | 物联网、日志采集、实时数据流 | 中等 |

| Apache Airflow | 调度编排 | Python DAG | 任务调度、工作流编排(不直接做 ETL,但调度 ETL 任务) | 中等偏高 |

| Kettle (PDI) | 批量 ETL | 图形化拖拽 | 传统数仓、中小规模 ETL | 较低 |

| DataX | 批量同步 | 配置文件 | 异构数据源之间的离线批量同步 | 低 |

| SeaTunnel | 批量/流式 | 配置文件 + Flink/Spark | 高性能海量数据同步 | 中等 |

| dbt | 数据转换 | SQL | 数仓内部 ELT 转换(搭配 ELT 模式) | 低 |

新手推荐路线

1. 先学用 Python 手写 ETL(类似本文最后的示例)——理解原理

2. 再学 Airflow——理解任务编排

3. 按需选择:如果公司用 Flink/Spark 生态,学 SeaTunnel;如果偏传统数仓,Kettle 也能快速上手


ETL 流程设计

1. 抽取策略(Extract)

全量抽取(Full Load)


# 伪代码
df = source.read("SELECT * FROM orders")
target.write(df)
  • **方式**:每次把源表全部数据读出来,覆盖写入目标
  • **优点**:实现简单,不会遗漏
  • **缺点**:数据量大时极慢,浪费存储和带宽
  • **适用**:数据量 < 10 万行、字典表(如地区表)

增量抽取(Incremental Load)


# 伪代码 - 基于时间戳
last_max = target.get_max("update_time")
df = source.read(f"SELECT * FROM orders WHERE update_time > '{last_max}'")
target.write(df)

常见增量方式:

| 方式 | 原理 | 优缺点 |

|------|------|--------|

| 时间戳 | 记录 updated_at 最大值 | 简单,但需源表有可靠的更新时间字段 |

| 增量日志(CDC) | 解析数据库 binlog(如 Debezium、Canal) | 实时、无侵入,但部署复杂 |

| 增量标志位 | 源表设 is_synced 标记 | 侵入性强 |

| 窗口对比 | 用主键范围或分片对比差异 | 适合无时间戳的表 |

**最佳实践**:小表全量 + 大表增量,是最常见的搭配。

2. 数据清洗(Transform)

数据清洗是 ETL 中最费时费力的一环,通常占整个 ETL 开发工作量的 60%~80%。为什么会这么耗时?因为现实世界的数据远比想象中的要混乱。

举几个真实的"脏数据"例子:

  • 用户姓名前后带空格:" 张三 "、"李四 "
  • 手机号格式不统一:"138-1234-5678"、"13812345678"、"+86 13812345678"
  • 金额字段出现负值,或者存成了字符串:"¥199.00"
  • 日期格式五花八门:"2024-01-15"、"01/15/2024"、"2024年1月15日"
  • 同一个分类在不同记录里写法不同:"手机"、"智能手机"、"Mobile Phone"

数据清洗就是要统一化、规范化、标准化这些散乱的数据。常见问题及处理:


import pandas as pd

df = pd.read_csv("raw_data.csv")

# ① 去除前后空格
df['name'] = df['name'].str.strip()

# ② 处理空值
df['age'] = df['age'].fillna(0)                     # 数值列填 0
df['email'] = df['email'].fillna('unknown@')         # 字符串列填默认值

# ③ 统一格式
df['phone'] = df['phone'].str.replace(r'\D', '', regex=True)  # 去除非数字

# ④ 去除重复行
df = df.drop_duplicates(subset=['user_id'])

# ⑤ 过滤异常值
df = df[df['amount'] >= 0]

# ⑥ 类型转换
df['created_at'] = pd.to_datetime(df['created_at'])

# ⑦ 衍生字段
df['year_month'] = df['created_at'].dt.strftime('%Y-%m')

3. 加载策略(Load)

| 策略 | 做法 | 适用场景 |

|------|------|---------|

| 全量覆盖 | TRUNCATE + INSERT | 小表、维度表 |

| 增量追加 | INSERT INTO ... | 日志表、不可变事件表 |

| 合并更新(UPSERT) | INSERT ... ON DUPLICATE KEY UPDATE | 订单表等需要更新的场景 |

| 拉链表(SCD Type 2) | 记录历史变化版本 | 用户档案等需追溯变化 |

对于初学者,先掌握"全量覆盖"和"增量追加"两种就够了,UPSERT 和拉链表是进阶内容。

完整 Python ETL 示例

下面我用 Python 写一个完整的 ETL 示例,把 MySQL(或 SQLite)中的订单数据抽取出来,做清洗转换,再加载到目标数据库中。

这个示例使用 **pandas** 做数据清洗,**SQLAlchemy** 做数据库读写。你可以直接在本地跑起来。

环境准备


pip install pandas sqlalchemy pymysql  # MySQL 驱动
# 或者
pip install pandas sqlalchemy duckdb   # 用 DuckDB 做本地测试(无需数据库服务)

示例代码


"""
etl_demo.py - 一个完整的 ETL 示例
从 SQLite 源库抽取订单数据 → 清洗转换 → 加载到目标库
"""

import pandas as pd
from sqlalchemy import create_engine, MetaData, Table, Column, Integer, String, Float, DateTime, Date
from datetime import datetime, timedelta
import logging

# ========== 配置日志 ==========
logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')
logger = logging.getLogger(__name__)


# ========== 第一步:Extract(抽取) ==========

def extract(source_conn_str: str, sql: str) -> pd.DataFrame:
    """
    从源数据库抽取数据
    """
    logger.info(f"[Extract] 开始从源库抽取数据...")
    engine = create_engine(source_conn_str)
    with engine.connect() as conn:
        df = pd.read_sql(sql, conn)
    logger.info(f"[Extract] 抽取完成,共 {len(df)} 行")
    return df


# ========== 第二步:Transform(转换/清洗) ==========

def transform(raw_df: pd.DataFrame) -> pd.DataFrame:
    """
    数据清洗与转换
    """
    logger.info(f"[Transform] 开始清洗数据,原始行数: {len(raw_df)}")

    df = raw_df.copy()

    # --- 数据清洗 ---

    # 1. 删除完全重复行
    before = len(df)
    df = df.drop_duplicates()
    logger.info(f"  去重: {before} → {len(df)}(删除 {before - len(df)} 行)")

    # 2. 处理空值
    df['customer_name'] = df['customer_name'].fillna('未知客户')
    df['amount'] = df['amount'].fillna(0.0)
    df['status'] = df['status'].fillna('pending')

    # 3. 去除字符串前后空格
    for col in ['customer_name', 'product', 'status']:
        if col in df.columns:
            df[col] = df[col].astype(str).str.strip()

    # 4. 金额必须 >= 0,负值取绝对值(也可标记为异常)
    df['amount'] = df['amount'].abs()

    # 5. 统一状态字段
    status_mapping = {'待处理': 'pending', '已完成': 'completed', '已取消': 'cancelled'}
    df['status'] = df['status'].map(status_mapping).fillna(df['status'])

    # --- 数据转换 ---

    # 6. 时间字段标准化
    df['order_date'] = pd.to_datetime(df['order_date'], errors='coerce')
    df['order_date_str'] = df['order_date'].dt.strftime('%Y-%m-%d')

    # 7. 衍生字段:订单月份
    df['order_month'] = df['order_date'].dt.to_period('M').astype(str)

    # 8. 计算字段:金额等级
    df['amount_level'] = pd.cut(
        df['amount'],
        bins=[0, 100, 500, float('inf')],
        labels=['小额', '中额', '大额']
    )

    logger.info(f"[Transform] 清洗完成,最终行数: {len(df)}")
    return df


# ========== 第三步:Load(加载) ==========

def load(target_conn_str: str, table_name: str, df: pd.DataFrame, if_exists: str = 'replace'):
    """
    加载数据到目标数据仓库
    if_exists: 'replace'(覆写) / 'append'(追加) / 'upsert'(需自定义)
    """
    logger.info(f"[Load] 开始加载数据到目标表 {table_name}...")
    engine = create_engine(target_conn_str)

    with engine.begin() as conn:
        df.to_sql(
            name=table_name,
            con=conn,
            if_exists=if_exists,
            index=False,
            dtype={
                'order_id': String(50),
                'customer_name': String(100),
                'product': String(200),
                'amount': Float(),
                'status': String(20),
                'order_date': DateTime(),
                'order_date_str': String(10),
                'order_month': String(7),
                'amount_level': String(10),
            }
        )

    logger.info(f"[Load] 加载完成,共写入 {len(df)} 行到 {table_name}")


# ========== 主流程 ==========

def run_etl():
    """
    执行完整的 ETL 流程
    """
    # 数据库连接信息(这里用 SQLite 做演示,换 MySQL 只需改连接串)
    SOURCE_DB = "sqlite:///source_orders.db"
    TARGET_DB = "sqlite:///dw_orders.db"

    # 源数据查询 SQL
    EXTRACT_SQL = """
        SELECT order_id, customer_name, product, amount, status, order_date
        FROM orders
        WHERE order_date >= date('now', '-30 days')
    """

    TARGET_TABLE = "fact_orders"

    logger.info("=" * 50)
    logger.info("ETL 流程开始")
    logger.info("=" * 50)

    try:
        # Step 1: Extract
        raw_data = extract(SOURCE_DB, EXTRACT_SQL)

        if raw_data.empty:
            logger.warning("源数据为空,跳过本次 ETL")
            return

        # Step 2: Transform
        clean_data = transform(raw_data)

        # Step 3: Load
        load(TARGET_DB, TARGET_TABLE, clean_data, if_exists='append')

        logger.info("=" * 50)
        logger.info("ETL 流程成功完成 ✅")
        logger.info("=" * 50)

    except Exception as e:
        logger.error(f"ETL 流程失败: {e}")
        raise


# ========== 辅助:初始化测试数据 ==========

def init_test_data():
    """在 source 库中生成测试数据供演示"""
    engine = create_engine("sqlite:///source_orders.db")

    # 生成 100 条示例订单
    np = __import__('numpy', globals(), locals(), [], 0)
    records = []
    for i in range(1, 101):
        records.append({
            'order_id': f'ORD{i:05d}',
            'customer_name': f'客户{i}' if i % 10 != 0 else None,
            'product': ['手机', '电脑', '耳机', '键盘', '鼠标'][i % 5],
            'amount': round(abs(np.random.normal(300, 150)), 2),
            'status': np.random.choice(['pending', 'completed', 'cancelled', '待处理'], p=[0.3, 0.5, 0.1, 0.1]),
            'order_date': datetime.now() - timedelta(days=np.random.randint(0, 60))
        })

    df_test = pd.DataFrame(records)
    df_test.to_sql('orders', engine, if_exists='replace', index=False)
    print("测试数据已写入 source_orders.db")


if __name__ == '__main__':
    # 首次运行先初始化测试数据
    import os
    if not os.path.exists('source_orders.db'):
        init_test_data()

    run_etl()

运行结果


python etl_demo.py

你会看到类似这样的日志输出:


2026-06-22 10:00:00 - INFO - ==================================================
2026-06-22 10:00:00 - INFO - ETL 流程开始
2026-06-22 10:00:00 - INFO - ==================================================
2026-06-22 10:00:00 - INFO - [Extract] 开始从源库抽取数据...
2026-06-22 10:00:00 - INFO - [Extract] 抽取完成,共 100 行
2026-06-22 10:00:00 - INFO - [Transform] 开始清洗数据,原始行数: 100
2026-06-22 10:00:00 - INFO -   去重: 100 → 98(删除 2 行)
2026-06-22 10:00:00 - INFO - [Transform] 清洗完成,最终行数: 98
2026-06-22 10:00:00 - INFO - [Load] 开始加载数据到目标表 fact_orders...
2026-06-22 10:00:00 - INFO - [Load] 加载完成,共写入 98 行到 fact_orders
2026-06-22 10:00:00 - INFO - ==================================================
2026-06-22 10:00:00 - INFO - ETL 流程成功完成 ✅
2026-06-22 10:00:00 - INFO - ==================================================

这个示例虽然很简,但麻雀虽小五脏俱全——它覆盖了 ETL 的核心流程:从源库抽取 → 清洗去重空值处理 → 类型转换字段衍生 → 加载到目标库。你只需要把连接串改成 MySQL/PostgreSQL,就可以直接用在真实项目中了。


ETL 项目最佳实践 & 注意事项

最后分享一些我在实际项目中积累的经验:

1. 可观测性(Observability)

ETL 任务必须有完善的日志和告警


# 建议记录的关键指标
etl_metrics = {
    "source_rows": 10000,       # 源数据行数
    "dropped_rows": 150,        # 丢弃行数
    "loaded_rows": 9850,        # 加载行数
    "duration_sec": 45.2,       # 耗时
    "status": "success",        # 状态
}
  • 用 Prometheus + Grafana 监控任务指标
  • 失败时通过飞书/钉钉/邮件告警
  • 推荐工具:Airflow 自带日志 + 自定义回调函数

2. 幂等性(Idempotency)

同一个 ETL 任务无论跑多少次,结果应该一样。

  • 全量加载:先 `TRUNCATE` 再写
  • 增量加载:用 `UPSERT` 而不是 `INSERT`,避免重复

3. 数据质量检查

在 ETL 的每个环节都要加校验点

| 阶段 | 检查项 |

|------|--------|

| Extract 后 | 行数 > 0?字段数对不对? |

| Transform 后 | 空值比例是否正常?金额是否有负数? |

| Load 后 | 目标行数 = 预期行数? |

4. 增量抽取的"水印"管理


# 推荐用一张独立的 watermark 表管理增量点位
def get_watermark(table_name: str) -> datetime:
    """获取上次同步的最大时间戳"""
    query = f"SELECT max_value FROM etl_watermark WHERE table_name = '{table_name}'"
    ...

def set_watermark(table_name: str, max_value: datetime):
    """更新水印"""
    ...

5. 做好数据血缘(Data Lineage)

当你的 ETL 任务多起来以后(十几个、几十个表),一定会遇到一个问题:某个下游报表数据不对,到底是谁的锅?

数据血缘就是要回答这个问题——记录每个数据字段从哪里来、经过了哪些转换、最终去了哪里。你可以用以下方式管理:

  • 简单场景:在代码注释中写明数据来源和业务含义
  • 中等场景:用自建的 `etl_meta` 元数据表,记录每个字段的 mapping 关系
  • 专业场景:用 Apache Atlas、DataHub 等元数据管理平台自动采集血缘
数据血缘做得好,排查问题的时间能从"半天"缩小到"十分钟"。

6. 避坑清单


❌ 不要在 Transform 阶段做太重的计算——ETL 引擎不是数仓
❌ 不要忽略源表 schema 变更——上游加了一列,你的 ETL 可能就崩了
❌ 不要把密码写在代码里——用环境变量或密钥管理服务
❌ 不要一次性加载所有数据——分批处理,控制内存
✅ 核心表一定要做数据质量监控
✅ 每个 ETL 任务都要有回溯重跑能力
✅ 尽量保持 ETL 逻辑简单——"能跑"不等于"好维护"

总结

这篇文章从零开始介绍了 ETL 的核心概念:

| 知识点 | 一句话记住 |

|--------|-----------|

| ETL | 从源库抽数据 → 清洗转换 → 写入数仓 |

| 数据仓库 | OLAP 系统,存的是分析用的聚合、历史数据 |

| 星型模型 | 一张事实表 + 多张维度表,最常用 |

| ETL vs ELT | 转换在外还是在内?数仓强则 ELT,源差则 ETL |

| 抽取策略 | 小表全量,大表增量 |

| 工具选择 | 入门用手写 Python 理解原理,进阶用 Airflow/dbt |

ETL 本身不是一个高深的概念,但它的难点在于真实数据的复杂性。你可能遇到编码问题、字段类型不匹配、源库 schema 变更、数据量爆发增长……这些都是要靠实战去积累经验的。

如果你是一个新手,我建议你:

1. 用上面那个 Python 示例,在自己的电脑上跑一遍

2. 换一个真实的数据库(比如本机 MySQL)做源和目标

3. 试着给数据加一些"脏数据"(空值、乱码、重复),看 ETL 能否处理好

只有亲手踩过坑,才能真正理解 ETL。

希望这篇文章对你有帮助。如果有任何问题或想法,欢迎在评论区交流讨论 🙌


*本文是「数据工程入门系列」的第一篇,后续会继续写数据建模、Airflow 实战、数据质量等内容,敬请关注。*