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

Reactor Kafka消费转换消息后并行存库及转发Topic的API问题求助

解决Reactor Kafka发送消息及并行执行问题

修正后的完整代码

package org.example;

import org.example.repository.MyReactiveRepository;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.CommandLineRunner;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.kafka.receiver.KafkaReceiver;
import reactor.kafka.receiver.ReceiverRecord;
import reactor.kafka.receiver.ReceiverOffset;
import reactor.kafka.sender.KafkaSender;
import reactor.kafka.sender.SenderRecord;

import java.util.UUID;

@Service
public class MyConsumerSender implements CommandLineRunner {

    private static final Logger LOGGER = LoggerFactory.getLogger(MyConsumerSender.class);
    // 目标Topic,建议通过@Value从配置文件注入
    private static final String TARGET_TOPIC = "your-target-topic";

    @Autowired
    private KafkaReceiver<String, String> kafkaReceiver;

    @Autowired
    private KafkaSender<String, String> kafkaSender;

    @Autowired
    private MyReactiveRepository myReactiveRepository;

    @Override
    public void run(String... args) {
        consume().subscribe(
                result -> LOGGER.info("消息处理完成: {}", result),
                error -> LOGGER.error("处理流全局异常", error)
        );
    }

    Flux<String> consume() {
        return kafkaReceiver.receive()
                .flatMap(rec -> {
                    // Step 2: 业务转换为大写
                    String transformedString = rec.value().toUpperCase();
                    ReceiverOffset receiverOffset = rec.receiverOffset();

                    // Step 3A: 保存到响应式数据库
                    Mono<String> saveToDb = myReactiveRepository.save(transformedString)
                            .doOnSuccess(saved -> LOGGER.info("已保存到数据库: {}", saved))
                            .doOnError(error -> LOGGER.error("数据库保存失败", error));

                    // Step 3B: 发送到目标Kafka Topic
                    Mono<Void> sendToKafka = kafkaSender.send(
                                    Mono.just(SenderRecord.create(
                                            TARGET_TOPIC,
                                            null,
                                            null,
                                            UUID.randomUUID().toString(), // 按业务生成Key,示例用UUID
                                            transformedString,
                                            null
                                    ))
                            )
                            .doOnSuccess(sendResult -> LOGGER.info("已发送到Kafka Topic: {}", TARGET_TOPIC))
                            .doOnError(error -> LOGGER.error("Kafka发送失败", error))
                            .then(); // 忽略发送结果,转为Void类型Mono

                    // 并行执行3A和3B,完成后提交Offset
                    return Mono.zip(saveToDb, sendToKafka)
                            .doOnSuccess(tuple -> receiverOffset.acknowledge())
                            .map(tuple -> transformedString);
                })
                .onErrorContinue((error, obj) -> LOGGER.error("单条消息处理失败,继续处理后续消息", error));
    }
}

关键问题说明与优化点

  1. 正确构造SenderRecord
    Reactor Kafka的kafkaSender.send()不直接接收字符串,必须传入SenderRecord对象,封装目标Topic、Key、Value等元数据。示例中用UUID生成唯一Key,你可以根据业务逻辑替换为原消息Key或其他业务标识。

  2. 并行执行逻辑
    使用Mono.zip(saveToDb, sendToKafka)实现3A和3B的并行执行,只有当两个操作都成功完成后,才会进入后续流程。如果不需要收集返回结果,也可以用Mono.when(saveToDb, sendToKafka),它仅等待操作完成。

  3. 消费Offset提交
    必须在数据库保存和Kafka发送都成功后,调用receiverOffset.acknowledge()提交Offset,避免消息重复处理。若任意一步失败,Offset不会提交,Kafka会在重平衡后重新投递该消息。

  4. 错误处理

  • 为每个操作添加doOnError记录详细错误日志,便于排查问题
  • 使用onErrorContinue保证单条消息处理失败时,整个消费流不会中断,继续处理后续消息
  • 在subscribe中添加全局错误处理,捕获流级别的异常
  1. 额外配置建议
  • 目标Topic建议通过@Value("${kafka.target-topic}")从配置文件注入,避免硬编码
  • 可调整flatMap的并发参数(如flatMap(..., 5))控制并行处理的消息数量
  • 确保KafkaSender的生产者配置正确(如acks、retries等),保证消息可靠性

内容的提问来源于stack exchange,提问作者PatPanda

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 08:47:46