无分组键时,Java批量处理中间件可用的线程安全数据结构咨询
无分组键下的批量请求处理方案
当没有特定分组键时,所有请求共享同一个批次池,核心思路是用线程安全的队列存储请求,结合数量触发和定时触发实现批量发送,同时通过原子操作和兜底逻辑保证请求不丢失、不重复。
核心组件与实现逻辑
1. 线程安全的请求存储
用LinkedBlockingQueue<Request>作为请求队列,它本身提供线程安全的入队/出队操作,支持阻塞式写入,避免请求丢失。
2. 双重触发机制
- 数量触发:每次提交请求后检查队列大小,达到10个时立即触发批量处理。
- 定时触发:用单线程定时任务每隔100ms检查队列,有请求就批量取出发送。
3. 原子性批量取数
使用队列的drainTo方法原子性取出最多10个请求,保证同一批请求不会被多个线程重复处理。
4. 兜底保障逻辑
- 提交请求时若线程中断,直接发送单个请求避免丢失。
- 批量发送失败时,将请求重新入队(带重试次数限制)或持久化到可靠存储,确保请求不丢失。
代码实现(Java)
import java.util.ArrayList; import java.util.List; import java.util.concurrent.*; class Request { // 自定义请求字段,比如请求参数、重试次数等 private int retryCount = 0; public int getRetryCount() { return retryCount; } public void incrementRetryCount() { retryCount++; } } public class BatchRequestProcessor { private static final int BATCH_SIZE = 10; private static final long TIMEOUT_MS = 100; private static final int MAX_RETRY = 3; // 最大重试次数 private final BlockingQueue<Request> requestQueue = new LinkedBlockingQueue<>(1000); // 限制队列容量,避免内存溢出 private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); private final ExecutorService workerPool = Executors.newFixedThreadPool(2); // 可根据后端并发能力调整 public BatchRequestProcessor() { // 启动定时任务,首次延迟100ms后每隔100ms执行一次 scheduler.scheduleAtFixedRate(this::processBatch, TIMEOUT_MS, TIMEOUT_MS, TimeUnit.MILLISECONDS); } // 对外暴露的请求提交方法,线程安全 public void submit(Request request) { try { requestQueue.put(request); // 检查是否达到批量阈值,触发处理 if (requestQueue.size() >= BATCH_SIZE) { workerPool.submit(this::processBatch); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); // 提交失败,直接发送单个请求兜底 sendSingleRequest(request); } } // 批量处理核心逻辑 private void processBatch() { List<Request> batch = new ArrayList<>(BATCH_SIZE); // 原子性取出最多BATCH_SIZE个请求,避免重复处理 requestQueue.drainTo(batch, BATCH_SIZE); if (!batch.isEmpty()) { sendBatchToBackend(batch); } } // 批量发送到后端 private void sendBatchToBackend(List<Request> batch) { try { // 这里替换为实际后端调用逻辑 // backendService.sendBatch(batch); System.out.println("Sent batch of " + batch.size() + " requests"); } catch (Exception e) { // 发送失败,处理重试或持久化 handleBatchFailure(batch); } } // 处理批量发送失败的请求 private void handleBatchFailure(List<Request> batch) { for (Request req : batch) { if (req.getRetryCount() < MAX_RETRY) { req.incrementRetryCount(); try { requestQueue.put(req); } catch (InterruptedException ex) { Thread.currentThread().interrupt(); // 重试入队失败,持久化请求到可靠存储(如Redis、本地文件) persistRequest(req); } } else { // 超过最大重试次数,持久化请求 persistRequest(req); } } } // 发送单个请求兜底 private void sendSingleRequest(Request request) { try { // backendService.sendSingle(request); System.out.println("Sent single request due to submit failure"); } catch (Exception e) { if (request.getRetryCount() < MAX_RETRY) { request.incrementRetryCount(); try { requestQueue.put(request); } catch (InterruptedException ex) { Thread.currentThread().interrupt(); persistRequest(request); } } else { persistRequest(request); } } } // 持久化请求到可靠存储 private void persistRequest(Request req) { // 实现持久化逻辑,比如写入Redis、本地文件或消息队列 System.out.println("Persisted request after max retries"); } // 优雅关闭资源,处理剩余请求 public void shutdown() { scheduler.shutdown(); workerPool.shutdown(); // 处理队列中剩余的所有请求 processBatch(); try { if (!scheduler.awaitTermination(1, TimeUnit.SECONDS)) { scheduler.shutdownNow(); } if (!workerPool.awaitTermination(1, TimeUnit.SECONDS)) { workerPool.shutdownNow(); } } catch (InterruptedException e) { scheduler.shutdownNow(); workerPool.shutdownNow(); Thread.currentThread().interrupt(); } } }
关键保障说明
- 线程安全:
LinkedBlockingQueue的put和drainTo都是原子操作,多线程提交和处理不会出现竞态条件。 - 无重复发送:
drainTo一次性取出批量请求,同一请求不会被多个线程重复处理。 - 无请求丢失:提交时阻塞入队,中断时兜底发送;发送失败时重试或持久化,确保请求最终被处理。
注意事项
- 队列容量:设置合理的队列上限,避免无界队列导致内存溢出。
- 线程池配置:根据后端服务的并发能力调整
workerPool的大小,避免批量请求压垮后端。 - 重试限制:必须设置最大重试次数,避免失败请求无限循环占用资源。
- 优雅关闭:程序退出前调用
shutdown方法,处理队列中剩余的请求。
内容的提问来源于stack exchange,提问作者Han Eui-Jun
相关产品推荐
相关产品推荐

