
使用 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.javaDSL 的入口stream()等方法在此构建流与 streams/src/main/java/org/apache/kafka/streams/processor 包下的Topology类Processor API 的拓扑描述。无论用哪套 API最终都会得到一份可执行的Topology描述。库与 Maven 依赖Kafka Streams 相关的库在 Maven 坐标与作用如下当前仓库对应版本为4.3.0Group 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 下包含StringSerializer、LongSerializer、ByteArraySerializer等常用实现。使用 Maven 时的pom.xml依赖片段示例dependency groupIdorg.apache.kafka/groupId artifactIdkafka-streams/artifactId version4.3.0/version /dependency dependency groupIdorg.apache.kafka/groupId artifactIdkafka-clients/artifactId version4.3.0/version /dependency dependency groupIdorg.apache.kafka/groupId artifactIdkafka-streams-scala_2.13/artifactId version4.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.protocolstreams时动态成员会主动离组而配置了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 KStreamString, String source builder.stream(INPUT_TOPIC); final KTableString, 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.resetearliest保证可基于预置数据重复运行演示。运行前需先用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.javaTestInputTopic位于 streams/test-utils/src/main/java/org/apache/kafka/streams/TestInputTopic.java。在其 Maven 依赖中需要引入dependency groupIdorg.apache.kafka/groupId artifactIdkafka-streams-test-utils/artifactId version4.3.0/version scopetest/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),仅供参考