队列仅保留最后一条消息问题咨询:多线程文件读写场景
解决队列仅保留最后一条消息的问题
我帮你梳理下问题的核心原因,然后给出具体的修复方案:
问题根源分析
从你给出的Main类代码片段来看,主要有两个关键问题:
- Spring Bean作用域错误:如果你的
Reader(生产者)和Writer(消费者)是Spring默认的单例Bean,那么每次循环调用context.getBean()拿到的都是同一个实例。当你给它们设置新的BlockingQueue时,后面的队列会直接覆盖前面的,导致前面文件的生产者线程其实在往被遗弃的队列里写消息,只有最后一个队列的消息能被消费者处理。 - 队列作用域不当:
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
相关产品推荐
相关产品推荐

