为什么要构建数据管道?

在前两篇文章中,我们手动加载数据、做清洗、分析、可视化。但如果这个过程需要每天重复做一次呢?比如每天早上都要给业务团队发一份昨天的销售报告。

手动做几件事很麻烦:

  1. 每天从不同系统导出数据
  2. 重复执行相同的清洗和分析代码
  3. 把结果发给不同的人
  4. 一旦某个环节出错,整个流程中断

数据管道(Data Pipeline)就是为了解决这些问题。它把数据处理流程自动化、标准化、可重复化。

ETL 与 ELT 架构

在构建数据管道前,需要了解两种主流架构。

ETL(Extract, Transform, Load)

数据先提取到中间层做转换,再加载到目标系统。

数据源 → 提取 → 转换 → 加载 → 数据仓库

适合场景:

  • 目标系统对数据格式有严格要求
  • 需要在加载前清洗和标准化数据
  • 传统数据仓库场景

ELT(Extract, Load, Transform)

数据先原始加载到目标系统,在目标系统内做转换。

数据源 → 提取 → 加载 → 数据湖/仓库 → 按需转换

适合场景:

  • 目标系统计算能力强(如云数据仓库)
  • 需要保留原始数据以备后续重新处理
  • 分析需求经常变化
# ETL 示例:提取 → 转换 → 加载
# Extract
def extract_from_csv(filepath):
    return pd.read_csv(filepath)

def extract_from_api(api_url, api_key):
    import requests
    headers = {'Authorization': f'Bearer {api_key}'}
    response = requests.get(api_url, headers=headers)
    return pd.DataFrame(response.json())

# Transform
def transform_sales_data(df):
    df = df.copy()
    df['order_date'] = pd.to_datetime(df['order_date'])
    df['amount'] = df['price'] * df['quantity']
    df = df.dropna(subset=['order_id'])
    # 去除异常值
    df = df[df['amount'] > 0]
    return df

# Load
def load_to_csv(df, output_path):
    df.to_csv(output_path, index=False)
    print(f'已写入 {output_path}')

def load_to_parquet(df, output_path):
    df.to_parquet(output_path, index=False)
    print(f'已写入 {output_path}')

数据提取(Extract)

从 CSV 提取

import pandas as pd

# 基础读取
df = pd.read_csv('sales_2025_01.csv')

# 处理大文件:分块读取
chunk_size = 10000
chunks = []
for chunk in pd.read_csv('large_sales.csv', chunksize=chunk_size):
    # 在每块上做初步过滤
    chunk = chunk[chunk['amount'] > 0]
    chunks.append(chunk)

df = pd.concat(chunks)

从 Excel 提取

# 读取所有工作表
xls = pd.ExcelFile('monthly_report.xlsx')
print(f'工作表: {xls.sheet_names}')

# 读取指定工作表
df_sales = pd.read_excel(xls, sheet_name='销售数据')
df_refunds = pd.read_excel(xls, sheet_name='退款数据')

# 读取指定范围(跳过表头行)
df = pd.read_excel('report.xlsx', sheet_name='Sheet1', skiprows=3)

从 API 提取

很多业务系统的数据通过 REST API 提供。

import requests
import json
from datetime import datetime, timedelta

def extract_from_api(base_url, endpoint, params=None, auth_token=None):
    headers = {}
    if auth_token:
        headers['Authorization'] = f'Bearer {auth_token}'
    
    url = f'{base_url}/{endpoint}'
    all_data = []
    page = 1
    
    while True:
        page_params = {'page': page, 'per_page': 100}
        if params:
            page_params.update(params)
        
        response = requests.get(url, headers=headers, params=page_params)
        
        if response.status_code != 200:
            print(f'API 请求失败: {response.status_code}')
            break
        
        data = response.json()
        if not data:
            break
        
        all_data.extend(data)
        page += 1
    
    return pd.DataFrame(all_data)

# 示例:获取昨天订单
yesterday = (datetime.now() - timedelta(days=1)).strftime('%Y-%m-%d')
orders = extract_from_api(
    'https://api.example.com',
    'orders',
    params={'start_date': yesterday, 'end_date': yesterday},
    auth_token='your_api_token_here'
)

从数据库提取

from sqlalchemy import create_engine, text

# 连接 PostgreSQL
pg_engine = create_engine('postgresql://user:password@localhost:5432/ecommerce')

# 连接 MySQL
mysql_engine = create_engine('mysql+pymysql://user:password@localhost:3306/ecommerce')

