从 ETL 到架构思维

前面三篇文章我们分别深入了 ETL 的抽取、转换和加载三个环节。沿着这条路径,我们掌握了"怎么做"ETL。但当你从实现单个流程转向设计整个系统时,就需要切换到架构思维——考虑的不再是一条数据怎么跑,而是整个数据平台怎么搭。

数据架构的选择决定了系统的扩展性、维护成本和团队效率。选择一个不适合的架构,可能在数据量增长时遭遇性能瓶颈,或者在业务需求变化时难以调整。本文将从架构层面讨论 ETL 的各种模式,并与 ELT 进行深入对比。

传统 ETL 三層架构

传统 ETL 采用严格的三层架构,数据按照"源系统 → staging 区 → 数据仓库 → 数据集市"的顺序流动。

┌──────────┐    ┌──────────┐    ┌──────────┐    ┌──────────┐
│ 源系统    │ →  │ Staging  │ →  │ 数据仓库  │ →  │ 数据集市  │
│ (OLTP)   │    │ (原始层)  │    │ (明细层)  │    │ (汇总层)  │
└──────────┘    └──────────┘    └──────────┘    └──────────┘
                     ↓
                ┌──────────┐
                │ 转换引擎  │
                │ (ETL 层)  │
                └──────────┘

Staging 区

Staging 层是数据进入数据仓库的第一站。在此阶段,数据保持与源系统几乎一致的结构,不做任何转换。它的主要作用是隔离——避免直接从源系统拉数据时对业务产生影响。

-- Staging 表通常按源系统结构原样存储
CREATE TABLE stg_orders (
    order_id VARCHAR(50),
    order_data JSONB,     -- 保留原始 JSON,不做解析
    customer_id VARCHAR(50),
    order_amount VARCHAR(20),  -- 保留字符串,留待转换层处理
    order_date VARCHAR(20),
    source_system VARCHAR(50),
    ingestion_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);

数据仓库明细层

经过 ETL 转换引擎处理后,数据进入数据仓库的明细层。这里的数据已经过清洗和标准化,结构被设计为适合分析查询的模式(通常是星型模型或雪花模型)。

数据集市层

数据集市是面向特定业务主题的汇总数据。例如,销售团队使用销售数据集市,财务团队使用财务数据集市。每个数据集市的数据粒度、更新频率和访问权限都可以不同。

现代 ELT 架构

ELT(Extract-Load-Transform)是 ETL 的变体,核心变化是调换了转换和加载的顺序。数据先从源系统抽取出来,直接加载到目标平台,然后在目标平台内部进行转换。

┌──────────┐    ┌──────────┐    ┌──────────┐
│ 源系统    │ →  │ 目标平台  │ →  │ 转换处理  │
│ (OLTP)   │    │ 原始数据  │    │ (在目标)  │
└──────────┘    └──────────┘    └──────────┘

ELT 的典型流程

-- Step 1: Extract - 从源系统抽取数据(使用 dbt 或 Airbyte)
-- 将原始 CSV 导入到 raw schema
COPY raw.orders FROM '/data/landing/orders.csv' WITH CSV HEADER;

-- Step 2: Load - 数据已加载到目标平台,保持原始格式
-- 无需转换,直接入库
SELECT * FROM raw.orders LIMIT 10;

-- Step 3: Transform - 在目标平台内完成转换
-- 使用 SQL 进行数据清洗和建模
CREATE TABLE dw.clean_orders AS
SELECT
    order_id::INTEGER,
    customer_id::INTEGER,
    CASE
        WHEN amount ~ '^\d+\.?\d*$' THEN amount::DECIMAL(10,2)
        ELSE 0.00
    END AS amount,
    order_date::DATE
FROM raw.orders
WHERE order_id IS NOT NULL;

ELT 的适用场景

