大数据架构

深入理解大数据处理架构:Lambda 与 Kappa 架构对比、实时流处理、数据湖与湖仓一体、批流一体,掌握从数据采集到价值输出的完整数据工程体系。

大数据架构

随着数据规模爆发式增长,企业需要专门的大数据架构来处理海量数据的采集、存储、计算与分析。本文系统讲解 Lambda 与 Kappa 架构的演进、实时流处理、数据湖与湖仓一体、批流一体等核心概念与技术选型的完整方法论。


1. 大数据处理的核心挑战

挑战传统数据库大数据系统
数据量GB 级PB 级
数据速度批量写入实时流式写入
数据种类结构化结构化 + 半结构化 + 非结构化
数据价值已知问题回答未知问题探索
扩展方式Scale-UpScale-Out(分布式)

2. Lambda 架构

Lambda 架构由 Twitter 的 Nathan Marz 提出,同时维护批处理层和实时处理层,最终通过服务层合并结果。

2.1 三层结构

      ┌─────────────────────────────────────────────┐
      │                    数据源                     │
      │  (日志、数据库、消息、IoT、API)                │
      └──────────────┬──────────────┬───────────────┘
                     │              │
                     ↓              ↓
      ┌──────────────────┐  ┌──────────────────┐
      │    批处理层       │  │    实时处理层      │
      │   (Batch Layer)  │  │ (Speed Layer)    │
      │                   │  │                   │
      │  全量数据存储      │  │  增量流处理        │
      │  HDFS/S3          │  │  Kafka + Flink    │
      │                   │  │                   │
      │  MapReduce/Spark  │  │  窗口聚合/CEP     │
      │  离线计算          │  │  低延迟结果       │
      └─────────┬─────────┘  └─────────┬─────────┘
                │                      │
                ↓                      ↓
      ┌─────────────────────────────────────────┐
      │              服务层(Serving Layer)       │
      │                                         │
      │  ┌──────────────┐  ┌──────────────┐    │
      │  │ 批处理视图    │  │ 实时处理视图  │    │
      │  │ (精确结果)    │  │ (近似结果)    │    │
      │  └──────────────┘  └──────────────┘    │
      │           ↓ 合并 ↓                      │
      │      ┌──────────┐                      │
      │      │ 最终视图  │  ← 查询时合并批结果+实时增量 │
      │      └──────────┘                      │
      └─────────────────────────────────────────┘

2.2 各层技术选型

层级技术选择说明
批处理存储HDFS、S3、GCS不可变、 Append-Only
批处理计算Spark、Hive、Presto吞吐优先,小时/天级延迟
实时消息Kafka、Pulsar高吞吐消息队列
实时计算Flink、Spark Streaming、Storm毫秒~秒级延迟
服务层存储HBase、Cassandra、Druid、ClickHouse低延迟点查

2.3 Lambda 的问题

Lambda 的核心痛点:

1. 双代码路径
   批处理和实时处理需维护两套逻辑(同一份计算写两遍)
   业务逻辑变更 → 需同时修改两处 → 容易不一致

2. 系统复杂度高
   维护两套独立系统 = 双倍运维成本

3. 合并查询复杂
   查询时需合并批结果和实时增量
   实时层数据需有 TTL(过期清理)

3. Kappa 架构

Kappa 架构由 LinkedIn 的 Jay Kreps 提出,主张只保留实时处理层,用流处理统一批处理和实时计算。

3.1 核心思想

      ┌─────────────────────────────────────────────┐
      │                    数据源                     │
      └───────────────────┬───────────────────────────┘
                          │
                          ↓
      ┌─────────────────────────────────────────────┐
      │              消息队列(Kafka)                 │
      │                                             │
      │  ┌─────────┐  ┌─────────┐  ┌─────────┐     │
      │  │ 原始数据 │  │ 原始数据 │  │ 原始数据 │ ...  │
      │  │  Event  │  │  Event  │  │  Event  │     │
      │  └─────────┘  └─────────┘  └─────────┘     │
      │                                             │
      │  特性:不可变、顺序、可重放                     │
      └───────────────────┬───────────────────────────┘
                          │
           ┌──────────────┼──────────────┐
           │              │              │
           ↓              ↓              ↓
      ┌─────────┐  ┌─────────┐  ┌─────────┐
      │实时应用  │  │ 历史重算 │  │ 离线分析 │
      │(低延迟) │  │(修正Bug)│  │(全量报表)│
      └─────────┘  └─────────┘  └─────────┘
      
      统一使用流处理引擎(Flink/Spark Streaming)
      历史重算:从 Kafka 最早 offset 重新消费