# 使用 SQL 查询提取数据
query = """
SELECT 
    o.order_id,
    o.customer_id,
    o.order_date,
    oi.product_id,
    oi.quantity,
    oi.unit_price,
    oi.quantity * oi.unit_price AS amount
FROM orders o
JOIN order_items oi ON o.order_id = oi.order_id
WHERE o.order_date >= :start_date
  AND o.order_date < :end_date
"""

df = pd.read_sql(
    text(query),
    pg_engine,
    params={'start_date': '2025-01-01', 'end_date': '2025-02-01'}
)

# 增量提取:只提取上次之后的增量
last_run = '2025-01-19 03:00:00'
incremental_query = f"""
SELECT * FROM orders
WHERE updated_at > '{last_run}'
"""
incremental_df = pd.read_sql(incremental_query, pg_engine)

数据转换(Transform)

转换是整个管道的核心环节。它不仅要清洗数据,还要做业务逻辑层面的计算。

class SalesTransformer:
    """销售数据转换器,封装所有转换逻辑"""
    
    def __init__(self, df):
        self.df = df.copy()
        self.validation_errors = []
    
    def clean_dates(self, date_column='order_date'):
        self.df[date_column] = pd.to_datetime(
            self.df[date_column], errors='coerce'
        )
        invalid_dates = self.df[date_column].isna().sum()
        if invalid_dates > 0:
            self.validation_errors.append(
                f'{invalid_dates} 条记录日期无效'
            )
        return self
    
    def remove_duplicates(self, subset='order_id'):
        before = len(self.df)
        self.df = self.df.drop_duplicates(subset=subset)
        after = len(self.df)
        if before > after:
            self.validation_errors.append(
                f'移除 {before - after} 条重复记录'
            )
        return self
    
    def handle_missing_values(self, strategy='drop'):
        if strategy == 'drop':
            self.df = self.df.dropna()
        elif strategy == 'fill_zero':
            self.df = self.df.fillna(0)
        elif strategy == 'fill_mean':
            numeric_cols = self.df.select_dtypes(include=[np.number]).columns
            self.df[numeric_cols] = self.df[numeric_cols].fillna(
                self.df[numeric_cols].mean()
            )
        return self
    
    def validate_data_quality(self):
        """数据质量检查"""
        checks = {
            '缺失值数量': self.df.isna().sum().sum(),
            '负金额订单': (self.df['amount'] < 0).sum() if 'amount' in self.df else 0,
            '异常数量': (self.df['quantity'] > 1000).sum() if 'quantity' in self.df else 0,
        }
        return checks
    
    def add_business_metrics(self):
        """添加业务指标"""
        self.df['order_month'] = self.df['order_date'].dt.to_period('M')
        self.df['order_weekday'] = self.df['order_date'].dt.dayofweek
        self.df['order_hour'] = self.df['order_date'].dt.hour
        return self
    
    def get_report(self):
        reports = {
            'total_records': len(self.df),
            'total_amount': self.df['amount'].sum() if 'amount' in self.df else 0,
            'unique_customers': self.df['customer_id'].nunique() if 'customer_id' in self.df else 0,
            'date_range': f"{self.df['order_date'].min()} ~ {self.df['order_date'].max()}",
            'validation_errors': self.validation_errors
        }
        return reports

# 使用示例
raw_df = pd.read_csv('raw_sales.csv')
transformer = SalesTransformer(raw_df)
clean_df = (transformer
    .clean_dates()
    .remove_duplicates()
    .handle_missing_values(strategy='drop')
    .add_business_metrics()
    .df)

print(transformer.validate_data_quality())

数据加载(Load)

写入不同目标

# 写入 CSV(适合小规模、人类可读)
clean_df.to_csv('daily_sales_report.csv', index=False)

# 写入 Parquet(列式存储,压缩率高,适合大规模数据)
clean_df.to_parquet('daily_sales.parquet', compression='snappy')

# 写入数据库
clean_df.to_sql(
    'daily_sales',
    pg_engine,
    if_exists='replace',  # 或 'append' 做增量
    index=False,
    method='multi',  # 批量插入,提升性能
    chunksize=1000
)

增量加载策略