ELT 架构在以下场景中尤其有优势:

  • 云数据仓库:Snowflake、BigQuery、Redshift 的计算能力非常强,直接在目标内做转换效率极高
  • 数据量大:不需要额外的转换集群,利用数仓的弹性扩缩能力
  • 团队偏 SQL:如果团队擅长 SQL 而不是 Python/Java,ELT 的学习成本更低
  • 需要数据探索:原始数据保留在数仓中,分析师可以随时探索和发现新的数据用途

ETL vs ELT:全面对比

对比维度 ETL ELT
转换位置 专门的转换引擎(Python/Spark) 目标平台内(SQL/dbt)
数据流向 源 → 转换引擎 → 目标 源 → 目标 → 转换
存储要求 目标只需要存最终数据 目标需要存原始数据+转换后数据
转换能力 支持复杂编程逻辑 受限于 SQL 表达能力
初始加载速度 慢(转换后才入库) 快(原始数据直接入库)
灵活性 转换逻辑变更需重跑整个流程 原始数据保留,随时可重新转换
工具代表 Informatica, Spark, Pandas dbt, Snowflake, BigQuery
团队要求 需要数据工程能力 SQL 能力强即可
治理能力 集中控制数据质量 需要额外治理原始数据访问

如何选择 ETL 还是 ELT?

没有绝对的优劣,只有适合不适合。以下决策框架可以帮助你判断:

选 ETL 的情况:

  1. 目标系统的计算能力有限(如传统本地部署的数仓)
  2. 转换逻辑非常复杂,需要编程语言的支持(如 NLP 处理、图像识别)
  3. 数据质量要求极高,需要在入库前完成严格校验
  4. 源系统数据频繁变动,需要复杂的 CDC 和实时处理
  5. 团队有较强的数据工程背景

选 ELT 的情况:

  1. 目标平台是云数据仓库(Snowflake、BigQuery、Redshift)
  2. 数据量巨大,不想在转换集群上花额外成本
  3. 团队以数据分析师和 SQL 专家为主
  4. 业务流程还在快速迭代,原始数据需要保留以便随时重新建模
  5. 使用 dbt 等现代数据转换工具

Lambda 架构

Lambda 架构是一种同时处理批处理和流数据的混合架构。它将数据流分为两条路径:批处理路径负责全量数据的精确计算,流处理路径负责实时数据的快速处理。

                    ┌──────────────┐
                    │ 所有数据进入    │
                    └──────┬───────┘
                           │
               ┌───────────┴───────────┐
               │                       │
        ┌──────▼──────┐        ┌──────▼──────┐
        │ 批处理层     │        │ 流处理层     │
        │(全量计算)   │        │(实时计算)   │
        └──────┬──────┘        └──────┬──────┘
               │                       │
        ┌──────▼──────┐        ┌──────▼──────┐
        │ 批处理视图   │        │ 实时视图     │
        └──────┬──────┘        └──────┬──────┘
               │                       │
               └───────────┬───────────┘
                           │
                    ┌──────▼──────┐
                    │ 服务层       │
                    │(合并查询)   │
                    └─────────────┘

Lambda 架构的 ETL 实现

from datetime import datetime, timedelta
import pandas as pd
from pyspark.sql import SparkSession

