为什么用 Python 做 ETL?

Python 是数据工程领域最受欢迎的语言之一。原因很简单:语法简洁、生态丰富、社区庞大。对于 ETL 开发来说,Python 的优势尤其明显。

Python 拥有几乎所有数据源对应的库。连接数据库有 SQLAlchemy 和 psycopg2,处理文件有 Pandas 和 PyArrow,调用 API 有 requests 和 httpx,操作云存储有 boto3 和 google-cloud-storage。你几乎不需要从零开始写任何底层代码。

另一个关键优势是灵活性。SQL 擅长做声明式的数据转换,但遇到复杂的业务逻辑、条件分支、循环处理时,Python 的表达力远超 SQL。ETL 流程中那些脏活累活——数据清洗、异常处理、重试机制、日志记录——用 Python 写起来得心应手。

核心库速览

Pandas:数据转换的主力

Pandas 是 Python ETL 中最核心的库。它提供了 DataFrame 数据结构,可以理解为一个内存中的表格,支持各种数据操作:过滤、聚合、关联、透视、窗口函数等。

`python import pandas as pd

从 CSV 读取数据

df = pd.read_csv(“sales.csv”)

数据清洗

df = df.dropna(subset=[“amount”]) df[“amount”] = df[“amount”].clip(lower=0) # 过滤负值

数据转换

df[“order_month”] = pd.to_datetime(df[“order_date”]).dt.to_period(“M”) df[“total_price”] = df[“quantity”] * df[“unit_price”]

聚合计算

monthly = df.groupby(“order_month”).agg( total_revenue=(“total_price”, “sum”), order_count=(“order_id”, “nunique”) ).reset_index() `

Pandas 的链式调用风格让数据转换代码非常简洁。但要注意,Pandas 是单机内存处理,数据量超过内存容量时会出问题。对于超大数据集,可以考虑 Dask 或 PySpark。

SQLAlchemy:数据库连接的瑞士军刀

SQLAlchemy 是 Python 生态中最成熟的 ORM 和数据库工具包。在 ETL 场景中,我们主要用它来建立数据库连接、执行 SQL 语句、读写 DataFrame。

`python from sqlalchemy import create_engine, text

创建数据库连接

engine = create_engine( “postgresql://user:password@host:5432/database” )

使用 Pandas 读取数据库表

df = pd.read_sql(“SELECT * FROM orders WHERE date >= ‘2026-01-01’”, engine)

执行原生 SQL

with engine.connect() as conn: result = conn.execute(text(“SELECT count(*) FROM orders”)) print(result.scalar())

将 DataFrame 写入数据库

df.to_sql(“orders_clean”, engine, if_exists=“replace”, index=False) `

SQLAlchemy 支持多种数据库方言:PostgreSQL、MySQL、SQLite、Oracle、SQL Server 等。切换数据库只需要改一下连接字符串,代码完全不用动。

requests:API 数据抽取

很多现代数据源通过 REST API 提供数据访问。requests 库是 Python 中最流行的 HTTP 客户端。

`python import requests import pandas as pd from datetime import datetime, timedelta

def extract_from_api(base_url: str, api_key: str, days_back: int = 30): """从 API 抽取数据""" headers = {“Authorization”: f"Bearer {api_key}"} params = { “start_date”: (datetime.now() - timedelta(days=days_back)).date().isoformat(), “end_date”: datetime.now().date().isoformat(), “page”: 1, “per_page”: 1000 }

all_records = []
while True:
    response = requests.get(
        f"{base_url}/orders",
        headers=headers,
        params=params
    )
    response.raise_for_status()
    data = response.json()

    all_records.extend(data["items"])

    if data["page"] >= data["total_pages"]:
        break
    params["page"] += 1

return pd.DataFrame(all_records)

`

boto3:云存储集成

如果数据源在 AWS S3 上,boto3 是必备工具。

`python import boto3 import pandas as pd from io import BytesIO

def extract_from_s3(bucket: str, prefix: str, aws_profile: str = “default”): """从 S3 读取 CSV 文件""" session = boto3.Session(profile_name=aws_profile) s3 = session.client(“s3”)

response = s3.list_objects_v2(Bucket=bucket, Prefix=prefix)
files = [obj["Key"] for obj in response.get("Contents", [])
         if obj["Key"].endswith(".csv")]

dfs = []
for file in files:
    obj = s3.get_object(Bucket=bucket, Key=file)
    df = pd.read_csv(BytesIO(obj["Body"].read()))
    dfs.append(df)

return pd.concat(dfs, ignore_index=True)

`

项目结构设计

一个规范的 Python ETL 项目应该有清晰的结构。下面是一个推荐的项目布局:

etl_project/ ├── config/ │ ├── __init__.py │ ├── settings.py # 全局配置 │ └── connections.py # 数据库连接管理 ├── extract/ │ ├── __init__.py │ ├── base.py # 抽取基类 │ ├── api_extractor.py # API 数据抽取 │ ├── db_extractor.py # 数据库数据抽取 │ └── file_extractor.py # 文件数据抽取 ├── transform/ │ ├── __init__.py │ ├── cleaners.py # 数据清洗 │ ├── aggregators.py # 聚合计算 │ └── validators.py # 数据验证 ├── load/ │ ├── __init__.py │ ├── db_loader.py # 数据库加载 │ └── file_writer.py # 文件输出 ├── utils/ │ ├── __init__.py │ ├── logger.py # 日志工具 │ └── retry.py # 重试机制 ├── tests/ │ ├── test_extract.py │ ├── test_transform.py │ └── test_load.py ├── config.yaml # 环境配置 └── requirements.txt

这种分层结构的好处是职责清晰。抽取、转换、加载三个环节完全解耦,每个模块可以独立测试和修改。如果哪天需要把数据源从 CSV 换成 API,只需要修改 extract 模块,transform 和 load 完全不受影响。

配置管理

不要把数据库密码、API Key 等敏感信息硬编码在代码里。推荐使用环境变量加 YAML 配置文件的方式:

`python

config/settings.py

import os import yaml from pathlib import Path

class Settings: def init(self, env: str = “dev”): config_path = Path(file).parent.parent / “config.yaml” with open(config_path) as f: self._config = yaml.safe_load(f)[env]

    # 敏感信息从环境变量读取,覆盖配置文件
    self.db_url = os.getenv("DB_URL", self._config.get("db_url"))
    self.api_key = os.getenv("API_KEY", self._config.get("api_key"))
    self.s3_bucket = self._config.get("s3_bucket")
    self.log_level = self._config.get("log_level", "INFO")

`

对应的 config.yaml 文件:

`yaml dev: db_url: postgresql://localhost:5432/dev_db s3_bucket: dev-data-lake log_level: DEBUG extract: batch_size: 1000 retry_times: 3

prod: db_url: postgresql://prod:5432/warehouse s3_bucket: prod-data-lake log_level: INFO extract: batch_size: 10000 retry_times: 5 `

错误处理模式

ETL 流程中错误是常态。网络超时、数据库连接断开、数据格式异常,这些情况随时可能发生。好的错误处理能让管道在遇到问题时优雅恢复,而不是直接崩溃。

重试机制

`python

utils/retry.py

import time import logging from functools import wraps

logger = logging.getLogger(name)

def retry(max_attempts: int = 3, delay: int = 5, backoff: float = 2.0): """重试装饰器,指数退避""" def decorator(func): @wraps(func) def wrapper(*args, **kwargs): last_exception = None for attempt in range(1, max_attempts + 1): try: return func(*args, **kwargs) except (ConnectionError, TimeoutError) as e: last_exception = e wait = delay * (backoff ** (attempt - 1)) logger.warning( f"第 {attempt} 次尝试失败: {e}," f"{wait} 秒后重试…" ) time.sleep(wait) raise last_exception return wrapper return decorator

使用示例

@retry(max_attempts=3, delay=2) def fetch_data_from_api(url: str): response = requests.get(url, timeout=10) response.raise_for_status() return response.json() `

数据质量检查

在转换环节加入数据验证,提前发现数据问题:

`python

transform/validators.py

import pandera as pa from pandera.typing import DataFrame

class SalesSchema(pa.DataFrameModel): """销售数据校验模型""" order_id: str = pa.Field(nullable=False) amount: float = pa.Field(ge=0, nullable=False) order_date: str = pa.Field(nullable=False) customer_id: str = pa.Field(nullable=False)

def validate_sales_data(df: pd.DataFrame) -> pd.DataFrame: """验证销售数据格式""" schema = SalesSchema validated_df = schema.validate(df) return validated_df `

日志与监控

ETL 管道通常无人值守运行,日志是排查问题的唯一线索。推荐使用 Python 标准库 logging,配合结构化日志格式:

`python

utils/logger.py

import logging import sys from datetime import datetime

def setup_logger(name: str, level: str = “INFO”): logger = logging.getLogger(name) logger.setLevel(getattr(logging, level))

handler = logging.StreamHandler(sys.stdout)
formatter = logging.Formatter(
    "[%(asctime)s] %(levelname)s - %(name)s - %(message)s",
    datefmt="%Y-%m-%d %H:%M:%S"
)
handler.setFormatter(formatter)
logger.addHandler(handler)
return logger

使用示例

logger = setup_logger(“etl_pipeline”) logger.info(“开始抽取数据,共 50000 条记录”) logger.warning(“发现 23 条缺失金额的记录,已过滤”) logger.error(“数据库连接超时,正在重试…”) `

