本文目录导读:

我来详细讲解Java队列处理提速的常见方案和案例。
基础场景:List作为队列的慢速问题
慢速实现
// 慢速:使用ArrayList作为队列
public class SlowQueueProcessor {
private List<String> queue = new ArrayList<>();
public void addTask(String task) {
queue.add(task); // O(1)
}
public String processTask() {
if (queue.isEmpty()) return null;
return queue.remove(0); // O(n) - 需要移动所有元素
}
}
快速实现
// 快速:使用LinkedList或ArrayDeque
public class FastQueueProcessor {
private Queue<String> queue = new LinkedList<>();
// 或者
// private Queue<String> queue = new ArrayDeque<>();
public void addTask(String task) {
queue.offer(task); // O(1)
}
public String processTask() {
return queue.poll(); // O(1)
}
}
多线程处理提速
基础版本(Serial处理)
public class SerialProcessor {
public void processTasks(List<Task> tasks) {
for (Task task : tasks) {
processTask(task); // 串行处理
}
}
private void processTask(Task task) {
// 耗时操作,如数据库查询、网络请求
task.execute();
}
}
多线程优化版本
public class ParallelQueueProcessor {
private BlockingQueue<Task> taskQueue = new LinkedBlockingQueue<>(1000);
private volatile boolean running = true;
private List<Thread> workers = new ArrayList<>();
public void startWorkers(int workerCount) {
for (int i = 0; i < workerCount; i++) {
Thread worker = new Thread(() -> {
while (running) {
try {
Task task = taskQueue.poll(100, TimeUnit.MILLISECONDS);
if (task != null) {
processTask(task);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
}, "worker-" + i);
worker.start();
workers.add(worker);
}
}
public void addTask(Task task) throws InterruptedException {
taskQueue.put(task); // 阻塞直到有空间
}
public void shutdown() {
running = false;
workers.forEach(Thread::interrupt);
}
private void processTask(Task task) {
task.execute();
}
}
批量处理优化
单个处理模式
public class SingleProcessor {
private Database database;
public void processTask(Task task) {
// 每次单个写入数据库
database.save(task);
// 建立数据库连接、执行SQL、关闭连接
// 频繁建立连接开销大
}
}
批量处理优化
public class BatchProcessor {
private BlockingQueue<Task> queue = new LinkedBlockingQueue<>();
private static final int BATCH_SIZE = 100;
private static final int FLUSH_INTERVAL = 1000; // 1秒
@PostConstruct
public void init() {
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
scheduler.scheduleAtFixedRate(this::flush,
FLUSH_INTERVAL, FLUSH_INTERVAL, TimeUnit.MILLISECONDS);
}
public void add(Task task) {
queue.offer(task);
// 达到批量大小立即处理
if (queue.size() >= BATCH_SIZE) {
flush();
}
}
private synchronized void flush() {
List<Task> batch = new ArrayList<>(BATCH_SIZE);
queue.drainTo(batch, BATCH_SIZE);
if (!batch.isEmpty()) {
// 批量写入数据库
database.batchSave(batch);
}
}
}
实际案例:日志处理系统
问题:高吞吐日志处理
// 存在问题:每个日志都直接写入文件
public class LogProcessor {
public void processLog(String log) {
// 每次都要打开文件、写入、关闭
try (FileWriter fw = new FileWriter("app.log", true);
BufferedWriter bw = new BufferedWriter(fw)) {
bw.write(log);
bw.newLine();
} catch (IOException e) {
e.printStackTrace();
}
}
}
优化方案:批量+异步写入
public class HighPerformanceLogProcessor {
private BlockingQueue<String> logQueue = new ArrayBlockingQueue<>(10000);
private volatile boolean running = true;
private FileWriter fileWriter;
private BufferedWriter bufferedWriter;
public HighPerformanceLogProcessor() throws IOException {
fileWriter = new FileWriter("app.log", true);
bufferedWriter = new BufferedWriter(fileWriter);
// 启动后台写线程
startWriterThread();
}
private void startWriterThread() {
Thread writer = new Thread(() -> {
List<String> batch = new ArrayList<>(100);
while (running) {
try {
// 等待第一条日志
String log = logQueue.poll(500, TimeUnit.MILLISECONDS);
if (log != null) {
batch.add(log);
// 收集更多日志
logQueue.drainTo(batch, 99); // 最多100条
// 批量写入
writeBatch(batch);
batch.clear();
} else {
// 超时,刷入缓冲区
bufferedWriter.flush();
}
} catch (Exception e) {
e.printStackTrace();
}
}
}, "log-writer");
writer.setDaemon(true);
writer.start();
}
private void writeBatch(List<String> logs) throws IOException {
for (String log : logs) {
bufferedWriter.write(log);
bufferedWriter.newLine();
}
bufferedWriter.flush();
}
public void addLog(String log) {
if (!logQueue.offer(log)) {
// 队列满了,可以丢弃或阻塞
// 这里选择异步丢弃
System.err.println("Log queue is full, dropping log: " + log);
}
}
public void shutdown() {
running = false;
try {
bufferedWriter.close();
} catch (IOException e) {
e.printStackTrace();
}
}
}
使用Disruptor(高性能无锁队列)
import com.lmax.disruptor.*;
import com.lmax.disruptor.dsl.Disruptor;
public class DisruptorProcessor {
// 定义事件
public static class TaskEvent {
private Task task;
public void setTask(Task task) {
this.task = task;
}
public Task getTask() {
return task;
}
}
// 事件工厂
public static class TaskEventFactory implements EventFactory<TaskEvent> {
public TaskEvent newInstance() {
return new TaskEvent();
}
}
// 事件处理器
public static class TaskEventHandler implements EventHandler<TaskEvent> {
public void onEvent(TaskEvent event, long sequence, boolean endOfBatch) {
processTask(event.getTask());
}
private void processTask(Task task) {
task.execute();
}
}
public void start() {
// RingBuffer的大小,必须是2的N次方
int bufferSize = 1024;
Disruptor<TaskEvent> disruptor = new Disruptor<>(
new TaskEventFactory(),
bufferSize,
DaemonThreadFactory.INSTANCE, // 线程工厂
ProducerType.SINGLE, // 单生产者
new BlockingWaitStrategy() // 等待策略
);
// 连接事件处理器
disruptor.handleEventsWith(new TaskEventHandler());
// 启动Disruptor
disruptor.start();
// 获取RingBuffer用于发布事件
RingBuffer<TaskEvent> ringBuffer = disruptor.getRingBuffer();
// 发布事件示例
for (int i = 0; i < 10000; i++) {
long sequence = ringBuffer.next();
try {
TaskEvent event = ringBuffer.get(sequence);
event.setTask(new Task(i));
} finally {
ringBuffer.publish(sequence);
}
}
disruptor.shutdown();
}
}
性能对比与建议
| 方案 | 适用场景 | 吞吐量 | 延迟 |
|---|---|---|---|
| 单线程处理 | 简单任务,数量少 | 低 | 低 |
| 多线程处理 | CPU密集型任务 | 中 | 中 |
| 批量处理 | IO密集型任务(DB、文件) | 高 | 高 |
| Disruptor | 超高吞吐,低延迟 | 最高 | 最低 |
- 使用合适的数据结构:ArrayDeque > LinkedList > ArrayList
- 异步处理:使用BlockingQueue解耦
- 批量操作:减少IO次数
- 合理配置线程数:CPU密集型 = CPU核心数,IO密集型 = CPU核心数 * 2
- 监控队列大小:防止内存溢出
- 考虑背压机制:控制生产者速度
选择哪种方案取决于你的具体场景和性能需求。