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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 16:35:24