使用 Kafka Streams 编写流处理应用:API 选型、生命周期管理与测试实践
2026/9/11 5:29:25 网站建设 项目流程

使用 Kafka Streams 编写流处理应用:API 选型、生命周期管理与测试实践

【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka

导读

本文基于 Apache Kafka 4.x 仓库中的官方开发者指南,系统讲解如何编写一个 Kafka Streams 应用:从kafka-streams依赖的引入、处理器拓扑(processor topology)的定义,到KafkaStreams实例的启动、优雅关闭与异常处理,再到基于test-utils的单元测试。阅读完成后,你将掌握在 Java/Scala 工程中搭建、运行并测试一个完整 Kafka Streams 应用的全部核心技能,并能对照仓库源码理解每个 API 调用背后的实际行为。

什么是 Kafka Streams 应用

任何使用 Kafka Streams 库的 Java 或 Scala 应用,都被视为 Kafka Streams 应用。Kafka Streams 应用的计算逻辑被定义为一个处理器拓扑(processor topology),它是一张由流处理器(节点)和流(边)构成的图:

  • 节点:流处理器(processor),代表对数据执行的具体计算步骤,例如过滤、映射、聚合;
  • :流(stream),代表数据在处理器之间流动的通道。

这张拓扑图可以用两套 API 来定义:

API定位适用场景
Kafka Streams DSL高层 API,开箱即用地提供mapfilterjoinaggregations等最常见的转换操作推荐 Kafka Streams 新手的起点,能覆盖绝大多数流处理需求;编写 Scala 应用时还可使用 Kafka Streams DSL for Scala 库,省去大量 Java/Scala 互操作样板代码
Processor API低层 API,允许你自行添加和连接处理器,并直接与状态存储(state store)交互需要比 DSL 更高灵活性、但愿意接受更多手工编码(更多代码行)的场景

从源码结构看,这两套 API 分别对应仓库中 streams/src/main/java/org/apache/kafka/streams/StreamsBuilder.java(DSL 的入口,stream()等方法在此构建流)与 streams/src/main/java/org/apache/kafka/streams/processor 包下的Topology类(Processor API 的拓扑描述)。无论用哪套 API,最终都会得到一份可执行的Topology描述。

库与 Maven 依赖

Kafka Streams 相关的库在 Maven 坐标与作用如下(当前仓库对应版本为4.3.0):