class LambdaPipeline:
    """简化的 Lambda 架构 ETL 示例"""

    def __init__(self):
        self.spark = SparkSession.builder.appName("LambdaETL").getOrCreate()

    def batch_layer(self, date):
        """批处理层:全量历史数据计算"""
        # 读取全量历史数据
        historical = self.spark.read.parquet("/data/lake/sales/")
        # 执行复杂聚合
        batch_result = (
            historical
            .groupBy("product_id", "month")
            .agg({"amount": "sum", "quantity": "sum"})
            .withColumnRenamed("sum(amount)", "total_amount")
            .withColumnRenamed("sum(quantity)", "total_quantity")
        )
        batch_result.write.mode("overwrite").parquet("/data/batch_view/")
        print(f"批处理层更新完成,处理历史数据")

    def speed_layer(self, today=None):
        """流处理层:当天实时数据"""
        if today is None:
            today = datetime.now().strftime("%Y-%m-%d")

        # 模拟读取实时流数据
        stream_data = self.spark.read.parquet(f"/data/stream/{today}/")
        speed_result = (
            stream_data
            .groupBy("product_id")
            .agg({"amount": "sum", "quantity": "sum"})
        )
        speed_result.write.mode("overwrite").parquet("/data/speed_view/")
        print(f"流处理层更新完成,处理 {today} 的实时数据")

    def serving_layer(self, query_date):
        """服务层:合并批处理和实时结果"""
        batch = pd.read_parquet("/data/batch_view/")
        speed = pd.read_parquet("/data/speed_view/")

        # 合并结果(实时覆盖批处理中的同一天数据)
        merged = pd.concat([batch[~batch["date"].isin([query_date])], speed])
        return merged

Lambda 架构的优缺点

优点:同时保证数据的完整性和实时性;批处理可以修正流处理中的错误。

缺点:维护两套代码逻辑,复杂度加倍;批处理和流处理的结果可能不一致;开发和运维成本高。

Kappa 架构

Kappa 架构是对 Lambda 架构的简化。它认为可以通过单一的流处理管道同时满足实时和批处理的需求,不需要维护两套代码。

┌──────────┐     ┌──────────┐     ┌──────────┐
│ 数据流    │ →   │ 流处理    │ →   │ 结果存储  │
│ (Kafka)  │     │ (Flink)  │     │ (数据湖)  │
└──────────┘     └──────────┘     └──────────┘
                     │
                     ↓ 重放历史数据
              ┌──────────────┐
              │ 得到全量结果   │
              └──────────────┘

Kappa 架构的核心思路

Kappa 架构的关键洞察是:既然流处理框架(如 Flink、Kafka Streams)可以配置状态大小和回溯历史数据,那么只要数据源(如 Kafka)保留了完整的变更日志,就完全可以重放历史数据来重新计算结果,而不需要单独维护一个批处理管道。

# Kappa 架构的 ETL 思路:所有数据通过统一的流处理管道
def kappa_pipeline(env, checkpoint_path):
    """简化的 Kappa 架构 Flink 处理流程"""
    # 所有数据来自统一的数据流
    data_stream = env.from_source(
        KafkaSource.builder()
            .set_topics("all_events")
            .set_starting_offsets(OffsetsInitializer.committed_offsets(
                OffsetsInitializer.earliest()
            ))
            .build(),
        WatermarkStrategy.noWatermarks(),
        "kafka-source"
    )

    # 统一的处理逻辑,既处理实时数据也处理历史回放
    result = (
        data_stream
        .key_by(lambda event: event["product_id"])
        .window(TumblingProcessingTimeWindows.of(Time.hours(1)))
        .aggregate(SalesAggregator())
    )

    result.sink_to(ClickhouseSink())
    return env.execute("kappa-etl-pipeline")

Kappa 架构的适用场景

Kappa 架构最适合数据流特性统一、不需要区分实时和批处理的场景。比如日志分析、点击流分析、监控指标处理等。它最大的限制是数据源必须支持消息回溯(如 Kafka 的日志保留功能)。

云原生 ETL 模式

Serverless ETL

云厂商提供了完全托管的 Serverless ETL 服务,开发者只需要定义数据源、转换逻辑和目标,无需管理任何基础设施。

# AWS Glue(Serverless Spark ETL)
import sys
from awsglue.transforms import *
from awsglue.context import GlueContext
from pyspark.context import SparkContext

sc = SparkContext()
glueContext = GlueContext(sc)

# 从 S3 读取数据
orders = glueContext.create_dynamic_frame.from_options(
    connection_type="s3",
    connection_options={"paths": ["s3://data-lake/raw/orders/"]},
    format="parquet"
)

