使用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

