Java队列处理提速案例怎么做

wen java案例 28

本文目录导读:

Java队列处理提速案例怎么做

  1. 基础场景:List作为队列的慢速问题
  2. 多线程处理提速
  3. 批量处理优化
  4. 实际案例:日志处理系统
  5. 使用Disruptor(高性能无锁队列)
  6. 性能对比与建议

我来详细讲解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 超高吞吐,低延迟 最高 最低
  1. 使用合适的数据结构:ArrayDeque > LinkedList > ArrayList
  2. 异步处理:使用BlockingQueue解耦
  3. 批量操作:减少IO次数
  4. 合理配置线程数:CPU密集型 = CPU核心数,IO密集型 = CPU核心数 * 2
  5. 监控队列大小:防止内存溢出
  6. 考虑背压机制:控制生产者速度

选择哪种方案取决于你的具体场景和性能需求。

抱歉,评论功能暂时关闭!