【问题标题】:Multi threading on Kafka Send in Spring reactor KafkaKafka 上的多线程在 Spring reactor Kafka 中发送
【发布时间】:2021-11-10 04:06:03
【问题描述】:

我有一个反应式 kafka 应用程序,它从一个主题读取数据、转换消息并写入另一个主题。我在主题中有多个分区,因此我创建了多个消费者以并行读取主题。每个消费者在不同的线程上运行。但看起来 kafka send 运行在同一个线程上,即使它是从不同的消费者调用的。 我通过记录线程名称进行测试以了解线程工作流程,每个消费者的接收线程名称不同,但在 kafka 发送 [kafkaProducerTemplate.send] 上,所有消费者的线程名称 [线程名称:producer-1] 都是相同的.我不明白它是如何工作的,我希望它对发送的所有消费者也有所不同。有人可以帮助我了解这是如何工作的。

@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();
}

【问题讨论】:

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


    【解决方案1】:

    生产者的所有发送都在单线程调度程序上运行(通过.publishOn())。

    DefaultKafkaSender.doSend()

    您应该为每个消费者创建一个发送者。

    【讨论】:

    • 你能帮我理解工作流程吗?每个消费者都在调用 kafkaProducerTemplate.send(),那不会创建不同的 sender 实例吗?
    • 否;每个模板都有一个发件人。您需要为每个消费者创建一个单独的模板,或者创建一个池并签出一个实例并在完成后将其返回到池中。
    • 感谢您的解释。重申我的理解:在我的代码中,KafkaConsumerTemplate 有多个消费者实例,因为 .receive() 创建了不同的接收器,并且每个接收器将订阅不同的分区。但是只有一个发送者与 KafkaProducerTemplate 相关联,因此必须为每个消费者创建不同的模板。我希望 .send() 的工作方式类似于 .receive(),但看起来我错了。你能确认这是正确的吗?
    • 模板是围绕响应式发送方和接收方的非常轻量级的包装器;实际功能由该库(reactor-kafka)提供。是的,它看起来确实有点不对称,但如果你仔细想想,receive() 是一个长时间运行的操作,而send() 通常非常短(异步)。对于大多数应用程序来说,只需要一个生产者即可。请参阅 KafkaProducer ...a single producer instance across threads will generally be faster... 的 javadocs。如果您使用事务,则必须使用单独的生产者,因为一次只能激活一个事务。
    猜你喜欢
    • 1970-01-01
    • 2021-12-13
    • 2019-08-10
    • 1970-01-01
    • 1970-01-01
    • 2020-04-25
    • 2021-06-12
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多