Spring Webflux实现MongoDB存数成功后发送数据到Kafka
实现方案
你当前使用的是Spring WebFlux响应式编程模型,要保证MongoDB保存成功后才触发Kafka发送,只需要在Mono流中按执行顺序串联操作即可,不要直接在流中间写同步阻塞逻辑。
前置准备
优先注入响应式Kafka生产者模板适配非阻塞流,避免阻塞Netty事件循环线程:
@Resource private ReactiveKafkaProducerTemplate<String, StockDTO> reactiveKafkaTemplate;
如果暂时用普通阻塞式KafkaTemplate,发送逻辑必须调度到独立弹性线程池执行。
代码实现
根据业务对消息可靠性的要求,选以下两种实现方式即可:
场景1:要求Kafka发送成功才返回接口结果(高可靠)
把Kafka发送逻辑作为流的一环串联,Mongo保存成功后才会执行发送,发送完成才返回结果给前端:
public Mono<StockDTO> createStock(StockDTO stockDTONBody) { return mongoTemplate.save(stockDTONBody) .flatMap(savedStock -> reactiveKafkaTemplate.send("stock-create-topic", savedStock.getUuid(), savedStock) .doOnSuccess(sendResult -> log.info("库存数据已发送至Kafka,对应UUID:{}", savedStock.getUuid())) // 配置降级:Kafka发送失败不阻断接口返回,避免Kafka故障影响主流程 .onErrorResume(e -> { log.error("Kafka发送失败,对应库存UUID:{}", savedStock.getUuid(), e); return Mono.empty(); }) .thenReturn(savedStock) ); }
场景2:Mongo保存成功即返回,异步发送Kafka(高吞吐)
适合对吞吐量要求高、可接受极端情况下少量消息丢失的场景,不会等待Kafka发送结果就直接返回:
public Mono<StockDTO> createStock(StockDTO stockDTONBody) { return mongoTemplate.save(stockDTONBody) .doOnSuccess(savedStock -> { // 若使用阻塞KafkaTemplate,需用publishOn切换到弹性线程池再执行发送 reactiveKafkaTemplate.send("stock-create-topic", savedStock.getUuid(), savedStock) .doOnError(e -> log.error("Kafka发送失败,对应库存UUID:{}", savedStock.getUuid(), e)) .subscribe(); }); }
注意事项
- 禁止在响应式流中直接执行无调度的阻塞Kafka发送操作,会导致事件循环线程阻塞,服务吞吐量大幅下降
- 若要求消息绝对不丢失,除了选择场景1的串流写法,还要配置Kafka生产者参数
acks=all、合理的重试次数、开启幂等性 - 建议对Kafka发送失败的消息做本地日志落盘/死信队列处理,方便后续补偿
内容的提问来源于stack exchange,提问作者Prakitidev Verma
相关产品推荐
相关产品推荐

