【问题标题】:How to create multiple instances of KafkaReceiver in Spring Reactor Kafka如何在 Spring Reactor Kafka 中创建多个 KafkaReceiver 实例
【发布时间】:2021-12-13 05:20:12
【问题描述】:

我有一个反应式 kafka 应用程序,它从一个主题读取数据并写入另一个主题。该主题有多个分区,我想创建与主题中的分区相同数量的消费者(在同一个消费者组中)。据我了解,thread .receive() 只会创建一个 KafkaReceiver 实例,该实例将从主题中的所有分区中读取。所以我需要多个接收器来并行读取不同的分区。

为此,我想出了以下代码:

@Bean
    public ReceiverOptions<String, String> kafkaReceiverOptions(String topic, KafkaProperties kafkaProperties) {
        ReceiverOptions<String, String> basicReceiverOptions = ReceiverOptions.create(kafkaProperties.buildConsumerProperties());
        return basicReceiverOptions.subscription(Collections.singletonList(topic))
                .addAssignListener(receiverPartitions -> log.debug("onPartitionAssigned {}", receiverPartitions))
                .addRevokeListener(receiverPartitions -> log.debug("onPartitionsRevoked {}", receiverPartitions));
    }

    @Bean
    public ReactiveKafkaConsumerTemplate<String, String> kafkaConsumerTemplate(ReceiverOptions<String, String> kafkaReceiverOptions) {
        return new ReactiveKafkaConsumerTemplate<String, String>(kafkaReceiverOptions);
    }

    @Bean
    public ReactiveKafkaProducerTemplate<String, List<Object>> kafkaProducerTemplate(
            KafkaProperties properties) {
        Map<String, Object> props = properties.buildProducerProperties();
        return new ReactiveKafkaProducerTemplate<String, List<Object>>(SenderOptions.create(props));
    }


public void run(String... args) {

        for(int i = 0; i < topicPartitionsCount ; i++) {
            readWrite(destinationTopic).subscribe();
        }
    }}


public Flux<String> readWrite(String destTopic) {
        return kafkaConsumerTemplate
                .receiveAutoAck()
                .doOnNext(consumerRecord -> log.info("received key={}, value={} from topic={}, offset={}",
                        consumerRecord.key(),
                        consumerRecord.value(),
                        consumerRecord.topic(),
                        consumerRecord.offset())
                )
                .doOnNext(consumerRecord -> log.info("Record received from partition {} in thread {}", consumerRecord.partition(),Thread.currentThread().getName()))
                .doOnNext(s-> sendToKafka(s,destTopic))
                .map(ConsumerRecord::value)               
                .onErrorContinue((exception,errorConsumer)->{
                    log.error("Error while consuming : {}", exception.getMessage());
                });
    }

public void sendToKafka(ConsumerRecord<String, String> consumerRecord, String destTopic){
   kafkaProducerTemplate.send(destTopic, consumerRecord.key(), transformRecord(consumerRecord))
                    .doOnNext(senderResult -> log.info("Record received from partition {} in thread {}", consumerRecord.partition(),Thread.currentThread().getName()))
                    .doOnSuccess(senderResult -> {
                        log.debug("Sent {} offset : {}", metrics, senderResult.recordMetadata().offset());
                    }
                    .doOnError(exception -> {
                        log.error("Error while sending message to destination topic : {}", exception.getMessage());
                    })
                    .subscribe();
}

当我对此进行测试时,它似乎工作正常,创建了多个并行处理分区的 Kafka Receiver 实例。 我的问题是,这是创建多个实例的最有效方法吗?反应式 Kafka 中还有其他方法可以做到这一点吗?

【问题讨论】:

    标签: apache-kafka spring-webflux spring-kafka project-reactor reactive-kafka


    【解决方案1】:

    你做的是对的。

    【讨论】:

    • 感谢您的确认。我对这种方法的一个问题是,如果其中一个实例发生故障并停止从其各自的分区中消耗,那么该分区是否会被任何其他实例拾取?我想确保我不会丢失任何数据。
    • 是的;重新平衡将重新分配分区。如果您想停止并重新启动接收器,您可以设置一些属性来延迟重新平衡。
    • 我的理解是,如果未另行指定,重新平衡将自动发生。我的理解正确吗?
    • 没错,是的。
    猜你喜欢
    • 2019-08-10
    • 1970-01-01
    • 1970-01-01
    • 2012-07-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多