为何一个Kafka Streams实例会阻塞另一个启动?Scala Kafka Streams问题
问题根源与解决方案
首先,你的问题核心在于两个KafkaStreams实例共用了同一个APPLICATION_ID_CONFIG,这是Kafka Streams的设计禁忌,直接导致了实例间的资源冲突。
为什么会出现阻塞/无法消费的情况?
Kafka Streams中,APPLICATION_ID_CONFIG是应用的唯一身份标识,它承担了两个关键作用:
- 作为消费者组ID:同一个ID的实例会被归为同一个消费者组,Kafka的消费者组机制会保证每个分区只被组内一个实例消费。当你启动第二个实例时,消费者组会触发重平衡,可能导致其中一个实例无法获取到目标分区的消费权限,自然就处理不了消息。
- 关联状态存储主题:Kafka Streams的状态存储(比如聚合、窗口计算的状态)是存在以这个ID命名的内部主题中的,两个实例同时操作同一批状态主题,会引发读写冲突,进一步干扰流的正常运行。
为什么调换启动顺序会暂时正常?
这只是巧合:第一个启动的实例会先完成消费者组的加入和分区分配,抢占了person主题的分区消费权;第二个启动的实例在重平衡过程中可能没能成功获取到分区,所以只有第一个启动的流能正常工作。但这种状态并不稳定,长时间运行后可能会出现重平衡循环、实例崩溃等问题。
如何解决?
有两种推荐的解决方案:
方案1:为每个KafkaStreams实例分配独立的Application ID
不要共用同一个Properties配置,为每个流单独创建配置并修改APPLICATION_ID_CONFIG:
// 基础配置 val baseConfig: Properties = { val p = new Properties() p.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092") p.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass) p.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass) p } // 为wordSplit流创建独立配置 val wordSplitConfig = new Properties() wordSplitConfig.putAll(baseConfig) wordSplitConfig.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-application-split") // 为JSON处理流创建独立配置 val jsonProcessConfig = new Properties() jsonProcessConfig.putAll(baseConfig) jsonProcessConfig.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-application-json") // 分别使用独立配置创建实例 val streams1 = wordSplit("lines", "wordCount", wordSplitConfig) val streams2 = readAndWriteJson("person", "personName", jsonProcessConfig) // 对应的方法也要修改为接收配置参数 private def wordSplit(intopic: String, outTopic: String, config: Properties) = { val builder = new StreamsBuilderS() val produced = Produced.`with`(Serdes.String(), Serdes.String()) val textLines: KStreamS[String, String] = builder.stream(intopic) val data: KStreamS[String, String] = textLines.flatMapValues(value => value.toLowerCase.split("\\\\W+").toIterable) data.to(outTopic, produced) new KafkaStreams(builder.build(), config) } private def readAndWriteJson(intopic: String, outTopic: String, config: Properties) = { val builder = new StreamsBuilderS() val produced = Produced.`with`(Serdes.String(), Serdes.String()) val textLines: KStreamS[String, String] = builder.stream(intopic) val data: KStreamS[String, String] = textLines.mapValues(value => { val person = Try(parse(value).extract[Person]).toOption println("1::", person) val personNameAndEmail = person.map(a => PersonNameAndEmail(a.name, a.email)) println("2::", personNameAndEmail) write(personNameAndEmail) }) data.to(outTopic, produced) new KafkaStreams(builder.build(), config) }
方案2:合并拓扑到单个KafkaStreams实例(更推荐)
如果你的两个流都是同一个应用的一部分,Kafka Streams支持在同一个StreamsBuilder中定义多个拓扑,然后用单个KafkaStreams实例启动。这种方式可以共享线程资源,避免多实例的冲突和资源浪费:
object Boot extends App { implicit val formats: DefaultFormats.type = DefaultFormats val config: Properties = { val p = new Properties() p.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-application") p.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092") p.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass) p.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass) p } val builder = new StreamsBuilderS() // 构建wordSplit拓扑 val textLines: KStreamS[String, String] = builder.stream("lines") val wordData = textLines.flatMapValues(_.toLowerCase.split("\\\\W+").toIterable) wordData.to("wordCount", Produced.`with`(Serdes.String(), Serdes.String())) // 构建JSON处理拓扑 val jsonLines = builder.stream("person") val jsonData = jsonLines.mapValues(value => { val person = Try(parse(value).extract[Person]).toOption println("1::", person) val personNameAndEmail = person.map(a => PersonNameAndEmail(a.name, a.email)) println("2::", personNameAndEmail) write(personNameAndEmail) }) jsonData.to("personName", Produced.`with`(Serdes.String(), Serdes.String())) // 启动单个KafkaStreams实例 val streams = new KafkaStreams(builder.build(), config) streams.start() Runtime.getRuntime.addShutdownHook(new Thread(() => { streams.close(10, TimeUnit.SECONDS) })) }
这样两个流会在同一个实例中并行处理,既避免了冲突,又能更高效地利用资源。
内容的提问来源于stack exchange,提问作者Shivansh
相关产品推荐
相关产品推荐