Group IDArtifact ID版本说明
org.apache.kafkakafka-streams4.3.0(必需)Kafka Streams 基础库
org.apache.kafkakafka-clients4.3.0(必需)Kafka 客户端库,内置序列化器/反序列化器
org.apache.kafkakafka-streams-scala4.3.0(可选)用于编写 Scala Kafka Streams 应用的 DSL 库;不使用 SBT 时,需在 artifact ID 后追加对应 Scala 版本后缀(_2.12_2.13

关于序列化器/反序列化器(Serdes)的选型细节,可进一步阅读数据类型与序列化。kafka-clients中内置的序列化器位于仓库 clients/src/main/java/org/apache/kafka/common/serialization 下,包含StringSerializerLongSerializerByteArraySerializer等常用实现。

使用 Maven 时的pom.xml依赖片段示例:

<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-streams</artifactId> <version>4.3.0</version> </dependency> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>4.3.0</version> </dependency> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-streams-scala_2.13</artifactId> <version>4.3.0</version> </dependency>

在应用代码中使用 Kafka Streams

你可以在应用代码的任何位置调用 Kafka Streams,但通常这些调用发生在main()方法(或其变体)中。定义处理拓扑的基本要素如下。

第一步:创建 KafkaStreams 实例

首先必须创建一个KafkaStreams实例,其构造函数的两个核心参数是:

  • 第一个参数:拓扑对象——DSL 场景下为StreamsBuilder#build()的返回值,Processor API 场景下为Topology对象;
  • 第二个参数java.util.Properties实例,定义该拓扑的专属配置(集群地址、默认序列化器、安全设置等)。
import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.kstream.StreamsBuilder; import org.apache.kafka.streams.processor.Topology; // 使用 builder 定义实际的处理拓扑:从哪些输入主题读取、 // 调用哪些流操作(filter、map 等),后续章节会详细展开。 StreamsBuilder builder = ...; // 使用 DSL 时 Topology topology = builder.build(); // // 或者 // Topology topology = ...; // 使用 Processor API 时 // 通过配置告诉应用:Kafka 集群在哪、默认使用什么序列化器、 // 安全设置如何、等等。 Properties props = ...; KafkaStreams streams = new KafkaStreams(topology, props);

从 KafkaStreams.java 的源码可以看到,KafkaStreams提供了多个重载构造函数,除了(Topology, Properties)之外,还可以传入自定义的TopologyMetadata等扩展参数;构造完成后内部结构即完成初始化,但处理尚未开始。

第二步:显式启动流线程

构造完成后处理并不会自动开始,必须显式调用KafkaStreams#start()来启动 Kafka Streams 线程:

// 启动 Kafka Streams 线程 streams.start();

在 KafkaStreams.java 的实现中,start()会先将实例状态置为REBALANCING,清理过期状态目录、初始化本地状态存储,然后依次启动全局线程(如果存在)与所有流线程,并记录实际启动的线程数。需要说明的几点:

  • start()是异步非阻塞的,它在后台启动线程后立即返回;但如果拓扑包含全局状态存储(global stores),该方法会阻塞直到所有全局存储恢复完成。
  • start()只能调用一次,重复调用会抛出IllegalStateException
  • Broker 兼容性约束(来自源码注释):Kafka Streams 4.x 要求 broker 版本不低于 2.1,否则连接会因协议版本不支持而失败;当processing.guarantee设置为exactly_once_v2时,broker 必须是 2.5 或更高版本,否则应用会在首次 rebalance 时检测到并转入ERROR状态。这类兼容性问题通常在start()返回后异步暴露,建议通过状态监听器或异常处理器感知。

如果在其他地方还有该流处理应用的实例在运行(例如另一台机器上),Kafka Streams 会自动将任务从现有实例重新分配到刚启动的新实例上。相关机制参见流分区与任务与线程模型。

第三步:捕获未捕获异常

为了捕获任何意外异常,可以在启动应用之前设置java.lang.Thread.UncaughtExceptionHandler。每当流线程因意外异常而终止时,该处理器就会被调用:

streams.setUncaughtExceptionHandler((Thread thread, Throwable throwable) -> { // 在这里检查 throwable/exception,并执行适当的处理动作! });

在 Kafka Streams 4.x 中,setUncaughtExceptionHandler的签名实际接收的是StreamsUncaughtExceptionHandler(其handle方法返回StreamThreadExceptionResponse枚举,指示应用在异常后应继续、替换线程还是关闭)。从源码可以确认:

  • 该处理器只能在不晚于start()之前设置,否则抛出IllegalStateException
  • 处理器必须是线程安全的,因为它会被所有内部线程共享,并可能从任何遭遇异常的线程上被调用;
  • 全局线程(global thread)与普通流线程的异常都会路由到该处理器。

更完整的故障感知还可以借助状态监听器setStateListener,它会在实例状态变化时回调onChange(newState, oldState),例如用于感知 broker 兼容性失败导致的ERROR状态。

第四步:停止应用与优雅关闭

停止应用实例时调用KafkaStreams#close()

// 停止 Kafka Streams 线程 streams.close();

源码中close()(KafkaStreams.java)会通知所有线程停止并等待其 join,是一个阻塞调用。关闭行为会依据当前使用的分组协议自适应:经典协议下 consumer 留在消费组内;使用 Streams 协议(group.protocol=streams)时,动态成员会主动离组,而配置了group.instance.id的静态成员继续留在组内、由 broker 在会话超时后移除。

为了让应用响应 SIGTERM 实现优雅关闭,官方推荐添加 shutdown hook 并在其中调用KafkaStreams#close()。Java 示例:

// 添加 shutdown hook 以停止 Kafka Streams 线程。 // 也可以选择为 close 提供超时时间。 Runtime.getRuntime().addShutdownHook(new Thread(streams::close));

应用停止后,Kafka Streams 会将该实例上运行的所有任务迁移到剩余可用实例上。

一个完整的可运行示例

仓库 streams/examples/src/main/java/org/apache/kafka/streams/examples/wordcount/WordCountDemo.java 完整演示了上述全流程:构建StreamsBuilderbuilder.build()new KafkaStreams(...)→ 注册 shutdown hook →streams.start()→ 通过CountDownLatch阻塞主线程。其核心拓扑代码如下:

final StreamsBuilder builder = new StreamsBuilder(); final KStream<String, String> source = builder.stream(INPUT_TOPIC); final KTable<String, Long> counts = source .flatMapValues(value -> Arrays.asList(value.toLowerCase(Locale.getDefault()).split("\\W+"))) .groupBy((key, value) -> value) .count(); // 需要覆盖 value 的序列化器为 Long 类型 counts.toStream().to(OUTPUT_TOPIC, Produced.with(Serdes.String(), Serdes.Long()));

对应的运行配置(WordCountDemo.java)展示了几个关键配置项及其默认值:

props.putIfAbsent(StreamsConfig.APPLICATION_ID_CONFIG, "streams-wordcount"); props.putIfAbsent(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.putIfAbsent(StreamsConfig.STATESTORE_CACHE_MAX_BYTES_CONFIG, 0); props.putIfAbsent(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.StringSerde.class); props.putIfAbsent(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.StringSerde.class); props.putIfAbsent(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");

其中application.id是该应用在 Kafka 集群中的唯一标识(也是内部主题与消费组命名的前缀),bootstrap.servers指向集群地址,default.key/value.serde指定默认序列化器,auto.offset.reset=earliest保证可基于预置数据重复运行演示。运行前需先用bin/kafka-topics.sh创建输入主题streams-plaintext-input与输出主题streams-wordcount-output,并用bin/kafka-console-producer.sh写入数据。

测试 Streams 应用

Kafka Streams 自带test-utils模块(对应仓库 streams/test-utils 目录),用于辅助测试流处理应用,详细用法参见测试指南。

test-utils提供了一组开箱即用的测试双件,可以做到不依赖真实 Kafka 集群即可驱动拓扑运行并断言结果:

  • TopologyTestDriver:在单进程中直接执行拓扑,模拟 broker 行为;
  • TestInputTopic/TestOutputTopic:向输入主题写入测试数据、从输出主题读取并断言结果,支持按时间戳推进(advanceTime)以测试窗口、会话等时间相关逻辑。

仓库中的TopologyTestDriver位于 streams/test-utils/src/main/java/org/apache/kafka/streams/TopologyTestDriver.java,TestInputTopic位于 streams/test-utils/src/main/java/org/apache/kafka/streams/TestInputTopic.java。在其 Maven 依赖中需要引入:

<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-streams-test-utils</artifactId> <version>4.3.0</version> <scope>test</scope> </dependency>

小结

编写一个 Kafka Streams 应用的完整链路可以概括为:引入kafka-streams依赖 → 用 DSL 或 Processor API 构建Topology→ 构造并start()KafkaStreams→ 设置异常处理器与 shutdown hook → 用test-utils验证拓扑逻辑。在整个生命周期中,start()负责异步拉起流线程并触发任务分配,close()负责优雅停机与任务迁移,二者都是只能执行一次的生命周期操作,结合setUncaughtExceptionHandlersetStateListener即可构建具备生产级健壮性的流处理应用。

【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询