如何合并两个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实例,同样会触发订阅覆盖问题,导致流无法正常拉取数据。
修复方案
并行消费多个主题时,二选一即可:
- 资源占用最低的方案:单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 - 需要独立维护消费流的方案:为每个消费流创建单独的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
相关产品推荐
相关产品推荐