def incremental_load(df, engine, table_name, unique_key, last_run):
    """增量加载:只插入新记录和更新已有记录"""
    
    # 读取已有数据
    existing = pd.read_sql(f'SELECT {unique_key} FROM {table_name}', engine)
    existing_keys = set(existing[unique_key])
    
    # 区分新增和更新
    new_records = df[~df[unique_key].isin(existing_keys)]
    update_records = df[df[unique_key].isin(existing_keys)]
    
    # 写入新增
    if len(new_records) > 0:
        new_records.to_sql(table_name, engine, if_exists='append', 
                          index=False, method='multi')
    
    # 更新已有
    if len(update_records) > 0:
        # 删除后重新插入
        with engine.connect() as conn:
            for _, row in update_records.iterrows():
                conn.execute(
                    f"DELETE FROM {table_name} WHERE {unique_key} = '{row[unique_key]}'"
                )
            update_records.to_sql(table_name, engine, if_exists='append',
                                 index=False, method='multi')
    
    return {
        'new': len(new_records),
        'updated': len(update_records)
    }

工作流编排

方式一:Cron 定时任务

最简单的自动化方式,适合单机任务。

# 每天凌晨 3 点运行销售报告管道
0 3 * * * cd /path/to/pipeline && python daily_sales_pipeline.py >> logs/pipeline.log 2>&1

# 每周一早上 6 点运行周报
0 6 * * 1 cd /path/to/pipeline && python weekly_report.py

# 每月 1 号凌晨 4 点运行月报
0 4 1 * * cd /path/to/pipeline && python monthly_aggregation.py

但 cron 有几个局限:没有依赖管理、没有失败重试、没有监控界面。

方式二:Apache Airflow

Airflow 是目前最流行的开源工作流编排工具。

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.bash import BashOperator
from datetime import datetime, timedelta
import pandas as pd

default_args = {
    'owner': 'data_team',
    'depends_on_past': False,
    'email_on_failure': True,
    'email_on_retry': False,
    'retries': 3,
    'retry_delay': timedelta(minutes=5),
}

dag = DAG(
    'daily_sales_report',
    default_args=default_args,
    description='每日销售报表管道',
    schedule='0 6 * * *',  # 每天早上 6 点
    start_date=datetime(2025, 1, 1),
    catchup=False,
    tags=['sales', 'daily'],
)

def extract_sales(**context):
    """从数据库提取昨日销售数据"""
    execution_date = context['execution_date']
    yesterday = execution_date - timedelta(days=1)
    date_str = yesterday.strftime('%Y-%m-%d')
    
    import sqlalchemy
    engine = sqlalchemy.create_engine('postgresql://...')
    
    df = pd.read_sql(
        f"SELECT * FROM orders WHERE date(created_at) = '{date_str}'",
        engine
    )
    
    # 保存为临时文件,供下一个任务使用
    df.to_parquet(f'/tmp/sales_raw_{date_str}.parquet')
    return f'/tmp/sales_raw_{date_str}.parquet'

def transform_sales(**context):
    """清洗和转换数据"""
    ti = context['ti']
    input_path = ti.xcom_pull(task_ids='extract')
    date_str = context['execution_date'].strftime('%Y-%m-%d')
    
    df = pd.read_parquet(input_path)
    
    # 应用转换逻辑
    df = df.dropna(subset=['order_id'])
    df = df[df['amount'] > 0]
    df['order_hour'] = pd.to_datetime(df['created_at']).dt.hour
    
    output_path = f'/tmp/sales_clean_{date_str}.parquet'
    df.to_parquet(output_path)
    return output_path

def aggregate_sales(**context):
    """计算每日指标"""
    ti = context['ti']
    input_path = ti.xcom_pull(task_ids='transform')
    
    df = pd.read_parquet(input_path)
    
    metrics = {
        'date': context['execution_date'].strftime('%Y-%m-%d'),
        'total_revenue': df['amount'].sum(),
        'order_count': len(df),
        'avg_order_value': df['amount'].mean(),
        'unique_customers': df['customer_id'].nunique(),
        'top_category': df.groupby('category')['amount'].sum().idxmax(),
    }
    
    metrics_df = pd.DataFrame([metrics])
    output_path = f'/tmp/daily_metrics_{context["execution_date"].strftime("%Y-%m-%d")}.parquet'
    metrics_df.to_parquet(output_path)
    return output_path

def load_to_database(**context):
    """将结果加载到报表数据库"""
    ti = context['ti']
    input_path = ti.xcom_pull(task_ids='aggregate')
    
    df = pd.read_parquet(input_path)
    engine = sqlalchemy.create_engine('postgresql://reporting:password@reporting-db:5432/reports')
    df.to_sql('daily_sales', engine, if_exists='append', index=False)

