You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何合并两个Kafka主题的ZStream为单流?并行消费无输出问题

ZIO Kafka 多主题并行/合并消费无输出问题

问题现象

使用 zio-kafka 同时消费test、topic两个 Kafka 主题时,无论是做流合并还是并行 fork 两个消费流,控制台都没有任何输出;但单独消费任意一个主题时功能完全正常。

复现代码

import org.slf4j.LoggerFactory
import zio.blocking.Blocking
import zio.clock.Clock
import zio.console.{Console, putStrLn}
import zio.kafka.consumer.{CommittableRecord, Consumer, ConsumerSettings, Subscription}
import zio.kafka.consumer.Consumer.{AutoOffsetStrategy, OffsetRetrieval}
import zio.kafka.serde.Serde
import zio.stream.ZStream
import zio.{ExitCode, Has, URIO, ZIO, ZLayer}

object Test2Topics extends zio.App {
  val logger = LoggerFactory.getLogger(this.getClass)

  val consumerSettings: ConsumerSettings =
    ConsumerSettings(List("localhost:9092"))
      .withGroupId(s"consumer-${java.util.UUID.randomUUID().toString}")
      .withOffsetRetrieval(OffsetRetrieval.Auto(AutoOffsetStrategy.Earliest))

  val consumer: ZLayer[Clock with Blocking, Throwable, Has[Consumer]] =
    ZLayer.fromManaged(Consumer.make(consumerSettings))

  val streamString: ZStream[Any with Has[Consumer], Throwable, CommittableRecord[String, String]] =
    Consumer.subscribeAnd(Subscription.topics("test"))
      .plainStream(Serde.string, Serde.string)

  val streamInt: ZStream[Any with Has[Consumer], Throwable, CommittableRecord[String, String]] =
    Consumer.subscribeAnd(Subscription.topics("topic"))
      .plainStream(Serde.string, Serde.string)
  val combined = streamString.zipWithLatest(streamInt)((a,b)=>(a,b))
  val program = for {
    fiber1 <- streamInt.tap(r => putStrLn(s"streamInt: ${r.toString}")).runDrain.forkDaemon
    fiber2 <- streamString.tap(r => putStrLn(s"streamString: ${r.toString}")).runDrain.forkDaemon
  } yield ZIO.raceAll(fiber1.join, List(fiber2.join))
  override def run(args: List[String]): URIO[zio.ZEnv, ExitCode] = {
    //combined.tap(r => putStrLn(s"Combined: ${r.toString}")).runDrain.provideSomeLayer(consumer ++ Console.live).exitCode
    program.provideSomeLayer(consumer ++ Console.live).exitCode
  }
}

问题根因

  • 单个Consumer实例的行为和原生 Kafka Consumer 完全一致,同一时间只能维护一组订阅关系:代码里两次调用Consumer.subscribeAnd,后执行的订阅会直接覆盖前一次的订阅配置,最终只会有一个主题被实际监听,另一个流永远拿不到数据。
  • 代码中使用ZIO.raceAll等待fiber,只要任意一个fiber完成就会终止整个程序。Kafka 消费流本身是常驻运行的,正常情况下永远不会完成,但订阅被覆盖后其中一个流会提前退出,直接触发整个程序终止,不会打印任何消费日志。
  • zipWithLatest操作要求两个流共享同一个Consumer实例,同样会触发订阅覆盖问题,导致流无法正常拉取数据。

修复方案

并行消费多个主题时,二选一即可:

  1. 资源占用最低的方案:单Consumer一次订阅所有主题,拿到流之后按记录的topic名做分流处理
    // 一次订阅两个主题
    val allTopicsStream = Consumer.subscribeAnd(Subscription.topics("test", "topic"))
      .plainStream(Serde.string, Serde.string)
    
    val program = allTopicsStream
      .tap(record => 
        if (record.record.topic() == "test") putStrLn(s"streamString: ${record.toString}")
        else putStrLn(s"streamInt: ${record.toString}")
      )
      .runDrain
    
  2. 需要独立维护消费流的方案:为每个消费流创建单独的Consumer实例,不要复用同一个Consumer做多次订阅
    // 为两个主题创建独立的Consumer实例
    val testTopicConsumer = ZLayer.fromManaged(Consumer.make(consumerSettings))
    val intTopicConsumer = ZLayer.fromManaged(Consumer.make(consumerSettings))
    
    val program = for {
      fiber1 <- streamInt.tap(r => putStrLn(s"streamInt: ${r.toString}")).runDrain
        .provideSomeLayer(intTopicConsumer ++ Console.live)
        .forkDaemon
      fiber2 <- streamString.tap(r => putStrLn(s"streamString: ${r.toString}")).runDrain
        .provideSomeLayer(testTopicConsumer ++ Console.live)
        .forkDaemon
      // 用join等待两个常驻流持续运行,不要用raceAll
      _ <- fiber1.join
      _ <- fiber2.join
    } yield ()
    

内容的提问来源于stack exchange,提问作者eprst2019

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.27 19:30:17