如何使用Reactive Kafka异步向指定分区发送字符串数据?
没问题!我来给你一步步讲解如何用Reactive Kafka实现异步发送字符串到指定Kafka分区,结合你已经熟悉的非响应式逻辑,快速上手Rx(Mono/Flux)的用法。
1. 先引入必要的依赖
如果是Spring Boot项目,直接引入Reactive Kafka的starter即可(Maven示例):
<dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka-reactive</artifactId> </dependency>
2. 配置Reactive Kafka生产者
在application.yml里配置Kafka基础信息,和你之前非响应式的配置大部分一致:
spring: kafka: bootstrap-servers: localhost:9092 producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer # 可选:配置ACK级别,比如all确保消息被所有副本确认 acks: all
3. 核心代码实现:异步发送到指定分区
我们用ReactiveKafkaProducerTemplate来封装发送逻辑,这是Reactive Kafka的核心工具类,下面是一个完整的生产者Service示例:
import org.apache.kafka.clients.producer.ProducerRecord; import org.springframework.kafka.core.reactive.ReactiveKafkaProducerTemplate; import org.springframework.stereotype.Service; import reactor.core.publisher.Mono; @Service public class ReactiveKafkaProducer { // 注入Reactive生产者模板 private final ReactiveKafkaProducerTemplate<String, String> producerTemplate; public ReactiveKafkaProducer(ReactiveKafkaProducerTemplate<String, String> producerTemplate) { this.producerTemplate = producerTemplate; } /** * 异步发送字符串到指定Kafka分区 * @param topic 目标主题 * @param partition 指定分区号 * @param message 要发送的字符串消息 * @return Mono<Void> 代表异步发送的完成信号 */ public Mono<Void> sendToSpecificPartition(String topic, int partition, String message) { // 1. 创建指定分区的ProducerRecord,这里key传null(如果不需要key的话) ProducerRecord<String, String> targetRecord = new ProducerRecord<>(topic, partition, null, message); // 2. 异步发送消息,返回Mono<SenderResult> return producerTemplate.send(targetRecord) // 发送成功后的回调:可以记录日志、更新状态等 .doOnSuccess(senderResult -> { System.out.printf("消息发送成功!分区:%d,偏移量:%d%n", senderResult.recordMetadata().partition(), senderResult.recordMetadata().offset()); }) // 发送失败后的回调:处理异常,比如告警、重试等 .doOnError(throwable -> { System.err.printf("消息发送失败:%s%n", throwable.getMessage()); }) // 转换为Mono<Void>,只关注发送完成的信号(忽略返回的元数据) .then(); } }
4. 如何调用这个异步方法
Reactive的核心是订阅触发执行,因为Mono是“冷”的,只有当你调用subscribe()时才会真正发送消息。
普通场景调用(非WebFlux)
比如在一个普通的服务类里调用:
// 注入上面的ReactiveKafkaProducer @Autowired private ReactiveKafkaProducer kafkaProducer; public void someBusinessMethod() { // 异步发送消息,非阻塞 kafkaProducer.sendToSpecificPartition("your-topic", 1, "Hello Reactive Kafka!") .subscribe(); // 必须调用subscribe()才会执行发送逻辑 }
WebFlux场景调用(响应式接口)
如果是WebFlux的Controller,可以直接返回Mono,由框架自动处理订阅:
import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; import reactor.core.publisher.Mono; @RestController @RequestMapping("/kafka") public class KafkaController { private final ReactiveKafkaProducer kafkaProducer; public KafkaController(ReactiveKafkaProducer kafkaProducer) { this.kafkaProducer = kafkaProducer; } @PostMapping("/send") public Mono<Void> sendMessage( @RequestParam String topic, @RequestParam int partition, @RequestParam String message) { // 直接返回Mono,WebFlux会自动订阅执行发送 return kafkaProducer.sendToSpecificPartition(topic, partition, message); } }
几个关键知识点
- 非阻塞异步:和你之前用的
KafkaTemplate.send()返回ListenableFuture不同,Reactive的方式完全基于Reactor的Mono/Flux,不会阻塞当前线程,适合高并发场景。 - 异常处理:除了
doOnError,还可以用onErrorResume来捕获异常并返回替代逻辑,比如重试或者返回默认结果:return producerTemplate.send(targetRecord) .onErrorResume(throwable -> { // 自定义异常处理逻辑,比如重试一次 System.err.println("重试发送消息..."); return producerTemplate.send(targetRecord); }) .then(); - 获取发送结果:如果需要拿到Kafka返回的元数据(分区、偏移量),可以去掉
.then(),直接返回Mono<SenderResult<RecordMetadata>>,调用方就能获取这些信息。
内容的提问来源于stack exchange,提问作者aarish_codev
相关产品推荐
相关产品推荐