def send_report(**context):
    """发送报表通知"""
    ti = context['ti']
    metrics_path = ti.xcom_pull(task_ids='aggregate')
    df = pd.read_parquet(metrics_path)
    metrics = df.iloc[0].to_dict()
    
    message = f"""
    【每日销售报告】
    日期: {metrics['date']}
    营收: ¥{metrics['total_revenue']:,.2f}
    订单数: {metrics['order_count']}
    平均客单价: ¥{metrics['avg_order_value']:.2f}
    活跃客户: {metrics['unique_customers']}
    """
    
    # 发送到钉钉/飞书/企业微信
    import requests
    webhook_url = "https://hooks.example.com/report"
    requests.post(webhook_url, json={"msgtype": "text", "text": {"content": message}})
    
    return "Report sent"

# 定义任务
t_extract = PythonOperator(
    task_id='extract',
    python_callable=extract_sales,
    dag=dag,
)

t_transform = PythonOperator(
    task_id='transform',
    python_callable=transform_sales,
    dag=dag,
)

t_aggregate = PythonOperator(
    task_id='aggregate',
    python_callable=aggregate_sales,
    dag=dag,
)

t_load = PythonOperator(
    task_id='load',
    python_callable=load_to_database,
    dag=dag,
)

t_send = PythonOperator(
    task_id='send_report',
    python_callable=send_report,
    dag=dag,
)

# 设置依赖关系
t_extract >> t_transform >> t_aggregate >> [t_load, t_send]

Airflow DAG 的核心概念

概念 说明 类比
DAG 有向无环图,定义任务依赖关系 流程图
Operator 单个任务的操作类型 步骤
Task Operator 的实例化 具体动作
Schedule 执行周期表达式 定时器
XCom 任务间数据传递 中间变量
Sensor 等待外部条件满足 触发器

构建完整的每日销售报表管道

下面是一个端到端的管道脚本示例,包含监控和告警。

#!/usr/bin/env python3
"""
daily_sales_pipeline.py - 每日销售数据管道
功能: 提取昨日销售数据 → 清洗转换 → 生成报表 → 发送通知
"""

import pandas as pd
import numpy as np
from datetime import datetime, timedelta
import logging
import sys
from pathlib import Path

# 日志配置
logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
    handlers=[
        logging.FileHandler(f'logs/pipeline_{datetime.now().strftime("%Y%m%d")}.log'),
        logging.StreamHandler()
    ]
)
logger = logging.getLogger(__name__)

class PipelineMonitor:
    """管道监控:记录执行状态和耗时"""
    
    def __init__(self, pipeline_name):
        self.pipeline_name = pipeline_name
        self.start_time = datetime.now()
        self.steps = []
    
    def step_start(self, step_name):
        self.current_step = step_name
        self.step_start_time = datetime.now()
        logger.info(f'[开始] {step_name}')
    
    def step_end(self, status='success'):
        elapsed = (datetime.now() - self.step_start_time).total_seconds()
        self.steps.append({
            'step': self.current_step,
            'status': status,
            'elapsed_seconds': elapsed
        })
        logger.info(f'[完成] {self.current_step} ({elapsed:.1f}s) - {status}')
    
    def summary(self):
        total_time = (datetime.now() - self.start_time).total_seconds()
        failed_steps = [s for s in self.steps if s['status'] == 'failed']
        
        return {
            'pipeline': self.pipeline_name,
            'total_time': total_time,
            'steps': self.steps,
            'success': len(failed_steps) == 0
        }

def run_daily_pipeline():
    monitor = PipelineMonitor('daily_sales_pipeline')
    
    try:
        # 步骤 1: 提取
        monitor.step_start('extract')
        yesterday = (datetime.now() - timedelta(days=1)).strftime('%Y-%m-%d')
        
        # 从多个源提取
        orders_df = extract_from_database(yesterday)
        events_df = extract_from_api(
            'https://api.example.com', 
            'events',
            params={'date': yesterday}
        )
        monitor.step_end()
        
        # 步骤 2: 转换
        monitor.step_start('transform')
        transformer = SalesTransformer(orders_df)
        clean_orders = (transformer
            .clean_dates()
            .remove_duplicates()
            .handle_missing_values()
            .add_business_metrics()
            .df)
        
        # 计算每日 KPI
        daily_kpis = {
            'date': yesterday,
            'revenue': clean_orders['amount'].sum(),
            'orders': len(clean_orders),
            'avg_order_value': clean_orders['amount'].mean(),
            'unique_customers': clean_orders['customer_id'].nunique(),
            'new_customers': len(clean_orders[clean_orders['is_new_customer'] == True]) 
                if 'is_new_customer' in clean_orders else 0,
        }
        logger.info(f'KPI: {daily_kpis}')
        monitor.step_end()
        
        # 步骤 3: 加载
        monitor.step_start('load')
        # 保存到报表数据库
        save_to_database(daily_kpis)
        # 保存到数据湖
        clean_orders.to_parquet(
            f'data_lake/sales/year={yesterday[:4]}/month={yesterday[5:7]}/day={yesterday[8:10]}/sales.parquet',
            compression='snappy'
        )
        monitor.step_end()
        
        # 步骤 4: 告警检查
        monitor.step_start('alerting')
        alerts = check_alerts(daily_kpis)
        for alert in alerts:
            send_alert(alert)
            logger.warning(f'告警: {alert}')
        monitor.step_end()
        
        summary = monitor.summary()
        logger.info(f'管道执行完成,耗时 {summary["total_time"]:.1f}s')
        return summary
        
    except Exception as e:
        monitor.step_end('failed')
        logger.error(f'管道执行失败: {e}', exc_info=True)
        send_alert(f'销售管道执行失败: {e}')
        raise

