项目概述

经过前面 14 篇文章的学习,你已经掌握了 ETL 的各个核心环节:数据抽取、转换、加载、SQL 技巧、Python 开发、现代工具链和工作流编排。现在是时候把它们全部串联起来了。

本章是一个端到端的实战项目。我们将为一个电商平台构建完整的数据管道,从原始数据采集到业务指标计算,全部自动化运行。

场景设定

假设你在一家中型电商公司工作。公司有以下几个数据源:

  1. MySQL 业务数据库:存储订单、用户、商品信息
  2. CSV 文件:存储在 S3 上的每日营销费用数据
  3. 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: fct_daily_sales description: “每日销售事实表” columns:

    • name: order_date tests:
      • not_null
    • name: revenue tests:
      • not_null
      • dbt_utils.accepted_range: min_value: 0 [/code]

第七步: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]

项目扩展

这个基础架构可以很容易地扩展到更多场景:

  1. 添加新的数据源:只需要写一个新的 extractor 类,其他环节不用改
  2. 实时流处理:将 Airbyte 的 CDC 模式 + Kafka + Flink 引入,实现实时 ETL
  3. 数据血缘追踪:集成 OpenLineage 或 dbt docs,自动生成端到端的数据血缘图
  4. 大模型辅助:利用 LLM 自动生成 dbt 模型和测试代码

小结

恭喜你完成了这个完整的 ETL 实战项目!我们从零搭建了一个电商数据管道,覆盖了多源数据抽取、Python 清洗转换、数据仓库加载、dbt 建模转换、Airflow 自动调度以及数据质量监控。这个架构虽然不是最复杂的,但它代表了现代数据工程团队构建 ETL 管道的标准模式。掌握这个模式后,你可以把它应用到任何业务场景中——无论是电商、金融、游戏还是 IoT。

Summary: 电商 ETL 实战,完整覆盖抽取、转换、加载与 Airflow 编排。