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

使用ConcurrentLinkedQueue与@Scheduled时数据丢失问题求助

排查定时聚合发送数据时的随机丢失问题

我开发的RESTful应用需处理海量数据,通过@Scheduled(fixedRateString = "2000")每2秒聚合队列中的数据发送至外部服务器,但出现随机数据未发送的情况,以下是核心代码:

@Scheduled(fixedRateString = "2000")
void sendToApi() {
    var apiRequest = prepareRequests();
    webClient.sendRequests(apiRequest);
}
private final ConcurrentLinkedQueue<ApiRequest> apiRequests = new ConcurrentLinkedQueue<>();

private ApiRequest prepareApiRequest() {
    List<ApiRequest> requestsToSend = new ArrayList<>();
    for (int i = 0; i < maxRequests ; i++) {
        if (this.apiRequests.peek() == null) {
            break;
        }
        requestsToSend.add(this.apiRequests.poll());
    }
    return new ApiRequest(requestsToSend);
}

可能的原因及排查步骤

  • 方法名不匹配导致逻辑未执行:
    注意到sendToApi中调用的是prepareRequests(),但实际定义的聚合方法是prepareApiRequest(),这会直接导致编译错误或运行时找不到方法,聚合逻辑根本没执行。如果是笔误,先修正方法名,确保聚合逻辑正常调用。

  • 队列操作的竞态条件导致数据丢失:
    当前代码先通过peek()判断队列是否有元素,再调用poll()取出。在多线程环境下,这两步之间存在时间窗:线程A执行peek()发现有元素,此时线程B调用poll()把元素取走,线程A再执行poll()会返回null,最终requestsToSend中会混入null元素。如果sendRequests方法忽略null或处理null时抛出异常,就会导致对应数据丢失。
    优化方式:直接用poll()判断,去掉peek():

    private ApiRequest prepareApiRequest() {
        List<ApiRequest> requestsToSend = new ArrayList<>();
        for (int i = 0; i < maxRequests ; i++) {
            ApiRequest req = this.apiRequests.poll();
            if (req == null) {
                break;
            }
            requestsToSend.add(req);
        }
        return new ApiRequest(requestsToSend);
    }
    
  • 未捕获异常导致数据丢失:
    sendToApi方法没有任何异常捕获逻辑,如果prepareApiRequest或webClient.sendRequests()抛出异常(比如网络超时、序列化错误),当前批次从队列取出的数据会因为未发送成功且没有重试/回滚逻辑,直接丢失。
    解决:添加异常捕获,失败时将数据放回队列(注意避免死循环),同时记录日志:

    @Scheduled(fixedRateString = "2000")
    void sendToApi() {
        ApiRequest apiRequest = null;
        try {
            apiRequest = prepareApiRequest();
            webClient.sendRequests(apiRequest);
        } catch (Exception e) {
            // 日志记录异常详情
            log.error("发送数据失败", e);
            // 将未成功发送的数据放回队列头部,避免丢失
            if (apiRequest != null && apiRequest.getRequests() != null) {
                apiRequest.getRequests().forEach(req -> apiRequests.addFirst(req));
            }
        }
    }
    
  • WebClient异步发送未等待完成:
    如果webClient.sendRequests()是异步调用(比如返回Mono/Flux但未调用block()),定时任务会直接结束,此时异步请求可能还未完成,若后续出现异常(如网络中断),无法感知和处理,导致数据丢失。
    排查:确认sendRequests是否同步执行,若是异步,需等待请求完成并处理结果:

    // 假设sendRequests返回Mono<Void>
    webClient.sendRequests(apiRequest).block();
    
  • maxRequests设置不合理:
    如果maxRequests值过大,单次从队列取出过多数据,导致sendRequests处理超时或内存溢出,进而引发异常,未发送的数据直接丢失。建议根据外部服务器的并发限制和自身处理能力,合理设置maxRequests的大小。

  • 定时任务线程池阻塞:
    Spring默认的@Scheduled线程池是单线程的,如果某次sendToApi执行时间超过2秒(比如网络慢导致发送耗时久),下一次任务会等待前一次完成。若期间队列数据持续堆积,可能导致后续处理时出现异常,或部分数据因积压过久被其他逻辑影响丢失。
    解决:自定义定时任务线程池,增加线程数:

    @Configuration
    public class SchedulerConfig implements SchedulingConfigurer {
        @Override
        public void configureTasks(ScheduledTaskRegistrar taskRegistrar) {
            ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler();
            scheduler.setPoolSize(5);
            scheduler.setThreadNamePrefix("data-sender-");
            scheduler.initialize();
            taskRegistrar.setTaskScheduler(scheduler);
        }
    }
    
  • 队列数据添加逻辑异常:
    排查向apiRequests队列添加数据的代码,确认是否存在线程安全问题(虽然ConcurrentLinkedQueue是线程安全的,但如果添加时抛出异常,元素可能未被正确加入队列),或者是否有重复添加/误删除的情况。

内容的提问来源于stack exchange,提问作者J.Doe

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 02:28:25