【问题标题】:Multiple StreamListeners to same Topic with Spring Cloud Stream connected to Kafka多个 StreamListeners 到同一个主题,Spring Cloud Stream 连接到 Kafka
【发布时间】:2020-10-10 03:32:55
【问题描述】:

我有 Spring Boot 应用程序,我正在使用 Spring Cloud Stream 连接到 Kafka。我正在尝试为同一个 kafka 主题设置两个单独的流侦听器方法。

@StreamListener("countries")
    @SendTo("aggregated-statistic")
    public KStream<?, AggregatedCountry> process(KStream<Object, Country> input) {
        return input
                .groupBy((key, value) -> value.getCountryCode())
                .aggregate(this::initialize,
                        this::aggregateAmount,
                        materializedAsPersistentStore("countries", Serdes.String(),
                                Serdes.serdeFrom(new JsonSerializer<>(),
                                        new JsonDeserializer<>(AggregatedCountry.class))))
                .toStream()
                .map((key, value) -> new KeyValue<>(null, value));
    }
    @StreamListener("countries")
    @SendTo("daily-statistic")
    public KStream<?, List<DailyStatistics>> daily(KStream<Object, Country> input) {
        return input
                .groupBy((key, value) -> value.getCountryCode())
                .aggregate(this::initializeDailyStatistics,
                        this::dailyStatistics,
                        materializedAsPersistentStore("daily", Serdes.String(),
                                Serdes.serdeFrom(new JsonSerializer<>(),
                                        new JsonDeserializer<>(List.class))))
                .toStream()
                .map((key, value) -> new KeyValue<>(null, value));
    }

但是当我启动 Spring Boot 应用程序时出现此错误。

Exception in thread "kafka-stream-f4f8166b-cbeb-42ca-b461-2b3a23885a5d-StreamThread-1" java.lang.IllegalStateException: Consumer was assigned partitions [kafka-stream-daily-repartition-0] which didn't correspond to subscription request [kafka-stream-countries-repartition, countries]
    at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.handleAssignmentMismatch(ConsumerCoordinator.java:218)
    at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.onJoinComplete(ConsumerCoordinator.java:264)
    at org.apache.kafka.clients.consumer.internals.AbstractCoordinator.joinGroupIfNeeded(AbstractCoordinator.java:424)
    at org.apache.kafka.clients.consumer.internals.AbstractCoordinator.ensureActiveGroup(AbstractCoordinator.java:358)
    at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.poll(ConsumerCoordinator.java:353)
    at org.apache.kafka.clients.consumer.KafkaConsumer.updateAssignmentMetadataIfNeeded(KafkaConsumer.java:1251)
    at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1216)
    at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1201)
    at org.apache.kafka.streams.processor.internals.StreamThread.pollRequests(StreamThread.java:963)
    at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:863)
    at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:819)
    at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:788)

我想我需要为每个 StreamListener 方法提供单独的应用程序 ID,但是如果我正在收听相同的主题,如何在 application.yml 文件中配置它?

【问题讨论】:

    标签: apache-kafka spring-kafka spring-cloud-stream


    【解决方案1】:

    您需要提供两个单独的输入绑定(并且它们都可以指向同一个主题)。您不能在多个StreamListeners 上使用相同的绑定名称。然后,您可以在输入绑定上为多个基于 StreamListener 的处理器设置 application.id。例如

    spring.cloud.stream.kafka.streams.bindings.countries1.consumer.applicationId
    

    spring.cloud.stream.kafka.streams.bindings.countries2.consumer.applicationId
    

    请参阅参考文档中的this section

    【讨论】:

    • 感谢您的快速回复。您的解决方案效果很好。
    【解决方案2】:

    您正在阅读“国家”主题两次,如果您从“国家”阅读一次会更好,并将数据发送到“每日统计”和“聚合统计”。

    读取两次与并发处理不同。如果要并发配置这个参数:

    春天: 云.流: 绑定: 国家: 目的地:国家-主题 消费者并发:6

    您可以使用如下拓扑:

    @StreamListener("国家") @SendTo({"每日统计", "汇总统计"})

    【讨论】:

      猜你喜欢
      • 2019-03-19
      • 2017-11-20
      • 2018-12-26
      • 2019-06-28
      • 1970-01-01
      • 2022-07-12
      • 2021-05-08
      • 2020-08-22
      • 1970-01-01
      相关资源
      最近更新 更多