测试 ETL 代码

ETL 代码同样需要测试。单元测试验证转换逻辑的正确性,集成测试验证数据库读写是否正常。

`python

tests/test_transform.py

import pandas as pd import pytest from transform.cleaners import clean_sales_data

def test_clean_sales_data_removes_negative_amounts(): input_df = pd.DataFrame({ “amount”: [100, -50, 200, -10, 0], “date”: [“2026-01-01”] * 5 }) result = clean_sales_data(input_df) assert len(result) == 3 # 只保留 >= 0 的记录 assert (result[“amount”] >= 0).all()

def test_clean_sales_data_handles_missing_dates(): input_df = pd.DataFrame({ “amount”: [100, 200], “date”: [“2026-01-01”, None] }) result = clean_sales_data(input_df) assert len(result) == 1 # 缺失日期的行被删除 `

完整示例:CSV 到 PostgreSQL

下面是一个完整的端到端 ETL 示例,从 CSV 文件读取销售数据,经过清洗和聚合,最终加载到 PostgreSQL 数据库。

`python

etl_pipeline.py

import pandas as pd from sqlalchemy import create_engine from config.settings import Settings from utils.logger import setup_logger from utils.retry import retry

logger = setup_logger(“etl_pipeline”)

def extract(file_path: str) -> pd.DataFrame: """步骤 1:抽取数据""" logger.info(f"从 {file_path} 抽取数据") df = pd.read_csv(file_path) logger.info(f"抽取到 {len(df)} 条记录") return df

def transform(df: pd.DataFrame) -> pd.DataFrame: """步骤 2:转换数据""" logger.info(“开始数据转换”)

# 删除缺失值
df = df.dropna(subset=["order_id", "amount", "order_date"])
logger.info(f"删除缺失值后剩余 {len(df)} 条")

# 过滤异常值
df = df[df["amount"] > 0]
logger.info(f"过滤负值后剩余 {len(df)} 条")

# 日期处理
df["order_date"] = pd.to_datetime(df["order_date"])
df["year_month"] = df["order_date"].dt.to_period("M")

# 计算衍生字段
df["total_amount"] = df["quantity"] * df["unit_price"]

# 聚合
result = df.groupby(["year_month", "category"]).agg(
    total_sales=("total_amount", "sum"),
    order_count=("order_id", "nunique"),
    avg_order_value=("total_amount", "mean")
).reset_index()

logger.info(f"转换完成,生成 {len(result)} 条聚合记录")
return result

@retry(max_attempts=3, delay=2) def load(df: pd.DataFrame, engine) -> None: """步骤 3:加载数据""" logger.info(“开始加载数据到 PostgreSQL”) df.to_sql( “monthly_sales_summary”, engine, if_exists=“replace”, index=False, method=“multi”, chunksize=1000 ) logger.info(“数据加载完成”)

def run_pipeline(): """运行完整 ETL 管道""" settings = Settings(env=“prod”) engine = create_engine(settings.db_url)

df = extract("data/raw_sales.csv")
df_transformed = transform(df)
load(df_transformed, engine)

if name == “main”: run_pipeline() `

性能优化建议

Python ETL 的性能瓶颈通常不在 Python 本身,而在 I/O 操作和数据处理方式上。以下几点可以显著提升性能:

  1. 批量处理:不要逐行处理数据,尽量使用 Pandas 的向量化操作。向量化操作比 for 循环快几十倍。

  2. 分块读取:大文件不要一次性读入内存,使用 chunksize 参数分块处理。

python chunks = pd.read_csv("huge_file.csv", chunksize=10000) for chunk in chunks: process(chunk)

  1. 选择合适的数据类型:Pandas 默认使用 int64 和 float64,如果数据范围较小,可以降级为 int32 或 float32,节省一半内存。

  2. 数据库批量写入:使用 method=“multi” 和合适的 chunksize,避免逐行 INSERT。

  3. 并行处理:对于可并行的转换任务,使用 concurrent.futures 或 Pandas 的 groupby 并行化。

小结

本文介绍了用 Python 构建 ETL 管道的核心技能。我们从 Pandas、SQLAlchemy、requests、boto3 四个核心库入手,讨论了项目结构设计、配置管理、错误处理、日志监控和测试方法,最后给出了一个完整的 CSV 到 PostgreSQL 的 ETL 示例。Python 的灵活性和丰富的生态使其成为 ETL 开发的首选语言,但也要注意内存管理和性能优化。

Summary: Python ETL 核心库、项目结构、错误处理与完整示例。