You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.14 07:24:25