3.2 重算机制

场景:发现实时处理逻辑有 Bug

Lambda:修改批处理代码 → 重新跑全量批作业
Kappa:
  1. 部署修正后的流处理 Job
  2. 新 Job 从 Kafka 最早的 offset 开始消费
  3. 新结果写入新的输出表
  4. 切换查询到新的输出表
  5. 旧 Job 停止

关键依赖:Kafka 需保留足够长的历史数据
  → 配合 S3/GCS 做冷存储(Tiered Storage)

3.3 Lambda vs Kappa

特性LambdaKappa
系统复杂度高(两套系统)低(一套系统)
运维成本
开发成本高(双代码路径)低(单代码路径)
结果精确性批处理精确 + 实时近似依赖流处理语义
历史重算批处理天然支持需消息队列长期保留数据
适用场景强一致性要求的离线报表以流为主的现代架构

4. 实时流处理

4.1 流处理核心概念

概念说明
Event Time事件实际发生的时间(数据携带的时间戳)
Processing Time数据被处理的时间(系统当前时间)
Ingestion Time数据进入流系统的时间
Watermark允许延迟到达的数据处理的进度标记
Window将无限流切分为有限块进行计算

4.2 窗口类型

Tumbling Window(滚动窗口):
  ┌────┐┌────┐┌────┐┌────┐
  │0-10││10-20││20-30││30-40│  不重叠,固定大小
  └────┘└────┘└────┘└────┘

Sliding Window(滑动窗口):
  ┌──────┐
   ┌──────┐
    ┌──────┐  窗口可重叠,slide < size
     ┌──────┐

Session Window(会话窗口):
  ┌──┐    ┌──────┐  ┌─┐
  └──┘    └──────┘  └─┘  由活动间隙触发,动态长度
     ↑gap↑        ↑gap↑

Global Window(全局窗口):
  ┌────────────────────────┐  整个流一个窗口,需 Trigger 触发计算
  └────────────────────────┘
// Flink DataStream API
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// 设置时间语义:Event Time
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);

DataStream<OrderEvent> orders = env
    .addSource(new KafkaConsumer<>("orders", new OrderDeserializationSchema()))
    .assignTimestampsAndWatermarks(
        WatermarkStrategy.<OrderEvent>forBoundedOutOfOrderness(Duration.ofSeconds(5))
            .withTimestampAssigner((event, timestamp) -> event.getOrderTime())
    );

// 按商品分组,统计每 10 秒成交额
orders
    .keyBy(OrderEvent::getProductId)
    .window(TumblingEventTimeWindows.of(Time.seconds(10)))
    .aggregate(new SumAggregateFunction())
    .addSink(new RedisSink<>(...));

env.execute("Real-time Order Analytics");

5. 数据湖与湖仓一体

5.1 数据湖(Data Lake)

原始格式存储海量异构数据的存储系统,支持结构化、半结构化、非结构化数据。

数据湖架构:

  ┌──────────────────────────────────────────┐
  │          数据采集层                       │
  │   Kafka / Flume / Logstash / Sqoop       │
  └───────────────────┬──────────────────────┘
                      ↓
  ┌──────────────────────────────────────────┐
  │          数据存储层(对象存储)              │
  │   ┌──────────┐ ┌──────────┐ ┌──────────┐│
  │   │ 原始区    │ │ 处理区    │ │ 服务区    ││
  │   │(Bronze)  │ │(Silver)  │ │(Gold)    ││
  │   │原始格式   │ │清洗转换   │ │业务视图   ││
  │   └──────────┘ └──────────┘ └──────────┘│
  │          S3 / OSS / GCS / HDFS           │
  └──────────────────────────────────────────┘
                      ↓
  ┌──────────────────────────────────────────┐
  │          计算引擎层                       │
  │   Spark / Flink / Presto / Trino         │
  └──────────────────────────────────────────┘

5.2 数据湖 vs 数据仓库

特性数据仓库(数仓)数据湖
数据类型结构化结构化 + 半结构化 + 非结构化
Schema写时定义(Schema-on-Write)读时定义(Schema-on-Read)
用户BI 分析师、业务人员数据科学家、工程师
用途报表、BI、已知问题机器学习、探索性分析
成本高(专有存储)低(对象存储)
性能优化查询快需额外优化

5.3 湖仓一体(Lakehouse)

