聊到“批处理”,不少程序员脑子里浮现的还是Windows时代那个画面:清理文件占用时写一个bat脚本,反复试探删除被进程锁住的文件,直到成功为止。大数据工程师嘴里的批处理是另一回事,但内核有点像——定时触发、攒一批数据、算完固定结果、然后释放资源走人。这套玩法曾经极其稳,但这两年,“转流处理”的呼声越来越大,而真正帮我把架构卡住位置的方案,是Kappa架构。
这篇文章从一个实际项目的复盘视角来写。我大概花了一年时间,把手上一套离线数仓Hive批处理链路,逐步迁移到基于Kafka加Flink的流处理管道,最后收敛成Kappa架构。如果你正犹豫要不要从批处理转流处理,或者已经定了要转但不知道第一步怎么迈,这篇复盘应该能让你少走弯路。我会按真实顺序讲:先理解批处理为什么会撑不住,再看Kappa凭什么能接盘,然后是按决策、实施、踩坑的完整过程。
1. 先从老底子说起:批处理为什么撑不住了
1.1 批处理的第一性原理:定时、全量、退出
批处理的运行模型极其简单,就三个词:定时、全量、退出。调度系统到点拉起任务,任务读取一个时间段内攒下的全部数据,算完之后把结果写到目标表,进程退出,等下一次调度。Hive离线报表、Spark批作业、DataX同步任务,本质都是这个模型。
这个模型稳了十多年,核心原因是它好审计、好重跑。数据不对,改完代码重跑一次今天的分区,结果就修正了。任务失败,调度系统自动把依赖下游的任务也停掉,不会出现“半个数仓更新”的脏状态。加上凌晨跑批的时间窗口是固定的,资源和运维都可以提前规划,团队按部就班,不容易出大事故。
但“固定时间窗口”恰恰也是批处理的天花板。它像食堂,每天早中晚固定三个饭点。平时你饿了就只能等,一到饭点所有人挤在一起取餐。只要规模可控,这套机制是高效的;可一旦有人喊出“我中午十一点半就要吃上饭”,食堂模式就露馅了。数据世界里的“随时要结果”,就是流处理出现的最直接动因。
顺带说一句,数据领域的批处理和日常技术圈理解的“批处理”往往是两个东西。日常聊到“解除文件占用批处理”,说的是Windows下一个脚本反复重试释放被占用的文件;大数据圈里的批处理更像一条夜间流水线,占着调度窗口跑完全量然后退出。两者都不灵活,但前者是“操作系统的临时约定”,后者是“数据处理的基本架构”。
1.2 撑不住批处理的三类瞬间
第一类瞬间是实时指标大屏。业务方上午十点开早会,要看截至当前时刻的GMV、订单量、活跃用户数。批处理再快,最快也是昨天凌晨算好的T-1数据。你可以硬折中,把跑批时间提前到每半小时一次,但半小时一次的“微批”本质还是把延迟卡在窗口边界上,数据分布的峰值一来,哪一批都慢。
第二类瞬间是在线风控和实时推荐。支付环节的欺诈检测,要求单笔交易在几百毫秒内完成风险评分;推荐系统要做用户实时行为特征,越新越好。这类场景物理上就没法用批处理实现,因为它模型本身就是“事件到了立刻判定”,不是“攒满了统一结算”。
第三类瞬间是凌晨跑批失败。这个做离线数仓的人几乎都经历过:凌晨五点某个上游表没有按时产出,下游任务全部卡死,重跑一次要六小时,整个白天的数据报表都是空白。批处理把这叫“任务依赖”,本质上却是一次连锁性的定时炸弹,任何一个环节抖动,整条链路都要跟着延迟。
这三类瞬间不会同时出现在每家公司身上,但只要出现其中两个,架构转型的立项理由就成立了。
2. 流处理的本质,以及Kappa架构从哪来
2.1 流处理的三个关键词:事件、时间、状态
流处理跟批处理最根本的区别,在于它不等人。数据是一条一条以事件的形式到达的,处理程序长年累月运行在那里,每个事件到了直接计算,算完送到下游,然后继续等下一个。批处理理解的“数据”是文件里的存量记录,流处理理解的“数据”是正在发生的一连串行为。
想让这种模式不出错,必须理解三个概念:
第一个是事件。事件是流处理里的基本单位,一条支付记录、一次点击行为、一个传感器读数,都可以编码成一条事件。事件里有业务字段,有时间戳,这是后面所有窗口计算的依据。
第二个是时间。流处理里有两套时间——事件发生时间(Event Time)和事件被处理的时间(Processing Time)。绝大多数业务场景应该基于事件时间计算,但网络延迟、数据乱序会带来“迟到数据”,所以引入了水位线(Watermark),把它理解成“迟到数据的最晚容忍时间线”即可。水位线设置太短,大量迟到数据被丢掉;设置太长,结果出得慢。
第三个是状态。流处理运算不是无状态的一锤子买卖。窗口聚合要记住这一窗口所有事件的总和,去重要记住已经见过的ID,这些记忆就是状态。状态要保存在引擎内部,还要定期快照,否则程序一重启,记忆全丢,计算结果就变成一笔糊涂账。
流处理这三件事,每一件都比批处理更复杂,但换来的好处是“数据一到就能出结果”,延迟从小时级压到秒级甚至毫秒级。这也是Kappa架构能成立的前提——只有流引擎能把重放、状态恢复、一致性做到位,架构才有底气把所有计算都放到流管道里。
2.2 Lambda架构的痛,是Kappa架构出现的直接原因
在Kappa之前,业内想实时化的第一选择是Lambda架构。Lambda的思路很直白:用批处理链路保证最终正确,用流处理链路保证实时展示,两条链路各自算,最后在服务层合并。听起来完美,落地之后全是眼泪。
两套代码要写两次,同一个业务指标,批处理一套Hive SQL,流处理一套Flink SQL。改个口径要动两个地方,上线要发两次,出问题要排查两条链路。更磨人的是,两条链路算出来的数字经常对不上,一边可能用了“用户下单时间”分组,另一边用了“支付完成时间”分组,最后业务拿着两个数字来质问:到底哪个是真的?
Lambda架构没有错,它解决的是那个年代流处理引擎不够可靠的现实问题。但维护成本的双倍付出、结果对账的持续消耗,让很多团队最终走向了Kappa的核心理念:只保留一条流处理链路,批处理能做到的修正、回溯、重算,全部用流处理的重放能力来完成。这是架构设计的减法,也是Kappa最吸引人的地方。
我把三种架构放在一起比较过,这样更直观:
| 维度 | 传统批处理 | Lambda架构 | Kappa架构 |
|---|---|---|---|
| 延迟 | 小时级 | 流链路秒级,批链路小时级 | 秒级 |
| 代码数量 | 一套 | 两套,各自维护 | 一套 |
| 历史重算方式 | 重跑批任务 | 重跑批任务 | 从日志重放事件 |
| 数据一致性保证 | 容易 | 难,双链路经常对不上 | 相对容易 |
| 运维复杂度 | 低 | 高 | 中等 |
2.3 “日志即真相”:Kafka在Kappa里的角色
Kappa架构有一个精神内核,叫“日志即真相”。这句话最早是Jay Kreps在提出Kappa概念时说的,理解它基本就理解了Kappa。
所有业务数据,不管来自数据库Binlog、前端埋点还是后端日志,都先进入一个统一的、不可变的消息日志——通常就是Kafka。这个日志里的每一条消息,只代表“事实曾经发生过”,不代表“该发给谁”。下游是谁都不知道,消息只是在日志里静静地躺着。
流处理作业的角色也从此发生转变,它不再是一个“任务”,而是这个日志的读者。作业从某个offset开始往后读,边读边算,把结果写到服务层。这里最关键的一点是,日志是有序、可重放的。你想从三天前重新读一遍,把消费组的offset重置到三天前即可;你想修正一个计算逻辑,改完代码把日志重放一遍,全量结果就重建出来了。
这就是Kappa能把批处理那条链路砍掉的底气——批处理的价值不在于批处理本身,而在于“随时可以从头重新计算”的能力。Kappa把这个能力从“重跑任务”变成了“重放日志”,实现的成本和灵活度都大幅优化。Kafka在这里不再只是消息队列,它被真正当成存储层来用。
3. 动手之前,先想清楚这三件事
3.1 延迟目标决定架构形态
做架构转型最忌讳一上来就问“用Flink还是Spark”,正确的第一个问题是:你到底需要多少秒的延迟。
如果业务只要分钟级结果,比如每五分钟更新一次报表,那完全不用全量流化。把现有批处理任务改成十五分钟调度一次,配合Kafka攒一波数据,基本也够用,稳定性还高。如果业务要求十秒以内的准实时,才开始考虑流处理管道。而真正需要毫秒级响应的,只有在线风控、实时推荐这类场景,它们本来就不该走批处理。
我见过不少团队把实时大屏的需求当成全量流化的理由,大屏看起来炫,实际上数据链路慢个三十秒业务方根本感知不到。为了一个三十秒延迟的展示需求,去重构一套数据架构,投入产出比非常不划算。
实际操作中,可以按延迟要求把现有任务分三档:第一档是必须秒级响应的,比如风控特征、在线推荐特征,优先流化;第二档是分钟级可以接受的,比如报表、大屏,先看能不能用调度频率解决问题;第三档是离线分析、月度结算、财务对账,老老实实留在批处理里。只流化真正需要实时的部分,比一口吃掉整个数仓靠谱得多。
3.2 可重放能力:Kappa架构的命门
Kappa的所有优势都建立在一个前提上:你能随时把历史事件重新读一遍。所以动手之前,必须想清楚Kafka里的数据能留多久、能不能重放、重放的成本可不可控。
Kafka默认的日志保留时间通常只有七天,七天之前的消息会被清理掉。如果历史消息都没了,重放根本无从谈起。Kappa落地时要主动调整保留策略,比如把核心事件Topic的保留时间调到三十天或九十天,成本主要落在Broker的磁盘上。另一个手段是启用Kafka的Log Compaction,只保留每个主键的最新一条记录,适合保存状态快照类数据,能大幅减少重放时的数据量。
还要估算重建视图的耗时。假设你的核心事件每秒十万条,三十天总量约两千六百亿条,重放一遍如果只能跑到每秒五十万条,那就要好几个小时。这个时间业务能不能等,必须在设计期就知道。一旦重放耗时按天计算,Kappa的“修正能力”就变成了摆设。
这里有个经验可以分享:不要只依赖Kafka原始日志存储。定期把流作业自身的关键状态做一次快照存到HDFS,重放时从快照开始算,而不是从最早的offset开始。这样既保住可重放能力,又把修正耗时压缩到分钟级。
3.3 引擎选型:为什么落地我选了Flink
Kappa架构的流处理引擎,业内基本是Flink和Spark Structured Streaming二选一。我的结论是:对要做实时数仓和状态化计算的场景,Flink更适合。
| 对比项 | Flink | Spark Structured Streaming |
|---|---|---|
| 处理模型 | 原生事件流,逐条处理 | 微批为主,连续处理模式可选 |
| 延迟级别 | 毫秒级 | 秒级 |
| 状态管理 | 原生支持,状态后端成熟 | 依赖外部队列或存储 |
| 精确一次语义 | 内置checkpoint,成熟 | 支持但依赖配置 |
| SQL支持程度 | Flink SQL,窗口语义丰富 | Spark SQL,生态兼容性好 |
Flink的优势在于它从底层就是为“无穷事件流”设计的,窗口、水位线、状态、检查点这些流处理概念是原生的。Spark Structured Streaming把流拆成一段一段微批,处理模型跟批处理更接近,团队上手快,但窗口边界、乱序处理的灵活度差一些,延迟也做不到毫秒级。
选型不能只看技术,也要看团队底子。如果团队全是Spark背景、现有平台也是Spark生态,硬换Flink的学习成本非常高。但如果你要做的是长期、复杂的实时数仓,我建议直接上Flink,它把状态一致性这件事做得最正经。状态后端、增量快照、Exactly-once语义,这些能力在批转流过程中会反复用到,Flink对开发者的支持会舒服得多。
4. 从批处理到流处理:一次门店GMV指标的完整改造
4.1 盘点现有批任务:分清哪些值得流化,哪些继续保留
正式动手前,我先用一周时间盘点了线上所有批处理任务,把一百多个作业按来源、去向、延迟要求列了一张表。这不只是为迁移做准备,更重要的是逼自己回答一个问题:哪些任务是必须实时,哪些只是“觉得应该实时”。
最终分类结果是:真正必须流化的只有实时大屏、风控特征、在线推荐特征三类,加起来不到十个任务。剩下的离线报表、对账、月度统计全部保留在批处理链路。这个动作极其重要,它阻止了一场“把所有批处理都改成流处理”的灾难。
流处理不是银弹,它是为高时效场景准备的。批处理的吞吐优势和简单重跑能力,在离线分析里无可替代。Kappa架构并不排斥批处理,真正合理的落地形态是:核心业务事件全部进入Kafka,实时链路用流处理算,离线链路从同一份Kafka数据用批处理算,两条链路共享同一份数据源,而不是各自维护一套数据输入。
4.2 事件入湖:先把所有关键数据接进Kafka
盘点完之后,第一件事不是写Flink SQL,而是把所有关键数据源接入Kafka。这个动作叫事件入湖,它是Kappa架构的地基。
以门店GMV指标为例。原来是MySQL里的订单表,每天凌晨用Sqoop同步到Hive分区表。迁移后,订单表通过Debezium或Canal订阅Binlog,变更事件实时写入Kafka的dwd_order_payTopic。这样一个订单从产生到进入Kafka,延迟压到了百毫秒级。
建Topic时要注意分区数和副本数的设置。分区数直接决定并行处理的上限,我按下游Flink并行度来配,分区数是并行度的三倍,保证处理压力可以横向扩展。副本数设置三副本,配合Kafka自身的ISR机制防止Broker宕机丢数据。核心Topic的保留时间设为三十天,开启Log Compaction,为后续重放留出空间。
这个阶段最容易被低估的工作是Schema管理。消息格式统一用JSON还是Avro,字段变更怎么兼容,消费方如何感知,都要提前定好规范。我们后续深受其益——Flink SQL里建表时直接复用Schema Registry里的结构,字段对齐问题少了一大半。
4.3 Hive SQL改造成Flink SQL的实战代码
数据进Kafka后,就进入最核心的改造环节。从批处理代码改成流处理代码,我以最常写的报表指标——门店每日GMV为例演示。
原来Hive离线任务长这样,每天按日期分区聚合:
INSERT OVERWRITE TABLE dws_store_gmv_daily SELECT store_id, order_date, SUM(amount) AS gmv FROM dwd_order_pay WHERE dt = '${bizdate}' GROUP BY store_id, order_date;改造后的Flink SQL长这样,它是持续运行的,每五分钟自动滚动一个窗口聚合,实时写结果:
CREATE TABLE kafka_source ( store_id STRING, amount DECIMAL(10, 2), event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL '30' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'dwd_order_pay', 'properties.bootstrap.servers' = 'kafka-broker:9092', 'properties.group.id' = 'kappa_report_group', 'scan.startup.mode' = 'earliest-offset', 'format' = 'json' ); CREATE TABLE es_sink ( store_id STRING, window_end TIMESTAMP(3), gmv DECIMAL(20, 2), PRIMARY KEY (store_id, window_end) NOT ENFORCED ) WITH ( 'connector' = 'elasticsearch-7', 'hosts' = 'http://es-node:9200', 'index' = 'store_gmv_realtime' ); INSERT INTO es_sink SELECT store_id, TUMBLE_END(event_time, INTERVAL '5' MINUTE) AS window_end, SUM(amount) AS gmv FROM kafka_source GROUP BY store_id, TUMBLE(event_time, INTERVAL '5' MINUTE);这段SQL里有几个关键点值得展开。
水位线设的是事件时间减三十秒,意思是允许最多三十秒的乱序数据。窗口用的是滚动窗口,每五分钟关闭一个,然后输出该门店这个五分钟内的GMV。这里最容易踩的坑是分组时必须带窗口辅助函数TUMBLE_END(event_time, INTERVAL '5' MINUTE),不带它Flink SQL会拿整个流当无限大窗口做聚合,状态只涨不清,结果永远出不来。
还有一点,scan.startup.mode我设置的是earliest-offset,保证作业第一次启动时能从最早可消费位置开始,把历史数据也扫一遍。如果你只要实时增量,可以改成latest-offset。第一次全量消费通常是有意的,它能验证逻辑正确,也为后续数据订正打底。
4.4 一致性保障:checkpoint、状态后端与幂等写
流作业改造完,接下来最关键的是配置一致性参数。批处理天然具备“要么成功重跑,要么失败停靠”的特性,流处理没有重跑机会,它必须靠运行时的状态快照机制来保证“挂了不丢”。
取舍的核心是Flink的Checkpoint。我用的配置如下:
execution.checkpointing.interval: 60s execution.checkpointing.mode: EXACTLY_ONCE execution.checkpointing.min-pause: 30s state.backend: rocksdb state.checkpoints.dir: hdfs://namenode/flink-checkpointscheckpoint间隔设的是60秒,每六十秒做一次状态快照。第一次调整时很想把它改成十秒甚至一秒,后来发现这是极其典型的错误。快照频率太高,状态后端在持续做快照上传,计算线程反而被拖慢,任务吞吐下降严重。60秒是稳定性与恢复粒度的平衡点。
状态后端选择的是RocksDB而不是堆内存。堆内存状态恢复快,但一旦状态量超过JVM堆限制就会Full GC,整个流作业卡死。RocksDB把状态存在内存加磁盘,可以扛更大的状态量,代价是读写性能略慢,配合增量快照机制后,性能压力完全可以接受。
还有一个必须单独强调的坑:Flink的Exactly-once只保证引擎内部状态的一致性,并不保证下游写入不重复。ES写入用幂等ID覆盖、Hive结果用分区覆盖写,下游幂等才能真正做到端到端不重不丢。这两个一起配合,出错重放时才能保证业务侧看到的永远是干净数据。
4.5 新旧平行运行与平滑切换
流作业上线后,我没有立刻停掉离线批处理链路,而是让两条链路并行跑了两周。这相当于做了一场持续两周的对账。
对账的方式不算复杂。每日凌晨两点,跑一个批量脚本,把前一天的实时汇总结果和离线任务产出的结果做差值对比。差异在万分之五以内,说明链路基本稳住;差得大,就要逐小时排查是流处理少算了还是批处理口径不同。这里要说一句,最坑人的其实是“口径不一致”而不是“链路故障”。实时链路如果按支付时间聚合、离线链路按订单创建时间聚合,两边永远对不上,跟技术无关。
平行期稳定后,开始切下游消费者。把BI看板、数据API的读表源从离线表切到实时结果,线上大屏同时在流和批两个数据源各画一遍,观察半天,确认表现一致后再把离线任务停掉。
这步我强烈建议不要着急。切流最怕的是两边结果有细微差别而业务已经上线了。我有一次就是因为少了“退款单剔除”的口径,实时数据和离线数据对不上,最终只能回滚。宁可多跑一周对账,也别为了节省一天时间引入线上事故。
5. 上线三个月,我踩过的三个坑
5.1 回压风暴:一个下游变慢,整条链路跟着堵
上线第三周,某个核心作业的Kafka消费延迟突然从秒级涨到几百万条。打开Flink Web UI一看,Source算子是正常的,ES Sink算子背压指标显示100%,问题出在下游写入。
排查下来,并不是计算逻辑卡住了,而是ES集群里某个索引的主分片发生热点,写入吞吐骤降。Flink的背压机制一瞬间传导到整条链路,上游Kafka的Lag指数级增长,所有下游报表数据全部滞后。
回压的排查思路其实就一个:沿着背压指示从后往前找瓶颈算子。Flink UI里每个算子的BackPressured比例就是最直观的风向标。定位到瓶颈后再分析,是下游存储写不进去,还是某个算子计算量太大,或是不知不觉间一个超大Key把数据都集中到了同一个Subtask上。
那次解决的方案是给ES索引加了两倍分片,同时把Flink的批量写入参数调大,让数据攒一批再发,减少小请求对ES的冲击。效果立竿见影,Lag快速清零。后来我总结成一条经验:流处理链路的性能瓶颈八成在下游外部系统,而不是Flink本身,平台侧的任务要先把外部依赖的容量压测做足。
5.2 检查点失败:想恢复状态,没那么简单
月中的一个早晨,值班同事在群里扔了一张告警截图,某个作业的Checkpoint成功率掉到了40%。任务没有挂,但每次快照都超时失败,这比挂掉更难排查。
刚开始怀疑是状态太大。当时每个五分钟窗口的结果明细都保存在状态里,RocksDB本地没问题,但上传到HDFS时快照体积超过两个GB,检查点超时时间不够用。我开了增量快照,把全量上传改成只上传变化部分,快照体积从两个GB降到了两百MB,成功率立刻回满。
这里有一个适用于所有团队的排查顺序。先看超时时间是不是太短,再看快照体积是不是过大,然后看状态后端和HDFS之间的带宽。最大的可能是状态设计有问题,比如把没必要存的状态算子放太多。检查点不是越勤越好的,它应该负责“能恢复”,而不是每时每刻都在全量保存。
5.3 数据订正:Kappa如何“撤回”一天的错误结果
上线两个月后遇到一次真正的口径事故。业务方发现我们统计GMV时把退款单也算进去了,导致前两天的实时报表数据偏高。
如果还在批处理时代,解法很简单——改一行SQL,重跑那两天分区。但重跑分区意味着要等好几个小时,而实时报表的高峰期已经过了。Kappa时代的解法逻辑类似,但路径完全不同。
我先在作业里把“剔除退款单”的逻辑加上,然后新启了一个独立消费组,将该消费组的Offset重置到两天前,从那个点到当前时刻重新消费一遍所有订单事件,重算这两个天的每五分钟窗口。重放时,线上原有作业继续跑,不影响实时数据。重放完成后再切回一个消费组,利用ES结果表的幂等写入,把这两天错误的历史窗口结果一一覆盖掉。
这个能力是批处理给不了的——想修正哪一段数据,就从哪一段重放,修正期间线上服务照常运行。你把Kafka的Log Compaction和状态快照用起来,还能进一步缩小这段“需要重放的数据量”,让修正从小时级压缩到分钟级。这也是我一直坚持接入Kafka而不是用普通消息队列的原因:普通消息队列的使命是“成功投递”,Kafka在Kappa里承担的是“完整留存并允许重读”。
6. 这次转型教会我的几件事
技术上的经验前面已经写了很多,最后说几句更偏“个人体会”的东西。
批处理转流处理,难度最大的是思维模式的转换,不是代码。批处理工程师习惯“今天不对,明天重跑”,流处理没有明天,它永远在跑,每一次失败都必须立刻恢复并且不丢数据。这种思维转换需要团队一起完成,我花了至少三个月,才让团队每个人都能脱口而出“这个作业的离线重跑机制是什么,消费位点要盯哪几个指标”。
监控体系也要跟着换。批处理的监控核心是“作业有没有失败”,流处理的核心是“消费延迟在不在涨、检查点成功率够不够高、回压有没有出现”。调度报警那一套不能直接搬过来,反而要写很多自定义的可视化面板,把消费位点、背压指标、状态大小每天盯一遍。
还有一个现实建议。如果你的团队还在纠结要不要转型,最好的办法是找一个真正需要实时的场景打样,做成一条最小的流处理管道,把Kafka接入、Flink作业、结果对账整条链路跑通。打样阶段你会发现,架构上的问题远不如“口径一致性”和“可重放机制”这两个问题来得实际。把这两个问题解决好,Kappa基本就落地了一半。
如果只能给出一条建议,我会毫不犹豫说:先把数据接进Kafka,再考虑怎么算。数据不丢、可重放,后面一切转型都好说。至于架构是Kappa还是Lambda,反而只是第二步的选择题。