9 minutes
实战:构建完整 ETL 数据管道
项目概述
经过前面 14 篇文章的学习,你已经掌握了 ETL 的各个核心环节:数据抽取、转换、加载、SQL 技巧、Python 开发、现代工具链和工作流编排。现在是时候把它们全部串联起来了。
本章是一个端到端的实战项目。我们将为一个电商平台构建完整的数据管道,从原始数据采集到业务指标计算,全部自动化运行。
场景设定
假设你在一家中型电商公司工作。公司有以下几个数据源:
- MySQL 业务数据库:存储订单、用户、商品信息
- CSV 文件:存储在 S3 上的每日营销费用数据
- REST API:第三方物流服务商的运单状态接口
你的任务是把这些数据汇集到 PostgreSQL 数据仓库中,清洗转换后生成业务报表,并每天自动运行。
技术栈
| 环节 | 工具 | 用途 |
|---|---|---|
| 数据抽取 | Python + SQLAlchemy + requests | 从数据库、API、文件抽取数据 |
| 数据加载 | Python + PostgreSQL | 加载到数据仓库 |
| 数据转换 | dbt | 在仓库内完成转换 |
| 工作流编排 | Airflow | 调度整个流程 |
| 数据质量 | dbt test + Great Expectations | 数据验证 |
| 监控告警 | Airflow + Slack | 失败通知 |
第一步:项目结构
ecommerce-etl/
├── dags/ # Airflow DAG 定义
│ └── ecommerce_pipeline.py
├── dbt_project/ # dbt 项目
│ ├── dbt_project.yml
│ ├── models/
│ │ ├── staging/
│ │ ├── intermediate/
│ │ └── marts/
│ └── tests/
├── scripts/ # Python ETL 脚本
│ ├── extract/
│ ├── transform/
│ └── load/
├── config/
│ ├── config.yaml
│ └── connections.yaml
├── docker-compose.yaml # Airflow + PostgreSQL
├── requirements.txt
└── README.md
[/code]
## 第二步:搭建基础设施
### Docker 环境
version: ‘3.8’
services: postgres: image: postgres:15 environment: POSTGRES_USER: warehouse POSTGRES_PASSWORD: warehouse_pass POSTGRES_DB: ecommerce_wh ports: - “5432:5432” volumes: - postgres_data:/var/lib/postgresql/data - ./init_sql:/docker-entrypoint-initdb.d
airflow-webserver: image: apache/airflow:2.10.0 command: webserver ports: - “8080:8080” depends_on: - postgres environment: AIRFLOW__CORE__EXECUTOR: LocalExecutor AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres/airflow volumes: - ./dags:/opt/airflow/dags - ./scripts:/opt/airflow/scripts - ./config:/opt/airflow/config
airflow-scheduler: image: apache/airflow:2.10.0 command: scheduler depends_on: - postgres environment: AIRFLOW__CORE__EXECUTOR: LocalExecutor AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres/airflow volumes: - ./dags:/opt/airflow/dags - ./scripts:/opt/airflow/scripts - ./config:/opt/airflow/config
volumes: postgres_data: [/code]
初始化数据仓库 Schema
-- init_sql/01_create_schemas.sql
CREATE SCHEMA IF NOT EXISTS raw_data;
CREATE SCHEMA IF NOT EXISTS staging;
CREATE SCHEMA IF NOT EXISTS dw;
CREATE SCHEMA IF NOT EXISTS metrics;
-- init_sql/02_create_tables.sql
-- 原始数据表
CREATE TABLE raw_data.orders (
order_id VARCHAR(50),
customer_id VARCHAR(50),
product_id VARCHAR(50),
amount NUMERIC(12,2),
quantity INTEGER,
status VARCHAR(20),
order_date TIMESTAMP,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
CREATE TABLE raw_data.customers (
customer_id VARCHAR(50) PRIMARY KEY,
name VARCHAR(200),
email VARCHAR(200),
signup_date DATE,
city VARCHAR(100),
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
CREATE TABLE raw_data.marketing_spend (
spend_date DATE,
channel VARCHAR(50),
amount NUMERIC(12,2),
source_file VARCHAR(200),
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
-- 指标表
CREATE TABLE metrics.daily_overview (
report_date DATE PRIMARY KEY,
total_revenue NUMERIC(14,2),
order_count INTEGER,
customer_count INTEGER,
avg_order_value NUMERIC(10,2),
marketing_spend NUMERIC(12,2),
revenue_per_spend NUMERIC(10,2),
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
[/code]
## 第三步:数据抽取层
### 从 MySQL 抽取订单数据
scripts/extract/extract_mysql.py
import pandas as pd from sqlalchemy import create_engine, text from datetime import datetime, timedelta import logging
logger = logging.getLogger(name)
class MySQLExtractor: “““从 MySQL 业务数据库抽取数据”””
def __init__(self, connection_string: str):
self.engine = create_engine(connection_string)
def extract_orders(self, start_date: str, end_date: str) -> pd.DataFrame:
"""抽取指定日期范围的订单"""
query = text("""
SELECT
o.id AS order_id,
o.customer_id,
oi.product_id,
oi.unit_price * oi.quantity AS amount,
oi.quantity,
o.status,
o.created_at AS order_date
FROM orders o
JOIN order_items oi ON o.id = oi.order_id
WHERE o.created_at >= :start_date
AND o.created_at < :end_date
""")
with self.engine.connect() as conn:
df = pd.read_sql(query, conn, params={
"start_date": start_date,
"end_date": end_date
})
logger.info(f"从 MySQL 抽取了 {len(df)} 条订单记录")
return df
def extract_customers(self) -> pd.DataFrame:
"""全量抽取客户数据(客户表较小)"""
query = text("SELECT * FROM customers")
with self.engine.connect() as conn:
df = pd.read_sql(query, conn)
logger.info(f"从 MySQL 抽取了 {len(df)} 条客户记录")
return df
[/code]
从 S3 读取营销费用 CSV
# scripts/extract/extract_s3.py
import pandas as pd
import boto3
from io import BytesIO, StringIO
import logging
logger = logging.getLogger(__name__)
class S3Extractor:
"""从 S3 读取 CSV 文件"""
def __init__(self, bucket: str, prefix: str):
self.bucket = bucket
self.prefix = prefix
self.s3 = boto3.client("s3")
def extract_marketing_spend(self, date: str) -> pd.DataFrame:
"""读取指定日期的营销费用文件"""
key = f"{self.prefix}marketing_{date}.csv"
try:
obj = self.s3.get_object(Bucket=self.bucket, Key=key)
df = pd.read_csv(BytesIO(obj["Body"].read()))
df["source_file"] = key
logger.info(f"从 S3 读取了 {len(df)} 条营销记录")
return df
except self.s3.exceptions.NoSuchKey:
logger.warning(f"S3 文件不存在: {key},返回空数据")
return pd.DataFrame()
[/code]
### 调用物流 API
scripts/extract/extract_api.py
import requests import pandas as pd import logging from typing import Optional
logger = logging.getLogger(name)
class LogisticsAPIExtractor: “““从物流服务商 API 抽取运单状态”””
def __init__(self, api_url: str, api_key: str):
self.api_url = api_url
self.session = requests.Session()
self.session.headers.update({"Authorization": f"Bearer {api_key}"})
def extract_tracking(self, date: str) -> pd.DataFrame:
"""抽取指定日期的物流跟踪数据"""
all_records = []
page = 1
while True:
try:
response = self.session.get(
f"{self.api_url}/tracking",
params={"date": date, "page": page, "per_page": 1000},
timeout=30
)
response.raise_for_status()
data = response.json()
all_records.extend(data["items"])
if page >= data["total_pages"]:
break
page += 1
except requests.RequestException as e:
logger.error(f"API 请求失败: {e}")
raise
df = pd.DataFrame(all_records)
logger.info(f"从物流 API 抽取了 {len(df)} 条运单记录")
return df
[/code]
第四步:数据转换层(Python)
数据清洗
# scripts/transform/cleaners.py
import pandas as pd
import numpy as np
import logging
logger = logging.getLogger(__name__)
def clean_orders(df: pd.DataFrame) -> pd.DataFrame:
"""清洗订单数据"""
logger.info(f"清洗前订单数: {len(df)}")
# 删除完全重复的行
df = df.drop_duplicates(subset=["order_id", "product_id"])
# 处理缺失值
df["amount"] = df["amount"].fillna(0)
df["quantity"] = df["quantity"].fillna(1).astype(int)
# 过滤异常值
df = df[df["amount"] >= 0]
df = df[df["quantity"] >= 0]
# 标准化状态字段
status_mapping = {
"paid": "completed",
"shipped": "completed",
"delivered": "completed",
"cancelled": "cancelled",
"refunded": "refunded",
"pending": "pending"
}
df["status"] = df["status"].str.lower().map(status_mapping).fillna("unknown")
logger.info(f"清洗后订单数: {len(df)}")
return df
def clean_customers(df: pd.DataFrame) -> pd.DataFrame:
"""清洗客户数据"""
# 去除无效邮箱
df = df[df["email"].notna() & df["email"].str.contains("@")]
# 填充缺失城市
df["city"] = df["city"].fillna("未知")
# 去重,保留最新记录
df = df.sort_values("signup_date", ascending=False)
df = df.drop_duplicates(subset=["customer_id"])
return df
[/code]
### 业务指标计算
scripts/transform/metrics.py
import pandas as pd import logging
logger = logging.getLogger(name)
def calculate_daily_overview( orders_df: pd.DataFrame, marketing_df: pd.DataFrame ) -> pd.DataFrame: “““计算每日业务概览指标””” # 计算订单指标 orders_df[“order_date”] = pd.to_datetime(orders_df[“order_date”]).dt.date
daily_orders = orders_df.groupby("order_date").agg(
total_revenue=("amount", "sum"),
order_count=("order_id", "nunique"),
customer_count=("customer_id", "nunique"),
avg_order_value=("amount", "mean")
).reset_index()
# 计算营销费用
if not marketing_df.empty:
marketing_df["spend_date"] = pd.to_datetime(
marketing_df["spend_date"]
).dt.date
daily_marketing = marketing_df.groupby("spend_date").agg(
marketing_spend=("amount", "sum")
).reset_index()
else:
daily_marketing = pd.DataFrame(
columns=["spend_date", "marketing_spend"]
)
# 合并指标
result = daily_orders.merge(
daily_marketing,
left_on="order_date",
right_on="spend_date",
how="left"
)
result["marketing_spend"] = result["marketing_spend"].fillna(0)
# 计算投入产出比
result["revenue_per_spend"] = np.where(
result["marketing_spend"] > 0,
result["total_revenue"] / result["marketing_spend"],
0
)
return result
[/code]
第五步:数据加载层
# scripts/load/load_to_warehouse.py
from sqlalchemy import create_engine, text
import pandas as pd
import logging
logger = logging.getLogger(__name__)
class WarehouseLoader:
"""加载数据到 PostgreSQL 数据仓库"""
def __init__(self, connection_string: str):
self.engine = create_engine(connection_string)
def load_raw_data(self, table: str, df: pd.DataFrame, if_exists: str = "append"):
"""加载原始数据"""
if df.empty:
logger.warning(f"数据为空,跳过加载到 {table}")
return
df.to_sql(
table,
self.engine,
schema="raw_data",
if_exists=if_exists,
index=False,
method="multi",
chunksize=5000
)
logger.info(f"成功加载 {len(df)} 条数据到 raw_data.{table}")
def load_metrics(self, df: pd.DataFrame):
"""加载指标数据(UPSERT 模式)"""
if df.empty:
return
# 使用临时表 + MERGE 实现 UPSERT
df.to_sql(
"tmp_daily_overview",
self.engine,
schema="metrics",
if_exists="replace",
index=False,
method="multi"
)
with self.engine.connect() as conn:
conn.execute(text("""
MERGE INTO metrics.daily_overview AS t
USING metrics.tmp_daily_overview AS s
ON t.report_date = s.report_date
WHEN MATCHED THEN
UPDATE SET
total_revenue = s.total_revenue,
order_count = s.order_count,
customer_count = s.customer_count,
avg_order_value = s.avg_order_value,
marketing_spend = s.marketing_spend,
revenue_per_spend = s.revenue_per_spend,
updated_at = CURRENT_TIMESTAMP
WHEN NOT MATCHED THEN
INSERT (report_date, total_revenue, order_count,
customer_count, avg_order_value,
marketing_spend, revenue_per_spend)
VALUES (s.report_date, s.total_revenue, s.order_count,
s.customer_count, s.avg_order_value,
s.marketing_spend, s.revenue_per_spend)
"""))
conn.commit()
logger.info(f"成功加载 {len(df)} 条指标数据")
[/code]
## 第六步:dbt 转换模型
### dbt 项目配置
dbt_project/dbt_project.yml
name: ecommerce_etl version: 1.0.0 profile: ecommerce
model-paths: [“models”] test-paths: [“tests”] macro-paths: [“macros”]
models: ecommerce_etl: staging: materialized: view +schema: staging intermediate: materialized: ephemeral marts: materialized: table +schema: dw [/code]
Staging 模型
-- dbt_project/models/staging/stg_orders.sql
SELECT
order_id,
customer_id,
product_id,
amount,
quantity,
status,
order_date::DATE AS order_date,
CURRENT_DATE AS load_date
FROM raw_data.orders
WHERE order_date >= CURRENT_DATE - INTERVAL '2 days'
-- dbt_project/models/staging/stg_customers.sql
SELECT DISTINCT
customer_id,
name AS customer_name,
email,
signup_date,
city
FROM raw_data.customers
[/code]
### Intermediate 模型
– dbt_project/models/intermediate/int_customer_orders.sql SELECT c.customer_id, c.customer_name, c.city, COUNT(o.order_id) AS total_orders, SUM(o.amount) AS total_spent, MIN(o.order_date) AS first_order_date, MAX(o.order_date) AS last_order_date FROM {{ ref(‘stg_customers’) }} c LEFT JOIN {{ ref(‘stg_orders’) }} o ON c.customer_id = o.customer_id WHERE o.status = ‘completed’ GROUP BY 1, 2, 3 [/code]
Marts 模型(业务主题表)
-- dbt_project/models/marts/dim_customer.sql
SELECT
customer_id,
customer_name,
email,
city,
signup_date,
CASE
WHEN total_spent >= 10000 THEN 'VIP'
WHEN total_spent >= 5000 THEN 'Gold'
WHEN total_spent >= 1000 THEN 'Silver'
ELSE 'Regular'
END AS customer_tier,
total_orders,
total_spent,
first_order_date,
last_order_date,
CURRENT_DATE - last_order_date AS days_since_last_order
FROM {{ ref('int_customer_orders') }}
-- dbt_project/models/marts/fct_daily_sales.sql
SELECT
order_date,
product_id,
COUNT(DISTINCT order_id) AS order_count,
SUM(quantity) AS units_sold,
SUM(amount) AS revenue,
COUNT(DISTINCT customer_id) AS unique_customers
FROM {{ ref('stg_orders') }}
WHERE status = 'completed'
GROUP BY 1, 2
[/code]
### dbt 测试
– dbt_project/models/marts/schema.yml version: 2
models:
-
name: dim_customer description: “客户维度表,包含客户分层信息” columns:
- name: customer_id
tests:
- unique
- not_null
- name: customer_tier
tests:
- accepted_values: values: [‘VIP’, ‘Gold’, ‘Silver’, ‘Regular’]
- name: customer_id
tests:
-
name: fct_daily_sales description: “每日销售事实表” columns:
- name: order_date
tests:
- not_null
- name: revenue
tests:
- not_null
- dbt_utils.accepted_range: min_value: 0 [/code]
- name: order_date
tests:
第七步:Airflow 编排
终于到了将所有环节串起来的时刻。下面是一个完整的 Airflow DAG,它调用了我们前面编写的所有模块:
# dags/ecommerce_pipeline.py
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.bash import BashOperator
from airflow.providers.postgres.operators.postgres import PostgresOperator
from airflow.models import Variable
from datetime import datetime, timedelta
import os
import sys
sys.path.append("/opt/airflow/scripts")
from extract.extract_mysql import MySQLExtractor
from extract.extract_s3 import S3Extractor
from extract.extract_api import LogisticsAPIExtractor
from transform.cleaners import clean_orders, clean_customers
from transform.metrics import calculate_daily_overview
from load.load_to_warehouse import WarehouseLoader
default_args = {
"owner": "data_team",
"retries": 2,
"retry_delay": timedelta(minutes=5),
"email_on_failure": True,
"email": ["data@company.com"],
}
def extract_all(**context):
"""从所有数据源抽取数据"""
exec_date = context["ds"]
config = Variable.get("etl_config", deserialize_json=True)
# 1. 从 MySQL 抽取
mysql_ext = MySQLExtractor(config["mysql_conn"])
orders = mysql_ext.extract_orders(exec_date, context["next_ds"])
customers = mysql_ext.extract_customers()
# 保存到临时文件供后续任务使用
orders.to_parquet(f"/tmp/orders_{exec_date}.parquet")
customers.to_parquet(f"/tmp/customers_{exec_date}.parquet")
# 2. 从 S3 抽取营销费用
s3_ext = S3Extractor(config["s3_bucket"], config["s3_prefix"])
marketing = s3_ext.extract_marketing_spend(exec_date)
marketing.to_parquet(f"/tmp/marketing_{exec_date}.parquet")
# 3. 从 API 抽取物流数据
api_ext = LogisticsAPIExtractor(
config["logistics_api_url"],
config["logistics_api_key"]
)
tracking = api_ext.extract_tracking(exec_date)
tracking.to_parquet(f"/tmp/tracking_{exec_date}.parquet")
return {
"orders_count": len(orders),
"customers_count": len(customers),
"marketing_count": len(marketing),
"tracking_count": len(tracking)
}
def transform_all(**context):
"""转换所有数据"""
exec_date = context["ds"]
orders = pd.read_parquet(f"/tmp/orders_{exec_date}.parquet")
customers = pd.read_parquet(f"/tmp/customers_{exec_date}.parquet")
marketing = pd.read_parquet(f"/tmp/marketing_{exec_date}.parquet")
# 清洗
orders_clean = clean_orders(orders)
customers_clean = clean_customers(customers)
# 计算指标
metrics = calculate_daily_overview(orders_clean, marketing)
# 保存清洗后数据
orders_clean.to_parquet(f"/tmp/orders_clean_{exec_date}.parquet")
customers_clean.to_parquet(f"/tmp/customers_clean_{exec_date}.parquet")
metrics.to_parquet(f"/tmp/metrics_{exec_date}.parquet")
def load_all(**context):
"""加载所有数据到数据仓库"""
exec_date = context["ds"]
config = Variable.get("etl_config", deserialize_json=True)
loader = WarehouseLoader(config["warehouse_conn"])
orders = pd.read_parquet(f"/tmp/orders_clean_{exec_date}.parquet")
customers = pd.read_parquet(f"/tmp/customers_clean_{exec_date}.parquet")
metrics = pd.read_parquet(f"/tmp/metrics_{exec_date}.parquet")
loader.load_raw_data("orders", orders, if_exists="append")
loader.load_raw_data("customers", customers, if_exists="replace")
loader.load_metrics(metrics)
def send_slack_notification(**context):
"""发送完成通知"""
ti = context["ti"]
stats = ti.xcom_pull(task_ids="extract_data")
from slack_sdk.webhook import WebhookClient
webhook = WebhookClient(Variable.get("slack_webhook_url"))
message = (
f":white_check_mark: ETL 管道执行完成\n"
f"日期: {context['ds']}\n"
f"订单数: {stats['orders_count']}\n"
f"客户数: {stats['customers_count']}\n"
f"营销记录: {stats['marketing_count']}"
)
webhook.send(text=message)
with DAG(
dag_id="ecommerce_daily_etl",
default_args=default_args,
start_date=datetime(2026, 4, 1),
schedule="0 4 * * *",
catchup=True,
tags=["ecommerce", "etl"],
description="电商平台每日 ETL 数据管道"
) as dag:
extract = PythonOperator(
task_id="extract_data",
python_callable=extract_all,
)
transform = PythonOperator(
task_id="transform_data",
python_callable=transform_all,
)
load = PythonOperator(
task_id="load_data",
python_callable=load_all,
)
# 清理临时文件
cleanup = BashOperator(
task_id="cleanup_temp_files",
bash_command="rm -f /tmp/*_{{ ds }}.parquet",
)
# dbt 转换
dbt_run = BashOperator(
task_id="dbt_run",
bash_command="cd /opt/airflow/dbt_project && dbt run",
)
dbt_test = BashOperator(
task_id="dbt_test",
bash_command="cd /opt/airflow/dbt_project && dbt test",
)
notification = PythonOperator(
task_id="send_notification",
python_callable=send_slack_notification,
)
# 数据质量检查(在加载后执行 SQL 验证)
quality_check = PostgresOperator(
task_id="quality_check",
postgres_conn_id="warehouse",
sql="""
DO
DECLARE
revenue_diff NUMERIC;
BEGIN
-- 验证 raw_data 和 dw 的金额是否一致
SELECT ABS(
(SELECT COALESCE(SUM(amount), 0)
FROM raw_data.orders
WHERE order_date::DATE = CURRENT_DATE - 1)
-
(SELECT COALESCE(total_revenue, 0)
FROM metrics.daily_overview
WHERE report_date = CURRENT_DATE - 1)
) INTO revenue_diff;
IF revenue_diff > 0.01 THEN
RAISE EXCEPTION '数据不一致: 差异 %', revenue_diff;
END IF;
END ;
""",
)
# 设置依赖关系
extract >> transform >> load >> dbt_run >> dbt_test >> quality_check >> notification
load >> cleanup
[/code]
## 第八步:数据质量监控
在 ETL 管道中加入自动化的数据质量检查,是生产环境必不可少的环节。除了 dbt 的测试和 Airflow 中的 SQL 校验,我们还可以使用 Great Expectations 做更全面的数据验证:
quality/expectations.py
import great_expectations as ge import pandas as pd
def validate_orders(df: pd.DataFrame) -> bool: “““验证订单数据质量””” ge_df = ge.from_pandas(df)
# 基本期望
ge_df.expect_column_values_to_not_be_null("order_id")
ge_df.expect_column_values_to_be_between("amount", 0, 1000000)
ge_df.expect_column_values_to_be_in_set(
"status", ["completed", "pending", "cancelled", "refunded"]
)
# 统计期望
ge_df.expect_table_row_count_to_be_between(100, 100000)
# 执行验证
results = ge_df.validate()
if not results["success"]:
failures = [
r["expectation_config"]["expectation_type"]
for r in results["results"]
if not r["success"]
]
raise ValueError(f"数据质量检查失败: {failures}")
return True
[/code]
运行项目
# 1. 启动基础设施
docker compose up -d
# 2. 初始化 Airflow
docker compose exec airflow-webserver airflow db init
docker compose exec airflow-webserver airflow users create \
--username admin --password admin --role Admin \
--email admin@example.com --firstname Admin --lastname User
# 3. 在 Airflow Web UI (http://localhost:8080) 中:
# - 添加 PostgreSQL Connection(conn_id: warehouse)
# - 添加变量 etl_config(JSON 格式的配置)
# - 手动触发 DAG 测试运行
# 4. 查看 DAG 运行状态
docker compose exec airflow-webserver airflow dags list
docker compose exec airflow-webserver airflow tasks list ecommerce_daily_etl
[/code]
### 配置示例(Airflow Variable)
{ “mysql_conn”: “mysql+pymysql://user:pass@mysql-host:3306/ecommerce”, “warehouse_conn”: “postgresql://warehouse:warehouse_pass@postgres:5432/ecommerce_wh”, “s3_bucket”: “ecommerce-marketing-data”, “s3_prefix”: “daily_spend/”, “logistics_api_url”: “https://api.logistics.com/v1", “logistics_api_key”: “sk-xxxxx” } [/code]
项目扩展
这个基础架构可以很容易地扩展到更多场景:
- 添加新的数据源:只需要写一个新的 extractor 类,其他环节不用改
- 实时流处理:将 Airbyte 的 CDC 模式 + Kafka + Flink 引入,实现实时 ETL
- 数据血缘追踪:集成 OpenLineage 或 dbt docs,自动生成端到端的数据血缘图
- 大模型辅助:利用 LLM 自动生成 dbt 模型和测试代码
小结
恭喜你完成了这个完整的 ETL 实战项目!我们从零搭建了一个电商数据管道,覆盖了多源数据抽取、Python 清洗转换、数据仓库加载、dbt 建模转换、Airflow 自动调度以及数据质量监控。这个架构虽然不是最复杂的,但它代表了现代数据工程团队构建 ETL 管道的标准模式。掌握这个模式后,你可以把它应用到任何业务场景中——无论是电商、金融、游戏还是 IoT。
Summary: 电商 ETL 实战,完整覆盖抽取、转换、加载与 Airflow 编排。