ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Kafka+Flink实时流处理实战:从批处理思维到事件驱动架构

Kafka+Flink实时流处理实战:从批处理思维到事件驱动架构 1. 为什么批处理思路在实时场景下会翻车先讲个真实经历。之前我负责一个电商大促期间的实时监控项目业务方提了个需求每分钟统计一次全站订单金额、退款金额和支付成功率延迟不能超过10秒。我第一版方案很天真——沿用批处理的思路每隔10秒用SQL去数据库里扫一遍订单表然后聚合出结果。上线第一天就出事了。大促流量一起来订单表瞬间膨胀到几千万行那个查询直接跑了1分多钟数据库连接池被打满连带正常的交易查询都受了影响。后来仔细定位发现问题不是出在数据库性能上而是出在思维方式上批处理假设数据是静止的、可全量扫描的而实时数据流处理面对的是无穷无尽、一直在动的数据。这两者从根上就不是一回事。1.1 从定时算一遍到每条数据都算一遍的思维转变批处理的典型做法是T1离线计算今天凌晨跑昨天一整天的数据。它关心的是全量数据算完没有对延迟没有硬性要求。而实时流处理关心的是每一条数据到达之后多久能被消费、被聚合、被产出结果。延迟指标通常以秒甚至毫秒计。这个转变带来三个直接后果第一数据不再是可以反复读取的。离线任务可以今天跑失败明天重跑一遍数据源还在那儿。但实时数据流是连续不断的错过的数据可能就再也找不回来了。所以流处理框架必须有机制处理数据到达乱序数据重复到达数据丢失这些情况。第二计算模型从批量触发变成事件驱动。批处理是调度系统到点了开始跑一个任务流处理是来一条数据处理一条或者数据攒够一个窗口就触发一次计算。这就要求计算引擎本身是常驻服务时刻在监听数据源而不是靠外部调度器去踢一脚。第三全量这个概念消失了。你不可能在内存里存下无穷无尽的数据所以窗口Window就成了流处理最重要的基本单元——把无限的数据流切分成有限的一段一段每一段里做聚合。窗口怎么切、切多大、数据晚到了怎么处理这些在批处理里根本不用考虑在流处理里全是核心问题。打个比方批处理像会计月底做账把整月凭证搬过来慢慢算流处理像超市收银员顾客结完账立刻就得算清楚这一单收了多少钱。收银员的算法必须快、必须准还要能应对顾客突然改主意换商品这种乱序情况——对应到流处理里就是事件的乱序到达。1.2 实时流处理绕不开的三个时间问题我开始做实时项目时第一个被绕晕的概念就是时间。流处理里有三个时间很多人入门时都栽在这上面事件时间Event Time业务实际发生的时间比如用户在12点01分下单事件时间就是12:01。摄入时间Ingestion Time数据进入消息队列或计算引擎的时间。处理时间Processing Time计算引擎真正处理这条数据的时间。为什么这三个时间会不一致因为数据在网络上传输需要时间在消息队列里排队需要时间在计算引擎里等待调度也需要时间。高峰期的时候用户12点01分下的单可能12点05分才被处理引擎拿到。如果按处理时间来聚合这条订单就会被算进12:05那个窗口里而实际上是属于12:01的统计就错了。我刚上手时图省事直接用了处理时间结果大促当晚的数据统计和业务方手里的后台订单明细完全对不上人家后台是按事件时间看的。那一次我彻底明白了做实时统计只要涉及业务指标必须用事件时间不能偷懒。但用事件时间又带来新问题你怎么知道事件时间已经是全部到齐了数据有没有可能还在路上这就引出了水位线Watermark这个机制后面我会专门讲。2. 技术选型为什么我最终选了Kafka Flink这套组合第一版方案翻车之后我重新做了技术选型。当时市面上主流能打的选择是消息队列用Kafka计算引擎在Flink和Spark Streaming之间二选一。也有团队用Pulsar Flink或者干脆用Kafka Streams但综合社区成熟度、招聘难易度和踩坑成本的考虑我最终选了Kafka Flink这套组合。2.1 消息队列选型Kafka凭什么成了流处理的事实标准消息队列不是只有Kafka一家RabbitMQ、RocketMQ、Pulsar都能做。但实时数据流处理这个场景下Kafka有几个非常对口的特性高吞吐 低延迟Kafka的设计目标是顺序读写磁盘利用操作系统的Page Cache吞吐量能达到每秒几十万甚至上百万条消息同时延迟还能压在几十毫秒以内。这正好匹配实时数据流的大流量场景。分区模型天然适合并行消费Kafka的一个Topic可以拆成多个分区每个分区内的消息是有序的不同分区之间可以并行消费。Flink的并行度配置可以直接对齐Kafka的分区数做到一个分区对应一个并行子任务数据不打架。消息可以重放这是Kafka相比RabbitMQ最狠的一点。Consumer可以指定offset从任意位置重新消费Flink的检查点Checkpoint机制就依赖这个特性——任务故障重启后从最近的检查点恢复offset重新消费没处理完的数据做到不丢不重。我遇到过不少团队用RabbitMQ做实时流后面都栽了。RabbitMQ的模型是消息被消费了就删除而流处理需要的是数据流可以被多次消费、回溯、重放这两个场景根本不是一个路数。所以如果项目是正经做实时数仓或实时计算消息队列直接上Kafka别纠结。2.2 计算引擎选型Flink的优势和适用边界Spark Streaming和Flink都能做流处理但设计哲学不一样。Spark Streaming的核心是微批Micro-batch把流数据切成一个个小批次每个批次当作一个小的批任务来跑。Flink则是真流式True Streaming每条数据来了就处理计算引擎内部维护连续的状态和窗口。从结果上看Spark Streaming的吞吐也很高但其延迟本质上是受微批间隔限制的——一般设置在2秒到10秒之间很难压进亚秒级。Flink的纯流式处理则可以把端到端延迟做到毫秒级或秒级。我要做的实时订单统计要求延迟不超过10秒其实Spark Streaming也能满足但后来加了实时风控的需求要求3秒内判定异常订单Spark Streaming就有点吃力了。除了延迟Flink还有三个让我最终下决心的点原生支持事件时间和水位线Flink的窗口计算内置了事件时间语义和水位线机制处理乱序数据有现成的API不用自己造轮子。状态管理非常成熟实时计算经常需要维护这个用户过去5分钟的累计消费金额这种状态Flink提供了一整套有状态流处理API状态可以保存在RocksDB里还能定期做检查点持久化。精确一次Exactly-Once语义配合KafkaFlink能做到端到端的精确一次处理核心机制是检查点状态快照 事务性输出。当然Flink也有缺点——学习曲线陡、部署运维比Spark麻烦、JobManager挂了影响面大。但和它解决的核心问题比起来这些成本是值得的。如果项目流量不大、延迟要求也宽松分钟级、团队又只熟Spark生态那Spark Streaming也够用。没有万能的框架只有合适的选型。3. 一个完整的实战案例实时订单金额统计链路选型定了之后我搭了一条完整的实时链路。这个案例现在经常被我拿来当团队的培训demo讲这里完整分享出来包括代码和配置。3.1 场景定义和架构设计需求来自运营部门实时统计全站的支付订单金额、支付成功率、平均客单价按分钟粒度输出数据延迟不超过10秒同时要能实时识别同一用户在短时间内大量下单的异常行为。架构设计如下业务服务/埋点SDK ↓ 埋点日志/业务消息 KafkaTopic: order_event按订单ID哈希分区分发 ↓ Flink Consumer Flink事件时间 滚动窗口1分钟窗口聚合 状态检测 ↓ / ↓ Redis异常识别结果 MySQL/ClickHouse指标结果表为什么中间要加Kafka而不是让Flink直接对接业务服务三个原因一是削峰填谷大促流量高峰期Flink处理不过来的数据可以在Kafka里排队业务服务不会因为计算引擎抖动而被拖死二是解耦业务服务只需要保证消息发到Kafka不关心下游谁在消费三是为了重放Flink任务挂了重启可以从Kafka的offset重放数据把丢失的窗口补算回来。3.2 环境准备和依赖配置Flink跑起来一般有几种方式本地IDE运行开发调试用、Standalone集群、YARN或K8s部署。我在开发阶段用本地模式生产环境用K8s部署。下面给出开发阶段的Maven依赖配置properties flink.version1.17.1/flink.version scala.binary.version2.12/scala.binary.version /properties dependencies !-- Flink流处理核心 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java/artifactId version${flink.version}/version /dependency !-- Flink Kafka连接器 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-kafka/artifactId version${flink.version}/version /dependency !-- Flink窗口和状态处理需要 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-clients/artifactId version${flink.version}/version /dependency !-- 用Gson解析JSON消息 -- dependency groupIdcom.google.code.gson/groupId artifactIdgson/artifactId version2.10.1/version /dependency /dependencies依赖版本有个坑要提醒flink-connector-kafka的版本必须和Flink主版本严格对应而且要关注它内部依赖的Kafka客户端版本。Flink 1.17对应的Kafka连接器支持Kafka 2.x和3.x但如果业务端Kafka集群版本太老比如还是1.0建议先在测试环境用生产环境的Kafka版本验证一遍连接器兼容性别等到上线才发现消息死活消费不进来。3.3 核心代码实现下面是完整的Flink作业代码实现了订单金额统计和异常识别两个功能。import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner; import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.api.common.serialization.SimpleStringSchema; import org.apache.flink.connector.kafka.source.KafkaSource; import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator; import org.apache.flink.streaming.api.environment.CheckpointConfig; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows; import org.apache.flink.streaming.api.windowing.time.Time; import org.apache.flink.api.java.functions.KeySelector; import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction; import org.apache.flink.streaming.api.windowing.windows.TimeWindow; import org.apache.flink.util.Collector; import com.google.gson.JsonObject; import com.google.gson.JsonParser; public class OrderStatisticsJob { public static void main(String[] args) throws Exception { // 1. 创建执行环境并开启Checkpoint StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(5000); // 每5秒做一次检查点 env.getCheckpointConfig().setCheckpointStorage(hdfs:///flink-checkpoints/); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(2000); env.getCheckpointConfig().setCheckpointTimeout(60000); env.getCheckpointConfig().setMaxConcurrentCheckpoints(1); // 任务取消后保留检查点便于手动恢复 env.getCheckpointConfig().setExternalizedCheckpointCleanup( CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION); // 2. 配置Kafka Source KafkaSourceString source KafkaSource.Stringbuilder() .setBootstrapServers(kafka-1:9092,kafka-2:9092,kafka-3:9092) .setTopics(order_event) .setGroupId(flink-order-stat-group) .setStartingOffsets(OffsetsInitializer.latest()) // 生产环境统一从latest开始 .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); DataStreamString stream env.fromSource(source, WatermarkStrategy.noWatermarks(), Kafka Source); // 3. 解析JSON提取事件时间 SingleOutputStreamOperatorOrderEvent orderStream stream .map(json - parseOrderEvent(json)) .assignTimestampsAndWatermarks( WatermarkStrategy.OrderEventforBoundedOutOfOrderness( org.apache.flink.api.common.eventtime.SerializableTimestampAssigner, java.time.Duration.ofSeconds(10)) // 允许10秒乱序 .withTimestampAssigner((event, timestamp) - event.eventTime) ); // 4. 以订单ID为Key状态检测同一用户高频下单这里简化成按用户维度 SingleOutputStreamOperatorInteger anomalyCount orderStream .keyBy((KeySelectorOrderEvent, String) OrderEvent::getUserId) .process(new AnomalyDetectFunction()); // 5. 按事件时间滚动窗口1分钟统计一次全站订单金额 SingleOutputStreamOperatorOrderStatResult statResult orderStream .keyBy((KeySelectorOrderEvent, String) order - all) // 单key全站聚合 .window(TumblingEventTimeWindows.of(Time.minutes(1))) .process(new OrderStatWindowFunction()); // 6. 输出结果这里简化打印到日志生产环境应写回Kafka/MySQL/ClickHouse statResult.print(STAT-STAT-STAT-STAT); anomalyCount.print(ANOMALY); env.execute(Order Real-time Statistics Job); } }因为实际项目里我做了事件时间和乱序处理窗口计算用了ProcessWindowFunction这段逻辑是理解语义的关键单独拆出来看public static class OrderStatWindowFunction extends ProcessWindowFunctionOrderEvent, OrderStatResult, String, TimeWindow { Override public void process(String key, Context context, IterableOrderEvent elements, CollectorOrderStatResult out) { double totalAmount 0.0; int count 0; for (OrderEvent e : elements) { totalAmount e.amount; count; } long windowStart context.window().getStart(); long windowEnd context.window().getEnd(); OrderStatResult result new OrderStatResult( windowStart, windowEnd, System.currentTimeMillis(), count, totalAmount, count 0 ? 0 : totalAmount / count); out.collect(result); } }Context.window()可以直接拿到当前窗口的起止时间这一点在做结果落库时非常有用——因为下游要按窗口时间关联如果这里拿不到窗口边界后面还得想办法推断很麻烦。异常识别函数里我用了Flink的Keyed State判断同一个用户5分钟内下单次数超过阈值就输出告警public static class AnomalyDetectFunction extends KeyedProcessFunctionString, OrderEvent, Integer { private ValueStateInteger orderCountState; private ValueStateLong firstOrderTimeState; Override public void open(Configuration parameters) { ValueStateDescriptorInteger countDesc new ValueStateDescriptor(orderCount, Integer.class); orderCountState getRuntimeContext().getState(countDesc); ValueStateDescriptorLong timeDesc new ValueStateDescriptor(firstOrderTime, Long.class); firstOrderTimeState getRuntimeContext().getState(timeDesc); } Override public void processElement(OrderEvent value, Context ctx, CollectorInteger out) throws Exception { Integer count orderCountState.value(); if (count null) { count 0; firstOrderTimeState.update(value.eventTime); } orderCountState.update(count 1); Long windowStart firstOrderTimeState.value(); // 5分钟时间窗内判断 if (value.eventTime - windowStart 5 * 60 * 1000 count 1 5) { out.collect(value.getUserId()); } // 注册定时器5分钟窗口结束后清空状态 ctx.timerService().registerEventTimeTimer(windowStart 5 * 60 * 1000); } Override public void onTimer(long timestamp, OnTimerContext ctx, CollectorInteger out) throws Exception { firstOrderTimeState.clear(); orderCountState.clear(); } }这里特别要说一下状态清理如果不在onTimer里清空状态Keyed State会一直累积用户量一大内存就爆了。Flink虽然有默认的TTLTime-To-Live机制可以给ValueStateDescriptor设置状态过期时间但在这种时间窗口明确的场景里用事件时间定时器主动清理更精准。整个项目跑起来之后我做了Kafka生产速率和Flink消费速率的对比监控消费吞吐稳定维持在每秒7万条左右窗口延迟在3到6秒之间完全满足业务要求的10秒延迟。4. 上线后我踩过的坑和排查思路实话说案例代码能跑通只是开始真正考验人的是上线后的各种生产问题。我整理了三个最具代表性的坑每个都花了不少时间去排查写出来给大家省点时间。4.1 水位线设置不当导致的延迟问题第一个坑发生在压测阶段数据显示每分钟窗口的产出结果延迟很大有些窗口甚至过了5分钟才出结果。排查过程是这样的——一开始我怀疑是Kafka消费慢或者序列化问题但看监控发现Kafka Lag一直是0消费速率正常说明数据进来没问题问题出在窗口触发的时机上。然后我去查Flink UI上各个窗口的触发时间发现窗口触发普遍比窗口结束时间晚了2到3分钟这就要看水位线了。我当时的代码里用了WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(10))含义是允许数据最多晚到10秒。这个策略会生成一条水位线当前最大事件时间 - 10秒。当一个窗口的结束时间小于水位线时窗口才会触发计算。按这个逻辑窗口最多晚10秒触发才对。但Flink UI上显示的却是最大事件时间一直滞后于系统时间。再深挖发现原因是Kafka里积压了大量历史数据导致Flink刚开始消费时处理的是几天前的事件时间。水位线规则是当前最大事件时间减10秒当消费到的数据事件时间还没追平系统时间时水位线自然会滞后很久窗口也就迟迟不触发。解决办法有两种。第一种是调整Kafka的起始消费策略用setStartingOffsets(OffsetsInitializer.timestamp(...))跳过历史数据直接消费某个业务时间点之后的数据。第二种是给水位线设置一个初始值或者用一个PunctuatedWatermark策略定期推进水位线。我最终的做法是测试环境从最早开始消费压测全量数据但生产环境统一从latest开始避免历史数据影响水位线的追赶。这个坑的核心教训是水位线不是凭空生成的它依赖于已经消费到的数据的事件时间。如果数据源本身就有历史积压水位线一定滞后反过来影响窗口触发。4.2 背压问题Kafka消费被Flink反压拖慢第二个坑更隐蔽某天发现Kafka消费者组Lag持续上涨数据延迟越来越大Flink明明跑着就是消费不动。查了Flink UI上的Backpressure指标看到了最高的算子堆栈水位到99%——这是典型的**背压Backpressure**现象。背压的本质是Flink下游算子处理不过来上游算子的缓冲区被填满然后通过网络和反压机制把压力传导回Kafka Consumer导致消费速率下降。很多人一看到Lag上涨就以为是Kafka的问题其实Kafka是无辜的真正的瓶颈在下游计算。那次我排查发现瓶颈出现在异常检测算子里面状态存储用的是默认的堆内存状态后端TaskManager堆内存高峰期订单量一大状态对象大量堆积JVM频繁Full GC每次GC停顿好几秒整个处理链路就被卡死了。临时止损方案把TaskManager的并行度从12调到32压力立马分散了一些。长期方案把状态后端切换成RocksDB状态数据写到磁盘而不是堆内存GC压力显著下降state.backend: rocksdb state.checkpoints.dir: hdfs:///flink-checkpoints/ taskmanager.memory.managed.fraction: 0.8换完之后背压从99%降到了20%左右Lag也就跟着回落了。后来我把监控面板加上了算子背压指标 Kafka消费Lag GC次数三个维度联动观察再遇到类似问题一看面板就能定位是谁拖慢了谁。我再补充一个实用技巧排查背压时从Flink UI的Backpressure页签可以按算子维度看压力分布如果某个算子一直是High状态优先查看它的CPU使用率和GC情况。如果是CPU高考虑优化算子逻辑如果是GC高优先考虑RocksDB状态后端如果两者都不高再考虑资源并行度不够的问题。4.3 精确一次语义的配置细节第三个坑也是Flink最让人心动的精确一次语义。理论上Flink配合Kafka可以做到端到端Exactly-Once但前提是Sink端也要支持事务性写入。我一开始没注意这个把结果直接写MySQL用了普通的JDBC Sink结果窗口计算结果偶尔会重复写入——检查点恢复之后重复消费Kafka里未确认的消息同一个窗口被处理了两次。排查过程我先怀疑是检查点恢复太频繁看看日志里Task重启记录确实不少——某个算子因为网络抖动重启了两次。然后查看输出表的记录发现确实有主键重复的插入报错。后来查了Flink的官方文档确认要做端到端精确一次需要满足三个条件数据源支持按位点重放Kafka天然支持offset重放。计算引擎使用检查点快照记录状态和source offset。输出端实现基于两阶段提交的事务写入比如Kafka Sink自带的at-least-once和exactly-once模式或者自定义Sink实现TwoPhaseCommitSinkFunction。我当时的MySQL写入是自己实现的Sink只实现了invoke()方法在检查点恢复后重复写入的数据产生了脏数据。改造方案有两个方案一结果表改用去重表设计给每个窗口结果加一个window_start window_end的唯一主键重复写入时用INSERT ... ON DUPLICATE KEY UPDATE覆盖业务上可接受。方案二自定义Sink实现Flink的TwoPhaseCommitSinkFunction接口用事务包裹一批写入检查点完成时提交事务。我后来选了方案一不是因为方案二不好而是因为业务方对结果表里有过重复记录这件事并不敏感只要最终结果正确就行。方案二的复杂度高要处理事务超时、网络故障等一堆边角情况在小团队里维护成本太高。所以精确一次不是银弹要和业务需求结合起来选择实现粒度。5. 压测结果与资源规划经验项目上线前我给团队做了一次完整的压测顺便整理了一套实时流处理的资源规划方法这些经验分享一下。压测使用的是电商大促的流量峰值模板日均订单事件2000万条峰值每秒2万条业务消息加上埋点日志后Kafka的Topic总吞吐峰值大约在每秒5万到8万条之间。我用的Flink集群配置是组件配置TaskManager数量8台每台16核32GB内存并行度32对齐Kafka分区数状态后端RocksDB检查点间隔5秒Kafka分区数32Flink并行度直接设置为Kafka Topic的分区数这样每个并行子任务负责消费一个分区天然做到数据隔离、顺序不乱。在流量上涨想扩容时也要先把Kafka分区数扩上去再调Flink并行度顺序反了会白白浪费资源。压测结果如下稳定吞吐每秒6万条消息端到端延迟4到6秒含窗口聚合时间。峰值吞吐每秒9万条时出现轻微背压但检查点机制下数据不丢失。恢复效果模拟TaskManager宕机重启从检查点恢复耗时52秒恢复后Lag在1分钟内追平期间有一次窗口重复计算被主键去重成功兜住。我把这次压测的结论沉淀成了几条经验第一并行度对齐Kafka分区数是最优起点不是越大越好并行度大了反而增加调度和网络开销第二给Flink TaskManager预留足够的内存做网络缓冲否则高吞吐时网络缓冲区满了也会触发背压第三生产环境的检查点间隔不要设太短我用了5秒频率太高会让检查点本身成为瓶颈。6. 最后说点实在的做了一轮实时数据流处理项目下来最大的感受是这个领域入门门槛看起来不高代码框架封装得很好但真正难的是理解背后的语义——事件时间、水位线、窗口触发时机、检查点恢复机制、背压传播路径每一个概念背后都对应一类线上问题。我踩过的这三个坑本质上都是对某个语义理解不到位造成的。如果要给刚上手的人一个建议我会这样说先在本地把完整链路跑通包括Kafka、Flink、结果落库然后用真实业务数据压测主动制造故障杀掉TaskManager、停掉Kafka、注入乱序数据去观察表现。这些实操经验比看十篇文档都有用。实时流处理没有银弹理解了原理遇到问题基本上都能顺藤摸瓜地解决。
返回列表