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)); } }
关键问题说明与优化点
正确构造SenderRecord
Reactor Kafka的kafkaSender.send()不直接接收字符串,必须传入SenderRecord对象,封装目标Topic、Key、Value等元数据。示例中用UUID生成唯一Key,你可以根据业务逻辑替换为原消息Key或其他业务标识。并行执行逻辑
使用Mono.zip(saveToDb, sendToKafka)实现3A和3B的并行执行,只有当两个操作都成功完成后,才会进入后续流程。如果不需要收集返回结果,也可以用Mono.when(saveToDb, sendToKafka),它仅等待操作完成。消费Offset提交
必须在数据库保存和Kafka发送都成功后,调用receiverOffset.acknowledge()提交Offset,避免消息重复处理。若任意一步失败,Offset不会提交,Kafka会在重平衡后重新投递该消息。错误处理
- 为每个操作添加
doOnError记录详细错误日志,便于排查问题 - 使用
onErrorContinue保证单条消息处理失败时,整个消费流不会中断,继续处理后续消息 - 在
subscribe中添加全局错误处理,捕获流级别的异常
- 额外配置建议
- 目标Topic建议通过
@Value("${kafka.target-topic}")从配置文件注入,避免硬编码 - 可调整
flatMap的并发参数(如flatMap(..., 5))控制并行处理的消息数量 - 确保
KafkaSender的生产者配置正确(如acks、retries等),保证消息可靠性
内容的提问来源于stack exchange,提问作者PatPanda
相关产品推荐
相关产品推荐