结合数据湖的灵活性和数据仓库的性能与管理能力。

Lakehouse 关键特性:

1. 事务支持(ACID)
   Delta Lake / Apache Iceberg / Apache Hudi
   提供并发写、快照隔离、时间旅行

2. Schema 强制与演化
   可定义 Schema、自动演进、兼容旧数据

3. BI 性能
   物化视图、索引、缓存层

4. 开放格式
   Parquet(列式存储)+ 元数据层
   不绑定特定计算引擎

代表产品:
  - Databricks Delta Lake
  - Apache Iceberg(Netflix/Apple)
  - Apache Hudi(Uber)
  - Snowflake / BigQuery(外部表)

5.4 Delta Lake 示例

from delta import configure_spark_with_delta_pip
from pyspark.sql import SparkSession

builder = SparkSession.builder \
    .appName("DeltaLakeExample") \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")

spark = configure_spark_with_delta_pip(builder).getOrCreate()

# 写数据(ACID 事务)
df.write.format("delta").mode("overwrite").save("/data/orders")

# 读数据
spark.read.format("delta").load("/data/orders").show()

# 时间旅行:查询历史版本
spark.read.format("delta").option("versionAsOf", 0).load("/data/orders").show()

# 变更数据流(CDC)
spark.readStream.format("delta").load("/data/orders") \
    .writeStream.format("console").start()

6. 批流一体

同一套 API 和计算引擎同时处理批数据和流数据。

6.1 批流一体架构

            ┌─────────────┐
            │   数据源     │
            └──────┬──────┘
                   │
        ┌──────────┴──────────┐
        ↓                     ↓
  ┌─────────────┐      ┌─────────────┐
  │  有界数据集  │      │  无界数据流  │
  │  (Batch)    │      │  (Stream)   │
  └──────┬──────┘      └──────┬──────┘
         │                    │
         └──────────┬─────────┘
                    ↓
            ┌───────────────┐
            │   统一引擎      │
            │ Spark/Flink   │
            └───────┬───────┘
                    ↓
            ┌───────────────┐
            │   统一输出      │
            │  表/视图/API    │
            └───────────────┘
// 同一套代码,既可跑批也可跑流
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);

// 创建表(批流统一)
tableEnv.executeSql("""
    CREATE TABLE orders (
        order_id STRING,
        product_id STRING,
        amount DECIMAL(10,2),
        order_time TIMESTAMP(3),
        WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
    ) WITH (
        'connector' = 'kafka',
        'topic' = 'orders',
        'properties.bootstrap.servers' = 'kafka:9092',
        'format' = 'json'
    )
""");

// SQL 查询(批流统一语法)
tableEnv.executeSql("""
    SELECT 
        product_id,
        SUM(amount) as total_amount,
        TUMBLE_START(order_time, INTERVAL '1' HOURS) as window_start
    FROM orders
    GROUP BY 
        product_id,
        TUMBLE(order_time, INTERVAL '1' HOURS)
""").print();

// 批模式:读取有界数据(如文件),输出最终结果
// 流模式:读取 Kafka,持续输出更新的结果

7. 大数据技术选型指南

场景推荐方案说明
实时报表(秒级)Flink + ClickHouse/Druid低延迟聚合 + OLAP 查询
离线数仓Spark + Hive/Iceberg批处理 + 数据湖存储
日志分析ELK / ClickHouse全文检索 + 聚合分析
用户画像Flink + HBase/Redis实时标签更新 + 快速查询
推荐系统Spark ML + Flink 特征离线模型训练 + 实时特征
数据治理Apache Atlas + Great Expectations元数据管理 + 数据质量

8. 总结

大数据架构演进路线:

传统数仓(ETL + RDBMS)
  → Hadoop 生态(MapReduce + HDFS + Hive)
    → Lambda 架构(批处理 + 实时分离)
      → Kappa 架构(纯流处理统一)
        → 湖仓一体(Data Lakehouse)
          → 批流一体(统一引擎处理两种数据形态)

关键趋势:
1. 存储和计算分离(对象存储 + 弹性计算)
2. 实时化(从 T+1 到 T+0)
3. 开放格式(Parquet + Iceberg/Hudi/Delta)
4. 云原生(K8s + Serverless 大数据)
5. 数据网格(Data Mesh,去中心化数据管理)

选型原则:
- 没有最佳方案,只有最适合的方案
- 从 Lambda 起步,逐渐向 Kappa 或湖仓一体演进
- 开放的存储格式是避免厂商锁定的关键

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页