You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

队列仅保留最后一条消息问题咨询:多线程文件读写场景

解决队列仅保留最后一条消息的问题

我帮你梳理下问题的核心原因,然后给出具体的修复方案:

问题根源分析

从你给出的Main类代码片段来看,主要有两个关键问题:

  1. Spring Bean作用域错误:如果你的Reader(生产者)和Writer(消费者)是Spring默认的单例Bean,那么每次循环调用context.getBean()拿到的都是同一个实例。当你给它们设置新的BlockingQueue时,后面的队列会直接覆盖前面的,导致前面文件的生产者线程其实在往被遗弃的队列里写消息,只有最后一个队列的消息能被消费者处理。
  2. 队列作用域不当:BlockingQueue在for循环内部创建,每次循环都会生成新队列,但如果没有和对应的生产者/消费者实例绑定好,就会出现消息丢失或只保留最后一组消息的情况。

具体修复步骤

1. 把Reader和Writer改成原型作用域

因为每个文件都需要独立的生产者和消费者实例,每个实例绑定专属的队列。你可以通过注解或XML配置把它们的作用域设为prototype,这样每次获取Bean都会得到新的实例。

示例(注解方式):

@Component
@Scope("prototype") // 关键:每次getBean都会创建新实例
public class Reader implements Runnable {
    private BlockingQueue<String> queue;
    private String filePath;

    // 提供setter设置队列和文件路径
    public void setQueue(BlockingQueue<String> queue) {
        this.queue = queue;
    }

    public void setFilePath(String filePath) {
        this.filePath = filePath;
    }

    @Override
    public void run() {
        try (BufferedReader br = new BufferedReader(new FileReader(filePath))) {
            String line;
            while ((line = br.readLine()) != null) {
                queue.put(line); // 逐行写入队列
            }
            queue.put("EOF"); // 写入结束标记,告诉消费者任务完成
        } catch (IOException | InterruptedException e) {
            Thread.currentThread().interrupt();
            e.printStackTrace();
        }
    }
}

对应的Writer类:

@Component
@Scope("prototype")
public class Writer implements Runnable {
    private BlockingQueue<String> queue;
    private int waitTime;

    public void setQueue(BlockingQueue<String> queue) {
        this.queue = queue;
    }

    public void setWaitTime(int waitTime) {
        this.waitTime = waitTime;
    }

    @Override
    public void run() {
        try {
            String line;
            // 直到收到结束标记才停止消费
            while (!(line = queue.take()).equals("EOF")) {
                System.out.println(line);
                Thread.sleep(waitTime);
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            e.printStackTrace();
        }
    }
}

2. 修正Main类的循环逻辑

确保每个文件对应的队列、生产者、消费者一一绑定,并且线程列表放在循环外部,方便后续等待所有线程完成:

public class Main {
    public static void main(String[] args) {
        ApplicationContext context = new ClassPathXmlApplicationContext("applicationContext.xml");
        List<Thread> threadList = new ArrayList<>(); // 线程列表放在循环外

        for (int i = 0; i < args.length; i++) {
            String file = args[i];
            int queueSize = 10;
            int waitTime = 200;

            // 为当前文件创建专属队列
            BlockingQueue<String> queue = new LinkedBlockingQueue<>(queueSize);

            // 获取新的生产者和消费者实例
            Reader reader = context.getBean(Reader.class);
            Writer writer = context.getBean(Writer.class);

            // 绑定当前文件的队列和参数
            reader.setQueue(queue);
            reader.setFilePath(file);
            writer.setQueue(queue);
            writer.setWaitTime(waitTime);

            // 创建并启动线程
            Thread producerThread = new Thread(reader, "Producer-" + file);
            Thread consumerThread = new Thread(writer, "Consumer-" + file);
            threadList.add(producerThread);
            threadList.add(consumerThread);
            
            producerThread.start();
            consumerThread.start();
        }

        // 等待所有线程执行完毕
        for (Thread thread : threadList) {
            try {
                thread.join();
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                e.printStackTrace();
            }
        }

        // 关闭Spring上下文
        ((ClassPathXmlApplicationContext) context).close();
    }
}

修复后的效果

每个文件都会有独立的队列、生产者和消费者线程,不会出现队列被覆盖的情况,所有文件的内容都会被正确读取并消费,不会只保留最后一条消息。

内容的提问来源于stack exchange,提问作者Dred

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.26 09:41:09