大数据架构
随着数据规模爆发式增长,企业需要专门的大数据架构来处理海量数据的采集、存储、计算与分析。本文系统讲解 Lambda 与 Kappa 架构的演进、实时流处理、数据湖与湖仓一体、批流一体等核心概念与技术选型的完整方法论。
1. 大数据处理的核心挑战
| 挑战 | 传统数据库 | 大数据系统 |
|---|---|---|
| 数据量 | GB 级 | PB 级 |
| 数据速度 | 批量写入 | 实时流式写入 |
| 数据种类 | 结构化 | 结构化 + 半结构化 + 非结构化 |
| 数据价值 | 已知问题回答 | 未知问题探索 |
| 扩展方式 | Scale-Up | Scale-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
| 特性 | Lambda | Kappa |
|---|---|---|
| 系统复杂度 | 高(两套系统) | 低(一套系统) |
| 运维成本 | 高 | 低 |
| 开发成本 | 高(双代码路径) | 低(单代码路径) |
| 结果精确性 | 批处理精确 + 实时近似 | 依赖流处理语义 |
| 历史重算 | 批处理天然支持 | 需消息队列长期保留数据 |
| 适用场景 | 强一致性要求的离线报表 | 以流为主的现代架构 |
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 触发计算
└────────────────────────┘
4.3 Flink 流处理示例
// 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 │
└───────────────┘
6.2 Flink Table 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 或湖仓一体演进
- 开放的存储格式是避免厂商锁定的关键
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。