本文介绍一种基于 CompletableFuture 与磁盘持久化的长轮询数据同步方案,替代传统 SynchronousQueue#take() 的阻塞等待,避免因客户端断连导致线程挂起、消息丢失及资源泄漏问题。
本文介绍一种基于 `completablefuture` 与磁盘持久化的长轮询数据同步方案,替代传统 `synchronousqueue#take()` 的阻塞等待,避免因客户端断连导致线程挂起、消息丢失及资源泄漏问题。
在典型的长轮询(Long Polling)场景中,服务端需暂存待推送的数据,并在客户端发起请求时“即时”响应——但若直接依赖 SynchronousQueue.take() 等无界阻塞操作,一旦消费者(如 HTTP 请求线程)意外中断(如客户端关闭连接、超时、网络闪断),该线程将永远阻塞或在不可预期时机唤醒,导致:
- 消息被错误消费后丢失(take() 返回却无人处理);
- 线程无法回收,造成线程池耗尽;
- 服务重启后内存中未消费消息彻底丢失。
上述问题的根本症结在于:同步阻塞原语缺乏对消费者生命周期的感知能力。因此,理想的解决方案应满足三点:
✅ 消费者可主动取消等待(支持中断/超时/连接关闭);
✅ 生产者与消费者解耦,支持异步匹配与状态追踪;
✅ 消息具备持久化能力,保障服务宕机不丢数据。
文中提供的实现正是围绕这三点构建的轻量级事件驱动模型:
核心设计:双队列 + CompletableFuture 协同匹配
- consumers 队列:存储 CompletableFuture<LongPollingDto>,代表每个待响应的长轮询请求;
- messages 队列:存储 ContentAndFile(含 DTO 和对应临时文件路径),代表已到达但尚未分发的消息;
- match() 方法:非阻塞地尝试“撮合”可用的消费者与消息,完成 future.complete(dto) 并清理临时文件;
- onCompletion() 回调:在 DeferredResult 完成(无论成功、超时或客户端断开)时主动取消关联的 CompletableFuture,防止后续误触发。
@GetMapping("long-polling")
public DeferredResult<LongPollingDto> longPolling() throws IOException {
DeferredResult<LongPollingDto> result = new DeferredResult<>(30000L); // 可设超时
CompletableFuture<LongPollingDto> futureDto = longPollingService.getFutureDto();
executor.execute(() -> {
try {
result.setResult(futureDto.get(30, TimeUnit.SECONDS)); // 建议显式超时
} catch (TimeoutException e) {
log.warn("Long-polling request timed out");
result.setErrorResult("Timeout");
} catch (InterruptedException | ExecutionException e) {
Thread.currentThread().interrupt();
log.error("Error during future resolution", e);
} catch (CancellationException e) {
log.debug("Request cancelled (client disconnected)");
}
});
// 关键:客户端断连或结果完成时,强制取消 future,避免残留
result.onCompletion(() -> futureDto.cancel(true));
return result;
}持久化保障:文件落地 + 启动恢复
所有待分发消息均序列化为 JSON 文件写入磁盘(如 /tmp/sync/xxx.json),并在 ApplicationReadyEvent 中扫描并重入 messages 队列。这确保了:
- JVM 崩溃后,未消费消息仍可恢复;
- 多实例部署下可通过共享存储(如 NFS)或分布式锁协调,避免重复消费(需额外扩展)。
注意事项与优化建议
- ⚠️ match() 是递归调用,高并发下可能引发栈溢出,建议改为 while 循环;
- ⚠️ ConcurrentLinkedQueue 不保证强一致性,若需严格 FIFO 或批量匹配,可考虑 LinkedBlockingQueue 配合 poll() 循环;
- ✅ future.cancel(true) 会中断其内部线程(如有),但需确保 CompletableFuture 的计算逻辑本身响应中断(例如 HTTP 客户端需配置 readTimeout 和 connectionTimeout);
- ✅ 可结合 Spring WebFlux + Mono 进一步简化异步流管理,降低线程调度复杂度;
- ✅ 生产环境建议增加监控指标:pending_consumers_size、pending_messages_size、disk_usage_percent,便于容量预警。
该方案摒弃了“阻塞即同步”的惯性思维,转而以事件驱动、状态可观察、失败可恢复的设计原则,真正实现了长轮询场景下可靠性与伸缩性的统一。

















