
别的不说Kafka 都写到第 17 篇了流式处理这块再不动笔就说不过去了。很多朋友把 Kafka 当消息队列用得很溜但一到“流处理”就开始纠结到底是用 Kafka Streams还是上 Flink还是 Spark Streaming这期就把 Kafka Streams 的架构、适用场景和框架选型一次性聊透我自己在项目里踩过的坑和最终的选择思路也会一并放出来。1. 先弄清楚 Kafka Streams 到底是什么1.1 它不是一个独立运行的集群很多人第一次接触 Kafka Streams 时会下意识地把它和 Kafka Broker 搞混或者以为它是一个类似 Flink 那样需要单独部署的集群。这两个理解都不对。Kafka Streams 本质上是一个 Java 库它不启动独立的服务进程而是作为普通 Java 应用的一部分嵌入在你的业务进程里运行。这一点非常关键它决定了 Kafka Streams 的部署形态和运维方式与 Flink、Spark Streaming 有本质区别。你可以把 Kafka Streams 理解成一套“流处理 SDK”。你的应用引入kafka-streams依赖调用它的 API 编写处理逻辑然后像启动一个普通 Spring Boot 服务一样把它跑起来。它自己负责从 Kafka 消费数据、执行处理逻辑、把结果写回 Kafka 或者其他下游系统。整个过程中Kafka Streams 依赖 Kafka 本身来实现数据存储、容错和状态管理所以它不需要额外的调度器或者资源管理器。这里有一个很重要的理念Kafka Streams 是把 Kafka 从单纯的消息中间件升级成了流处理平台。数据源在 Kafka 里数据结果也在 Kafka 里中间的计算逻辑由 Kafka Streams 完成。这种设计简化了整个技术栈——你不需要为流处理再引入一套独立的计算集群而是在已有的 Kafka 基础设施之上就能完成实时计算。1.2 核心架构里的几个关键角色Kafka Streams 的架构基于几个核心概念理解了它们整个框架的运作机制就清晰了。Stream流一个无界、持续更新的数据序列。在 Kafka Streams 里一个 Stream 通常对应一个 Kafka Topic 中的数据。数据是持续流入的没有终点所以叫“无界”。Processor处理器处理流程中的一个计算节点。它接收上游的数据执行业务逻辑比如过滤、转换、聚合然后把结果发给下游 Processor 或 Sink。整个处理拓扑就是由多个 Processor 连接而成的有向无环图DAG。Topology拓扑由 Source Processor数据源、Processor处理节点、Sink Processor输出节点组成的完整处理流程。你可以通过TopologyAPI 手动构建也可以用更高级的Streams DSLDomain Specific Language来描述处理逻辑。DSL 底层会自动帮你构建拓扑。State Store状态存储这是 Kafka Streams 支持有状态计算的关键机制。比如你要按用户ID统计过去一小时的点击量就必须维护“当前累计值”这个状态。Kafka Streams 把状态存储在本地内存或 RocksDB并同步写入一个 Kafka 内部 Topic 中。这样当应用重启或发生故障时可以从 Kafka 中恢复状态保证状态不丢失。说得直白一点Kafka Streams 的核心架构就是“DSL/Topology 定义拓扑 Processor 处理数据 State Store 保存状态 Kafka 承担传输和持久化”。这套设计让它既保持了轻量又没有牺牲流处理应有的能力。1.3 它解决的到底是什么问题在实际项目中我们经常遇到这样的需求实时统计每秒钟的订单量、实时计算用户行为漏斗、持续监测设备上报数据中的异常值、把消息实时同步到数仓等等。传统做法是用 Kafka 做缓冲然后写一个或多个消费者程序把处理逻辑写在消费回调里。这个方案在小规模场景下没毛病但一旦逻辑变复杂问题就来了消费组内多个实例如何分担负载处理过程中需要维护的状态比如计数、最近N条记录放在哪里某个消费者挂了怎么恢复如何保证恰好一次语义这些本质上都是流处理框架要解决的问题。Kafka Streams 就是把这些能力封装好以库的形式提供给你让你不必重复造轮子。2. 使用场景什么时候真的需要 Kafka Streams2.1 实时ETL和数据管道清洗最典型的应用场景之一是实时 ETL。业务系统的原始日志、埋点数据、订单数据源源不断地进入 Kafka但原始数据的格式往往不适合直接给下游使用可能多了一些不需要的字段可能某些字段需要脱敏可能数据格式是 JSON 但下游希望是 Avro。你可以在 Kafka Streams 里写一个处理拓扑完成这些转换工作。比如我之前做过一个实时同步场景上游业务库的 binlog 被 Canal 解析后写入 Kafka但消息里有大量的 DDL 变更记录下游数仓不关心这些只关心真正的数据变更。用 Kafka Streams 写一个filter逻辑把 DDL 消息剔除再对字段名做统一映射转换然后把处理后的数据写入下游 Topic。整个过程毫秒级延迟且不需要额外部署 Flink 集群。这部分要点在于Kafka Streams 的 DSL 提供了map、filter、flatMap、selectKey等算子操作方式非常像 Java 8 的 Stream API上手成本低。对于通过数据清洗需求用 Kafka Streams 比引入一套独立流计算框架要轻量得多。2.2 实时聚合统计与风控监控举个电商场景的例子运营希望实时看到每个商品分类的每分钟销售额。如果数据量大按秒级甚至分钟级做聚合直接在数据库里 count 是不现实的。用 Kafka Streams 可以按窗口tumbling window即滚动窗口对订单流做聚合计算结果持续更新到状态存储中同时输出到下游结果 Topic。风控、异常检测也是典型场景。比如检测某个用户短时间内频繁下单或者某个 IP 短时间内大量请求。这种“时间窗口 计数 阈值判断”的模式Kafka Streams 的窗口聚合 API 正好适合。因为判断是否超过阈值需要保存历史数据这就要依赖 State Store。Kafka Streams 天然支持窗口状态日志留存周期也允许你指定实现这类逻辑基本不需要自己维护额外存储。2.3 事件驱动的实时应用有些场景不是简单地做数据统计而是要基于事件流触发业务动作。比如订单支付成功后触发积分发放、用户连续三次登录失败触发告警、设备状态异常时触发工单创建。这类需求的特点是事件之间通常有关联性需要根据多个事件组合做判断。用 Kafka Streams 处理这类场景的方式是定义处理拓扑用状态存储暂存待匹配的状态事件。比如“用户连续三次登录失败”这个逻辑按用户ID聚合统计在某个时间窗口内的失败次数达到阈值就向告警 Topic 输出一条事件。这样的处理方式天然具备背压和重试机制不需要你自己管理多线程和锁。2.4 不适合的场景不能因为 Kafka Streams 好用就什么都往上面堆。它毕竟是嵌入式库存在一定的局限性需要异构计算大规模如果你有非常复杂的计算逻辑、需要和外部系统大量交互或者是图计算、机器学习等计算密集型任务Kafka Streams 并不是最好的选择Flink 这类专用计算引擎会更合适。需要丰富的外部连接器生态Kafka Streams 的输入输出主要围绕 Kafka 和其 Connect 生态。如果你需要从 MySQL、HDFS、Elasticsearch 等多个系统直接读写数据Flink 的连接器丰富度更高。非 JVM 技术栈团队Kafka Streams 是 Java 库仅支持 JVM 语言。如果团队主力是 Go/Python虽然 Kafka Streams 有对应的轻量级实现比如 Kafka Streams for Python但成熟度和文档完整度远不如 Java 版。这种情况建议使用 Flink 的 Python 接口或者单独写消费者程序。3. 框架选型对比Kafka Streams、Flink、Spark Streaming 怎么挑3.1 三种框架的核心差异很多人在做流处理框架选型时最纠结的就是这三者。我先用一张表格把核心差异列出来再逐条分析。对比维度Kafka StreamsApache FlinkSpark Streaming部署模式嵌入应用进程无独立集群独立集群需 JobManager TaskManager独立集群基于 Spark 计算引擎编程语言Java / Scala / KotlinJava / Scala / PythonJava / Scala / Python / R流与批的模型纯流式流批一体微批处理延迟级别毫秒级毫秒级秒级状态存储本地 RocksDB Kafka TopicRocksDB / 内存 / HDFS 等需配置远端状态后端仅支持有状态算子稳定性依赖 Spark运维成本低跟随应用部署高需维护独立集群中高需维护 Spark 集群生态连接器主要依赖 Kafka Connect丰富支持大量外部系统丰富支持大量外部系统3.2 关键差异的深入解读先从部署模式和使用成本上找感觉。如果你的团队本来就维护着 Kafka 集群那么引入 Kafka Streams 几乎不增加运维成本因为它的“集群”就是 Kafka 本身。而引入 Flink 意味着你要额外修炼一套集群部署、监控、调优的技能。如果项目规模还没有大到需要专职流处理平台的程度Kafka Streams 是性价比极高的选择。再谈流处理语义。Flink 在这个方面名声在外它原生支持精确一次处理Exactly-once语义配合检查点机制能够提供良好的容错保证。Kafka Streams 同样能做到精确一次——它基于 Kafka 的幂等写和事务机制实现。但是需要说明的是精确一次语义是有成本的它带来额外的网络开销和事务管理开销。如果你的业务能容忍极低概率的重复比如实时大屏指标允许5秒内就近重复修正那用“至少一次”语义处理简单多了。最后看状态管理和数据落地。Kafka Streams 把状态存在应用本地磁盘默认是 RocksDB同时异步将 changelog 写入 Kafka 内部主题。因为状态就在本地所以状态读写延迟低状态恢复时则直接从 Kafka 重放 changelog这对中小规模状态来说足够快。Flink 的状态默认为增量检查点状态较大时恢复效率相对更高。但如果你需要跨任务共享状态或者有特别大的状态需求Flink 的伸缩性会更好。3.3 我的选型建议我自己做选型时一般遵循这样一个判断路径数据源和数据结果全部在 Kafka 里业务逻辑中等复杂度团队是 Java 技术栈 → 优先用 Kafka Streams需要从多个异构系统读取数据或需要精确的窗口计算和复杂事件处理 → Flink 更合适已有 Spark 技术栈和集群对延迟要求不是特别高实时任务只是次要在线分析的一部分 → Spark Structured Streaming只是简单的消息过滤、转发、格式转换 → 甚至不需要流处理框架直接用 Kafka Consumer API 写一段代码就够了不要为了用框架而用框架。有时候一段简单的消费者代码能解决的事引入一套流处理框架反而会让系统复杂化。Kafka Streams 的优势恰恰在于“轻”——你只需要在现有 Java 服务里加一个库就能把消费者逻辑进化成完整的流处理拓扑。4. 实操从零构建一个 Kafka Streams 应用4.1 环境准备与依赖引入我以一个最经典的“单词计数”场景为例带大家走一遍完整流程。这个场景虽小但涵盖了 Source、Processor、State Store、Window、Sink 这几个核心环节理解了它后续做复杂任务就有底了。创建 Maven 项目后在pom.xml里引入依赖dependency groupIdorg.apache.kafka/groupId artifactIdkafka-streams/artifactId version3.6.0/version /dependency dependency groupIdorg.slf4j/groupId artifactIdslf4j-simple/artifactId version1.7.36/version /dependency如果是 Spring Boot 项目也可以直接用spring-kafka里对 Kafka Streams 的集成支持在application.yml里配置spring.kafka.streams.*。这两种方式开发体验差不多关键还是理解拓扑构建逻辑。4.2 编写处理拓扑这里我用 Streams DSL 来构建拓扑。它的写法非常直观import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.*; import org.apache.kafka.streams.kstream.*; import java.util.Arrays; import java.util.Properties; public class WordCountDemo { public static void main(String[] args) { Properties props new Properties(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, word-count-app); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, localhost:9092); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass()); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass()); StreamsBuilder builder new StreamsBuilder(); KStreamString, String lines builder.stream(streams-plaintext-input); lines.flatMapValues(line - Arrays.asList(line.toLowerCase().split( ))) .groupBy((key, word) - word) .count() // 无窗口计数会跨批次累加 .toStream() .to(streams-word-count-output, Produced.with(Serdes.String(), Serdes.Long())); Topology topology builder.build(); KafkaStreams streams new KafkaStreams(topology, props); streams.start(); } }简单解读一下这段代码做了什么从名为streams-plaintext-input的 Kafka Topic 中读取数据。把每行文本按空格拆分转成多个单词。按单词本身作为 key 分组。对每组计数结果保存到本地状态存储。把计数结果持续输出到streams-word-count-outputTopic。这段逻辑放到 Flink 里可能需要写不少环境配置和数据源逻辑但 Kafka Streams 的 DSL 让它变得十分简洁。4.3 带窗口的统计实时订单金额统计如果要做不同时间窗口的统计窗口 API 是重点操作对象。以每分钟统计一次所有订单总金额为例核心代码如下import java.time.Duration; SerdeOrderInfo orderSerde new JsonSerde(OrderInfo.class); KStreamString, OrderInfo orders builder.stream(orders, Consumed.with(Serdes.String(), orderSerde)); KTableWindowedString, Long totalAmountPerMinute orders .groupBy((key, order) - order.getCategoryId()) .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(1))) .aggregate( () - 0L, (categoryId, order, total) - total order.getAmount(), Materialized.with(Serdes.String(), Serdes.Long()) );这里的TimeWindows.ofSizeWithNoGrace指定了一个1分钟的滚动窗口每个窗口内的订单金额会被累加。WindowedString类型的 key 包含了“分类ID 窗口起始时间”你可以通过window.key()获取分类IDwindow.window().start()获取窗口开始时间。需要注意窗口类操作默认支持乱序数据可以通过设置grace来配置迟到数据的容忍时长。比如TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5))的语义是允许窗口结束后5分钟内的数据继续进入该窗口计算。如果不配置窗口一到点就立即关闭。4.4 启动注意事项程序启动前记得先创建好输入和输出 Topic。Kafka 默认的auto.create.topics.enable虽然可以自动创建 Topic但分区数可能会和预期不一致最好手动创建kafka-topics.sh --create --topic streams-plaintext-input --partitions 3 --replication-factor 1 --bootstrap-server localhost:9092 kafka-topics.sh --create --topic streams-word-count-output --partitions 3 --replication-factor 1 --bootstrap-server localhost:9092这里先说一个我自己遇到过的坑Topic 的分区数最好和你预期应用的并发实例数保持一致或者有倍数关系。因为 Kafka Streams 默认一个分区对应一个处理线程如果分区数比实例数多负载才会均匀。如果分区数少于实例数一部分实例就会闲置。用控制台生产者发送几条测试数据再启动这个程序然后用控制台消费者查看输出就能看到分词计数的结果了。整个实测链路不需要任何数据库开发和验证成本确实低。5. 实战心得与避坑指南5.1 状态存储的配置很关键Kafka Streams 默认使用 RocksDB 作为本地状态存储。RocksDB 是基于 LSM 树的嵌入式数据库读写性能不错在某些场景下批量查询的性能优于内存存储但它也有内存占用的问题。如果应用部署在内存较小的容器里一定要关注 RocksDB 的内存限制。实际操作中可以通过RocksDBConfigSetter来调整 RocksDB 的内存参数。我一般会限制它最多使用容器内存的 60% 左右避免和其他组件抢内存导致 OOMpublic class CustomRocksDBConfig implements RocksDBConfigSetter { Override public void setConfig(final String storeName, final Options options, final MapString, Object configs) { options.setWriteBufferSize(64 * 1024 * 1024L); // 64MB write buffer options.setMaxWriteBufferNumber(3); options.setTargetFileSizeBase(64 * 1024 * 1024L); } }然后在配置里指定props.put(StreamsConfig.ROCKSDB_CONFIG_SETTER_CLASS_CONFIG, CustomRocksDBConfig.class);5.2 分区数量与线程模型的匹配Kafka Streams 的并行处理能力取决于num.stream.threads配置项。每个线程会分配一部分分区进行消费和处理。如果 Topic 有 6 个分区你在一台机器上启动了num.stream.threads6那么这台机器就能同时处理全部 6 个分区的数据。当应用需要水平扩展时多启动几个应用实例即可。Kafka 的消费组协调机制会自动负责分区在多个实例间的分配。但注意如果你期望两个实例各处理一半负载但 Topic 只有 2 个分区就会有一个实例闲置。规划 Topic 分区数时要结合未来可能的并发实例数一起设计。5.3 优雅关闭和检查点恢复Kafka Streams 应用关闭时如果不做优雅关闭处理可能导致状态保存不完整。推荐在应用里注册 JVM 关闭钩子Runtime.getRuntime().addShutdownHook(new Thread(streams::close));这样应用收到终止信号时会先执行一次状态落盘和提交操作再退出。虽然 Kafka Streams 本身的容错机制会在下次启动时从 changelog 主题恢复状态但优雅关闭能让恢复过程更快也能减少不必要的重复计算。5.4 背压和消费速率Kafka Streams 的消费速率和下游处理逻辑挂钩——如果处理逻辑变慢Consumer 会自动暂停拉取这是天然的背压机制。但如果你的处理逻辑里有外部系统调用比如查 MySQL就需要注意外部系统阻塞可能导致整个拓扑线程阻塞。这里我踩过一个比较深的坑某次在 Processor 内部通过 HTTP 调用下游推荐服务结果那个服务偶发超时居然把整个 KafkaStreams 线程线程池都给堵住了其他分区的数据处理全部停摆。后来改成“本地处理 异步批量上报”的模式把外部调用从处理主链路中摘出去问题才解决。这个经验是Kafka Streams 的拓扑内部最好不要有阻塞式的远程调用如果一定要有务必设置短超时并做好异步降级。5.5 监控指标不能只看消费延迟Kafka Streams 暴露了很多 JMX 指标比如process-rate、commit-rate、state-restore-rate等。我遇到不少朋友只盯着消费延迟却忽略了处理速率。如果你的处理速率小于消费速率积压是迟早的事情甚至可能因为状态过大导致磁盘暴涨。我常用的几个核心监控指标整理成表格方便参考指标名含义健康值参考process-rate每秒处理记录数保持大于生产速率commit-rate每秒提交事务数无持续为0的情况state-restore-rate状态恢复速率重启后逐渐归零rocksdb-block-cache-usageRocksDB块缓存使用量不超过内存上限的80%6. 常见问题与排查技巧6.1 明明往 Topic 发了数据但输出端什么都没有这个问题排查路径比较固定留意一下输入 Topic 和输出 Topic 是否建在了同一个 Kafka 集群上bootstrap.servers配置是否指向正确。这是最常见的低级错误但也是最容易疏忽的。再看应用的application.id是否和之前运行过的旧应用重复。如果重复消费者组用的是旧组的状态新实例可能直接接管了旧实例的分区分配但本地状态没跟上导致数据不处理。这种情况下换个新的application.id就好或者参与消费的计划逻辑保持一致。最后用控制台消费者直接看输入 Topic 里是否有数据。输入 Topic 没有数据一切都是白搭。6.2 状态无法恢复应用反复重启如果应用重启时长时间停留在state-restore阶段大概率要检查 changelog Topic 的数据量。KSTREAM-AGGREGATE-STATE-STORE-0000000001-changelog这类内部 Topic 记录了所有状态变更。如果状态量很大恢复确实需要时间。建议给内部 Topic 设置合理的retention.ms比如 7 天。太短可能导致 state store 需要从最开始的原始数据重新跑一遍太长则浪费存储。另外num.standby.replicas可以设置为 1这样状态在另一个实例上有备份恢复时可以走 standby 而不是从 changelog 全量重放速度明显提升。6.3 OOM 问题Kafka Streams 的 OOM 多数和 RocksDB 内存、堆内存配置不当有关。JVM 的堆主要留给数据处理和序列化对象RocksDB 使用的是堆外内存。如果你给的堆内存过大留给 RocksDB 的堆外空间就小反之亦然。配合前面的 RocksDB 配置类把堆外内存限制控制在合理范围内再用容器内存上限来兜底问题就不容易出现了。6.4 数据倾斜问题聚合类拓扑下有个常见的坑如果 key 的选择不当比如按“城市”分组头部城市的流量远大于其他城市那么某个分区的处理压力就会异常高。这种情况下需要引入滑动分区或者加盐策略来打散热点 key。比如累加统计时先按“城市 随机数%10”分组到了下游再做一轮合并这样既保证正确性又能让负载均衡。当然加盐之后下游要多做一步去盐聚合代码量会增加但在热点 key 场景下这个增加是值得的。Kafka Streams 这套框架说复杂确实有复杂度——状态存储、窗口机制、拓扑调优都有不少细节说简单也确实简单——DSL 的上手门槛比 Flink 要低不少尤其是和现有的 Kafka 基础设施无缝衔接。现在的流处理选型不是非黑即白在同一套技术栈里完全可以并存Kafka Streams 跑轻量级管道和业务状态流Flink 跑重量级实时数仓和复杂事件处理。先把手头的场景摸清楚再决定用哪个这才是做技术选型最靠谱的方式。