kafka-examples 三种滑动窗口均值实现对比:SimpleMovingAvg、FancyMovingAvg与Kafka Streams

kafka-examples 三种滑动窗口均值实现对比:SimpleMovingAvg、FancyMovingAvg与Kafka Streams kafka-examples 三种滑动窗口均值实现对比SimpleMovingAvg、FancyMovingAvg与Kafka Streams【免费下载链接】kafka-examplesSnippets and small examples demonstrating kafka features and configs项目地址: https://gitcode.com/gh_mirrors/kaf/kafka-exampleskafka-examples 是一个专为新手准备的 Kafka 入门示例仓库用同一个经典任务——滑动窗口移动平均值Moving Average——分别给出了三种由浅入深的实现SimpleMovingAvg最简消费者、FancyMovingAvg处理 Rebalance 的进阶消费者和 KafkaStreamsAvg用 Kafka Streams 框架算窗口聚合。读完本文你就能看懂三种 Kafka 滑动窗口均值实现的架构差异并知道该从哪一个学起。一、为什么滑动窗口均值是最佳 Kafka 入门示例一个能同时练到这些能力的例子就是绝佳的 Kafka 入门教材 消费者轮询poll与消费循环 消费位移offset的提交时机与方式 分区再均衡Rebalance处理⏱️ 窗口计算与状态管理 从裸写消费者到上框架的演进路线这三种实现恰好对应了 Kafka 生态的三个阶段手写消费者 → 手写增强消费者 → Kafka Streams 托管窗口。二、三种实现核心差异速览对比维度SimpleMovingAvgFancyMovingAvgKafkaStreamsAvg技术栈KafkaConsumer原生 APIKafkaConsumer Rebalance 监听Kafka Streams窗口实现CircularFifoBuffer固定长度环形缓冲同左TumblingWindows框架窗口聚合位移提交commitSync()每轮循环同步提交手动管理 mapcommitAsync()异步提交框架自动提交Rebalance 处理❌ 无✅ 分区被撤销时先提交位移框架自动处理状态持久化无重启即清零无✅ 内置状态存储代码量最少中等最少但依赖框架适合人群纯新手想搞懂底层机制的人想直接上生产思路的人三、SimpleMovingAvg最简 Kafka 消费者版SimpleMovingAvg 模块SimpleMovingAvg/里有两个版本正好展示了 Kafka 消费者 API 的历史演进旧版ZK 版SimpleMovingAvg/src/main/java/com/shapira/examples/zkconsumer/simplemovingavg/SimpleMovingAvgZkConsumer.java基于已废弃的 Old Consumer API需要连接 Zookeeper。了解历史即可不建议新学。新版推荐SimpleMovingAvg/src/main/java/com/shapira/examples/newconsumer/simplemovingavg/SimpleMovingAvgNewConsumer.java基于现代KafkaConsumerAPI。新版的套路非常标准值得新手逐行精读用一个容量为 window 的CircularFifoBuffer环形缓冲保存最近 N 个数poll(1000)轮询拉取消息把整数加入缓冲后重算均值每轮循环结束调用commitSync()同步提交位移注册 shutdown hook通过consumer.wakeup()优雅退出——这是处理 Ctrl-C 的经典写法。运行方式见SimpleMovingAvg/run_new_consumer.sh构建 fat-jar 后指定 brokers、group.id、topic 和窗口大小即可启动。 局限很明显窗口状态只活在内存里重启后归零也没处理多实例下的 Rebalance 问题。这正是第二个示例要补的课。四、FancyMovingAvg补齐 Rebalance 与手动提交FancyMovingAvg 的源码只有单文件FancyMovingAvg/src/main/java/com/shapira/examples/fancymovingavg/FancyMovingAvgConsumer.java它在 SimpleMovingAvg 基础上做了三件正经事关闭自动提交显式设置enable.auto.commitfalse位移完全由自己掌控异步提交把每个分区的最新位移存入MapTopicPartition, OffsetAndMetadata用commitAsync()提交并处理回调异常实现ConsumerRebalanceListener在onPartitionsRevoked分区被收走时立刻commitSync避免重复计算退出前的 finally 块再做一次同步兜底提交。源码注释里有一句点睛之笔大意如果你需要这些花活不如直接用 Kafka Streams它都帮你处理好了。——作者本人就在引导你进入下一节。五、KafkaStreamsAvg把窗口计算交给 Kafka StreamsKafkaStreamsAvg 模块KafkaStreamsAvg/展示了完全不同的思路不再自己写轮询循环而是用Kafka Streams 拓扑声明数据流。核心代码在KafkaStreamsAvg/src/main/java/com/shapira/examples/kstreamavg/StreamingAvg.java从ks_prices主题读入价格流与ks_names表leftJoin补全名称用aggregateByKeyTumblingWindows做窗口聚合窗口大小 10 秒结果写入ks_avg_prices主题。聚合逻辑封装在KafkaStreamsAvg/src/main/java/com/shapira/examples/kstreamavg/AvgAggregator.java实现了Aggregator接口的四个方法add累加、remove撤销、merge合并和initialValue中间状态用AvgValuecount sum表示——只存计数和总和两个数窗口均值 sum / count比存整个数组省内存得多。⚠️新手注意该示例使用的是KStreamBuilder这套早期 APIKafka 0.10 之前如今官方推荐StreamsBuildergroupByKey().aggregate()写法。学概念时照旧成立实际编码请以最新官方文档为准。六、三种实现选型建议新手路线第一次写 Kafka 消费者→ 精读 SimpleMovingAvg 新版掌握 poll 循环、提交与优雅退出三板斧️要上多实例生产集群→ 必须看 FancyMovingAvg手动提交 Rebalance 监听是面试和排障的高频考点️做窗口聚合、会话统计等流计算→ 直接上 Kafka StreamsKafkaStreamsAvg 的思路状态托管、容错、扩容全交给框架自己只写聚合逻辑。七、快速上手如何运行示例在 gitcode 上获取仓库代码git clone https://gitcode.com/gh_mirrors/kaf/kafka-examples以 SimpleMovingAvg 为例大致三步进入SimpleMovingAvg/目录执行mvn package打出 fat-jar各模块都带独立pom.xml可单独构建参考SimpleMovingAvg/run_new_consumer.sh中的启动命令替换为你的 brokers 地址、group.id、topic 和窗口大小另开一个生产者向同一 topic 写入整数消息观察控制台输出的滑动均值。仓库里还有 SimpleCounter 之类的生产者示例可以搭配着造数据AvroProducerExample和AvroConsumerExample则演示了 Avro 序列化场景适合作为进阶阅读。八、写在最后kafka-examples 用最朴素的一个需求串起了 Kafka 从手写消费者到流计算框架的完整成长路径。建议按 SimpleMovingAvg → FancyMovingAvg → KafkaStreamsAvg 的顺序逐个读源码每跳一级你对 offset、Rebalance 和状态管理的理解都会深一层——这就是它作为 Kafka 入门示例集合的最大价值。【免费下载链接】kafka-examplesSnippets and small examples demonstrating kafka features and configs项目地址: https://gitcode.com/gh_mirrors/kaf/kafka-examples创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考