# 转换:过滤无效订单
valid_orders = orders.filter(
    lambda row: row["amount"] is not None and row["amount"] > 0
)

# 写入数据仓库
glueContext.write_dynamic_frame.from_options(
    frame=valid_orders,
    connection_type="redshift",
    connection_options={
        "url": "jdbc:redshift://dw-cluster:5439/dw",
        "dbtable": "public.clean_orders",
        "user": "etl_user",
        "password": "***"
    }
)
# Google Cloud Dataflow(Apache Beam)
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions

options = PipelineOptions(project="my-project", region="us-central1")

with beam.Pipeline(options=options) as p:
    (p
     | "ReadFromGCS" >> beam.io.ReadFromText(
         "gs://data-lake/orders/*.csv"
     )
     | "ParseCSV" >> beam.Map(lambda line: parse_csv(line))
     | "FilterInvalid" >> beam.Filter(lambda x: x.amount > 0)
     | "Transform" >> beam.Map(transform_order)
     | "WriteToBigQuery" >> beam.io.WriteToBigQuery(
         table="project:dw.clean_orders",
         write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND
     )
    )

Data Mesh 数据网格

数据网格是一种分布式数据架构理念,它将数据视为产品,由各个业务域团队自行负责数据的生产、治理和交付。

在 Data Mesh 架构下,ETL 不再是集中式的。每个业务域都有自己的 ETL 管道,负责将本域的数据加工成可复用的数据产品。这些数据产品通过标准接口(如共享的 Parquet 文件或 API)提供给其他域消费。

┌─────────────────────────────────────────────────┐
│                 数据平台层                        │
│   (基础设施:存储、网络、IAM)                    │
├──────────┬──────────┬──────────┬─────────────────┤
│ 销售域    │ 营销域   │ 供应链域  │ 财务域          │
│ ETL 管道  │ ETL 管道 │ ETL 管道 │ ETL 管道        │
│ 数据产品   │ 数据产品  │ 数据产品  │ 数据产品        │
└──────────┴──────────┴──────────┴─────────────────┘

架构选型决策框架

面对这么多架构选择,如何做出决定?以下是一个实用的决策流程:

数据量多大?
├─ < 100GB/天 → 传统 ETL 或简单 ELT
└─ > 100GB/天 → 考虑云原生 ELT 或 Lambda

实时性要求多高?
├─ T+1(每天一次)  → 批处理 ETL/ELT 完全够用
├─ 分钟级延迟        → 微批次处理(Spark Streaming)
└─ 秒级延迟          → Lambda 或 Kappa 架构

团队技术栈?
├─ SQL 为主  → ELT + dbt
├─ Python 为主 → ETL + Pandas/Spark
└─ 混合       → 批处理用 ETL,实时用流处理

云环境?
├─ 多云/本地   → 传统 ETL 工具(Informatica、Airflow)
├─ 单一云厂商  → 云原生服务(Glue、Dataflow、Airbyte)
└─ Snowflake/BigQuery → 强烈推荐 ELT 模式

小结

本文从架构层面审视了 ETL 的多种实现模式。我们对比了传统 ETL 三层架构和现代 ELT 架构的区别,讨论了 Lambda 和 Kappa 两种大数据架构范式,以及云原生环境下的 Serverless ETL 和数据网格等新兴理念。没有一种架构是银弹,选择什么模式取决于你的数据规模、实时需求、团队能力和云基础设施。

通过这五篇文章,我们从 ETL 的基本概念出发,深入了抽取、转换、加载三个核心环节的具体实现,最后提升到架构设计层面,建立了从入门到精通的完整知识体系。希望这套课程能帮助你构建扎实的 ETL 技能,从容应对实际工作中的数据集成挑战。

Summary: ETL 架构模式、ELT 对比、Lambda/Kappa 架构与选型决策。