Spring Cloud:如何让Supplier仅发送单次Kafka事件而非重复发送?
问题定位与解决方案
核心问题分析
重复发送事件的根源在于以下几点:
- Supplier未移除已发送元素:
messageSupplier()中使用lists.peek()仅查看队列头部元素,但从未将其从列表中移除。只要列表不为空,每次Supplier触发时都会重复发送同一个元素。 - 列表元素未被清理:代码仅修改了元素的状态
setStatus(Status.SUCESS),但没有从lists中删除该元素,导致列表始终非空,Supplier持续触发发送。 - 冗余逻辑无效:
PublisherService的doOnNext中清理transactionsOfAccount的操作和列表元素清理无关,无法解决重复发送问题。
修复方案
1. 修改Supplier,确保发送后移除元素
将peek()改为poll(),发送完成后自动从列表头部移除元素,避免重复发送:
@Bean public Supplier<Message<Ticker>> messageSupplier() { return () -> { // 使用poll()获取并移除头部元素,而非peek()仅查看 Ticker ticker = tickerPublisher.lists.poll(); if (ticker != null) { Message<Ticker> msg = MessageBuilder .withPayload(ticker) .build(); log.info("Total Size is {}", tickerPublisher.lists.size()); log.info("Message: {}", msg.getPayload()); ticker.setStatus(Status.SUCESS); return msg; } else { return null; } }; }
2. 确保WebClient响应仅添加一次元素到列表
检查sendToKafka方法,确保每个WebClient响应只向lists添加一次元素:
// 示例sendToKafka方法,避免重复添加元素 private Mono<Ticker> sendToKafka(String tickerSymbol, Ticker data) { // 先检查列表中是否已存在该元素,避免重复添加 if (!tickerPublisher.lists.contains(data)) { tickerPublisher.lists.add(data); tickerPublisher.transactionsOfAccount.put(tickerSymbol, data); } return Mono.just(data); }
3. 修正RestController的返回逻辑
使用Sinks.One确保单次响应只被发送一次,避免重复订阅导致多次调用publisherMono:
private final Sinks.One<Ticker> sinkMono = Sinks.one(); @GetMapping(value = "/quote-mono", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Mono<Ticker> getQuoteMono(@RequestParam("symbol") String symbol) { // 仅当sink未被触发时,才调用publisherMono if (!sinkMono.isCancelled() && !sinkMono.isTerminated()) { tickerPublisher.publisherMono(symbol); } return sinkMono.asMono(); }
4. 清理PublisherService的冗余逻辑
移除无关的transactionsOfAccount.clear()操作,改为在元素发送完成后清理对应条目:
public void publisherMono(String ticker) { String path = ticker.toUpperCase() + "/prices/realtime?api_key=" + apiKey; this.webClient .get() .uri(path) .retrieve() .bodyToMono(Ticker.class) .flatMap(data -> sendToKafka(ticker, data)) .doOnNext(data -> { log.info("next events from published : {}", data); // 元素发送后,移除transactionsOfAccount中的对应条目 tickerPublisher.transactionsOfAccount.remove(ticker); }) .subscribe( data -> { log.info("data is {}", data); // 使用tryEmitValue确保仅发送一次 this.sinkMono.tryEmitValue(data); }, (err) -> log.info("Error: {}", err.getMessage()), () -> { log.info("Completed"); } ); }
关键注意事项
- 使用
Sinks.One替代普通Mono,保证单次响应只被触发一次。 - 队列操作优先使用
poll()而非peek(),确保元素发送后从队列中移除。 - 添加元素到队列前做存在性检查,避免重复添加导致的重复发送。
内容的提问来源于stack exchange,提问作者Nirav Kumar Patel
相关产品推荐
相关产品推荐