if __name__ == '__main__':
    run_daily_pipeline()

告警与监控

一个健壮的管道必须包含告警机制。

def check_alerts(kpis, thresholds=None):
    """检查指标是否超过阈值"""
    if thresholds is None:
        thresholds = {
            'revenue_min': 10000,
            'order_drop_pct': 0.3,  # 订单量相比7日均值下降超过30%
            'error_rate_max': 0.05,
        }
    
    alerts = []
    
    # 收入异常
    if kpis['revenue'] < thresholds['revenue_min']:
        alerts.append({
            'level': 'critical',
            'metric': 'revenue',
            'message': f"日收入异常偏低: ¥{kpis['revenue']:,.2f},低于阈值 ¥{thresholds['revenue_min']:,.2f}"
        })
    
    # 订单量骤降(需要和历史对比)
    # 此处简化处理,实际应从数据库读取历史数据
    if kpis['orders'] < 50:
        alerts.append({
            'level': 'warning',
            'metric': 'orders',
            'message': f"订单量偏低: {kpis['orders']} 单"
        })
    
    return alerts

def send_alert(alert, channels=['slack', 'email']):
    """多渠道发送告警"""
    for channel in channels:
        if channel == 'slack':
            # 发送到 Slack Webhook
            pass
        elif channel == 'email':
            # 发送邮件
            pass
        elif channel == 'sms':
            # 发送短信(仅 critical 级别)
            if alert['level'] == 'critical':
                pass

可重复性最佳实践

让管道真正可靠,需要养成以下习惯:

1. 代码版本化

# 所有管道代码使用 git 管理
git init && git add -A && git commit -m "feat: initial data pipeline"

# 每次修改走 PR 流程
git checkout -b feature/add-return-rate-metric
# ... 修改代码 ...
git commit -m "feat: add return rate to daily KPI"
# 创建 PR,经过审核后合并

2. 依赖锁定

# 使用 requirements.txt 或 poetry 锁定依赖版本
pip freeze > requirements.txt

3. 幂等性设计

管道应该可以多次安全运行,每次产生相同结果。

# 好的做法:先清理再写入
def load_with_idempotency(df, engine, table, date_column, date_value):
    """幂等加载:先删除当天数据再写入"""
    with engine.connect() as conn:
        conn.execute(f"DELETE FROM {table} WHERE {date_column} = '{date_value}'")
    df.to_sql(table, engine, if_exists='append', index=False)

4. 数据质量断言

def assert_data_quality(df):
    """管道中的断言关卡"""
    assert len(df) > 0, "数据集为空"
    assert df['amount'].sum() > 0, "总金额为负"
    assert df['order_id'].is_unique, "存在重复订单"
    assert not df['customer_id'].isna().any(), "存在缺失客户ID"
    print("✅ 所有数据质量检查通过")

小结

数据管道是数据分析从"一次性探索"走向"生产化"的关键一步。本篇文章涵盖了:

  • ETL 与 ELT 两种架构模式的差异和适用场景
  • 从 CSV、Excel、API、数据库等多种源提取数据的方法
  • 数据清洗和转换的工程化封装
  • 以 Airflow 为核心的工作流编排
  • 管道的监控、告警和可重复性实践

掌握了这些技能,你就可以把日常的数据处理工作自动化,把精力集中在更有价值的分析上。

在下一篇文章中,我们将探讨当数据量大到单机无法处理时,如何用 Spark 等大数据工具来应对。