使用 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,开箱即用地提供map、filter、join、aggregations等最常见的转换操作 | 推荐 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 ID | Artifact ID | 版本 | 说明 |
|---|---|---|---|
org.apache.kafka | kafka-streams | 4.3.0 | (必需)Kafka Streams 基础库 |
org.apache.kafka | kafka-clients | 4.3.0 | (必需)Kafka 客户端库,内置序列化器/反序列化器 |
org.apache.kafka | kafka-streams-scala | 4.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 下,包含StringSerializer、LongSerializer、ByteArraySerializer等常用实现。
使用 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 完整演示了上述全流程:构建StreamsBuilder→builder.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()负责优雅停机与任务迁移,二者都是只能执行一次的生命周期操作,结合setUncaughtExceptionHandler与setStateListener即可构建具备生产级健壮性的流处理应用。
【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考