5 minutes
Python ETL 实践
为什么用 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 操作和数据处理方式上。以下几点可以显著提升性能:
-
批量处理:不要逐行处理数据,尽量使用 Pandas 的向量化操作。向量化操作比 for 循环快几十倍。
-
分块读取:大文件不要一次性读入内存,使用 chunksize 参数分块处理。
python chunks = pd.read_csv("huge_file.csv", chunksize=10000) for chunk in chunks: process(chunk)
-
选择合适的数据类型:Pandas 默认使用 int64 和 float64,如果数据范围较小,可以降级为 int32 或 float32,节省一半内存。
-
数据库批量写入:使用 method=“multi” 和合适的 chunksize,避免逐行 INSERT。
-
并行处理:对于可并行的转换任务,使用 concurrent.futures 或 Pandas 的 groupby 并行化。
小结
本文介绍了用 Python 构建 ETL 管道的核心技能。我们从 Pandas、SQLAlchemy、requests、boto3 四个核心库入手,讨论了项目结构设计、配置管理、错误处理、日志监控和测试方法,最后给出了一个完整的 CSV 到 PostgreSQL 的 ETL 示例。Python 的灵活性和丰富的生态使其成为 ETL 开发的首选语言,但也要注意内存管理和性能优化。
Summary: Python ETL 核心库、项目结构、错误处理与完整示例。