大数据:Kappa架构
1. Kappa架构的诞生背景:为什么需要它?
要理解Kappa架构,必须先提它的“前辈”——Lambda架构。
-
Lambda架构:为了解决大规模数据的实时处理和批处理问题,Lambda架构采用了“双路并行”的策略。它包含三条层:
- 批处理层:处理所有历史数据,提供准确、全面的视图(如使用Hadoop、Spark)。高延迟,高准确性。
- 速度层:处理实时的新增数据,弥补批处理层的高延迟,提供近实时的视图(如使用Storm, Flink)。低延迟,近似结果。
- 服务层:合并批处理层和速度层的结果,提供给前端查询。
-
Lambda架构的痛点:
- 系统复杂:需要开发和维护两套独立的代码逻辑(批处理和流处理),它们可能产生不一致的结果。
- 运维成本高:需要管理两套不同的技术栈。
- 数据口径统一困难:确保批处理和流处理业务逻辑完全一致是一个巨大的挑战。
正是在这种背景下,Jay Kreps(LinkedIn的前工程师,Apache Kafka的联合创建者)在2014年提出了Kappa架构,其核心目标是:用一套流处理系统解决所有问题,简化架构。
2. Kappa架构的核心思想
Kappa架构的核心思想可以用一句话概括:
“所有数据都是流,批处理只是流的一个特例。”
它认为,没有必要维护两套独立的处理系统。相反,我们可以将历史数据重新作为流来播放,并通过单一的流处理引擎来处理实时数据和“重放”的历史数据。
3. Kappa架构的组成
一个典型的Kappa架构主要由以下两个核心组件构成:
-
消息队列/日志系统
- 角色:这是架构的核心与基石。它必须是一个支持高吞吐、可持久化、并能保留所有历史数据的分布式日志系统。
- 关键特性:数据重放能力。这是Kappa架构能够取代批处理的基础。你可以指定一个起始偏移量,重新消费历史数据。
- 典型技术:Apache Kafka 是最佳选择,因为它天然就是为了持久化日志而设计的。其他如Apache Pulsar也是可行的选项。
-
流处理引擎
- 角色:这是架构的大脑。它消费消息队列中的数据,执行计算逻辑(如聚合、连接、过滤等),并将结果输出到下游系统。
- 关键特性:需要支持有状态计算 和恰好一次的语义,以确保计算结果的准确性。
- 典型技术:Apache Flink(以其强大的状态管理和Exactly-Once语义而闻名)、Apache Spark Streaming(微批处理模型)、ksqlDB(用于Kafka的流式SQL引擎)等。
-
服务层/输出目标
- 角色:存储流处理引擎计算出的结果,并提供查询服务。
- 典型技术:可以是数据库(如Cassandra、HBase)、键值存储(如Redis)、数据仓库(如ClickHouse),或者直接生成一个新的Kafka Topic供其他服务消费。
4. Kappa架构的工作原理
让我们通过一个具体的例子来理解Kappa架构是如何工作的。
场景:一个电商平台需要实时计算每个商品的点击量。
首次部署与实时处理
- 数据采集:用户的每一次点击行为都被实时发送到Kafka的一个Topic中,例如
user_clicks。 - 流处理作业:一个Flink作业持续消费
user_clicksTopic中的数据。 - 实时计算:Flink作业维护一个“状态”,记录着每个商品的当前点击量。每来一条新的点击记录,就对相应商品的计数状态加一。
- 输出结果:Flink作业将每个商品的最新点击量实时写入一个数据库(如Redis)或另一个Kafka Topic(如
product_click_count)中,供前端应用查询。
当业务逻辑变更时(这是Kappa架构的精髓)
假设我们最初的计算逻辑有误,或者需要增加新的维度(例如,要区分来自不同地区的点击量)。
在Lambda架构中,你需要分别修改批处理层和速度层的代码,并保证两者逻辑一致,非常繁琐。
而在Kappa架构中,步骤如下:
- 停止旧作业: gracefully停止当前运行的Flink旧作业。
- 数据重放:
- 创建一个新的Flink作业,其中包含了新的、正确的业务逻辑。
- 让这个新作业从Kafka的
user_clicksTopic的起始偏移量(或者一个足够早的偏移量)开始消费。这就是所谓的“数据重放”。
- 处理与覆盖:
- 新作业会像处理实时数据一样,快速地将所有历史数据重新计算一遍。
- 新作业将计算结果输出到一个新的、临时的输出表或Topic中。
- 切换:当新作业追上实时数据后,将查询服务从旧的输出切换到新的输出上。
- 清理:停止并删除旧作业,删除旧的输出。
通过这种方式,我们仅用一套流处理逻辑,就完成了对历史数据的重新计算和业务逻辑的更新。
5. Kappa架构的优缺点
优点
- 架构简单:只有一套数据处理流水线,开发和维护成本显著降低。
- 逻辑一致:实时处理和历史数据处理使用同一份代码,从根本上避免了Lambda架构中数据口径不一致的问题。
- 低延迟:天然地为实时数据处理而设计,延迟极低。
- 可复现:由于所有原始数据都保存在日志中,任何计算都可以被精确地复现。
缺点
- 对消息队列要求极高:消息队列(如Kafka)需要具备巨大的存储容量来长期保存所有原始数据,成本较高。
- 数据重放可能很耗时:如果需要重放非常大量的历史数据,这个过程可能会很慢,期间需要双倍的计算资源。
- 流处理的局限性:对于复杂的、需要全量数据扫描的即席查询(Ad-hoc Query),或者迭代计算,纯流处理引擎可能不如批处理引擎高效。
- 运维复杂性转移:虽然系统组件变少了,但对核心组件(Kafka, Flink)的深度运维能力要求更高。
6. Kappa架构 vs. Lambda架构
| 特性 | Lambda架构 | Kappa架构 |
|---|---|---|
| 数据处理方式 | 批处理 + 流处理 两套路径 | 统一的流处理路径 |
| 系统复杂度 | 高(两套系统,代码逻辑可能不同) | 低(一套系统,一套代码) |
| 运维成本 | 高 | 相对较低 |
| 数据一致性 | 可能不一致(需合并层) | 强一致性 |
| 实时性 | 实时部分低延迟,全量结果高延迟 | 全链路低延迟 |
| 历史数据处理 | 批处理层天然支持 | 通过数据重放支持,可能较慢 |
| 技术栈 | Hadoop/Spark + Storm/Flink | Kafka + Flink/Spark Streaming |
7. 总结与适用场景
Kappa架构并非要完全取代Lambda架构,而是提供了一个更优雅、更简洁的替代方案,尤其适用于以下场景:
- 以实时数据为主的业务:如实时监控、实时推荐、实时风控、实时仪表盘等。
- 业务逻辑变更相对频繁,需要经常对历史数据进行重新计算。
- 团队技术栈相对统一,希望降低系统复杂性和运维成本。
- 数据源本身就是流式的,且对历史数据的回溯需求在可控范围内。
随着流处理技术(特别是Apache Flink)的日益成熟和“流批一体”概念的普及,Kappa架构的理念正在成为大数据架构演进的重要方向。许多新兴的数据系统都在朝着用一套API处理所有数据的目标努力,这正是Kappa架构思想的延续。
更多推荐
所有评论(0)