8 minutes
数据管道与自动化
为什么要构建数据管道?
在前两篇文章中,我们手动加载数据、做清洗、分析、可视化。但如果这个过程需要每天重复做一次呢?比如每天早上都要给业务团队发一份昨天的销售报告。
手动做几件事很麻烦:
- 每天从不同系统导出数据
- 重复执行相同的清洗和分析代码
- 把结果发给不同的人
- 一旦某个环节出错,整个流程中断
数据管道(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 等大数据工具来应对。