基于java.util.concurrent包实现多线程读写消息队列的技术需求问询
基于java.util.concurrent包实现多线程读写消息队列的技术需求问询
看起来你正在啃Java并发编程的生产者-消费者模型,这个消息队列的练习题挺经典的,刚好可以用java.util.concurrent包里的工具类完美解决,完全不用自己手写wait/notify那套容易踩坑的同步逻辑~
先帮你理清楚核心需求的实现思路,再补全代码并拆解关键点:
核心实现方案选择
java.util.concurrent包中的LinkedBlockingQueue是这个需求的绝佳选择:它本身就是线程安全的阻塞队列,内置了阻塞式的put()和take()方法——当队列满时,写线程调用put()会自动阻塞;当队列空时,读线程调用take()也会自动阻塞,完全替我们处理了线程间的同步等待,完美避开了禁止使用wait/notify的限制。
完整实现代码
import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.atomic.AtomicInteger; public class MessageQueueApp { private static final int MESSAGE_COUNT = 5; public static void main(String[] args) { // 处理命令行参数,获取线程数量N if (args.length != 1) { System.err.println("请传入一个整数参数,指定写/读线程的数量N"); System.exit(1); } int threadCount; try { threadCount = Integer.parseInt(args[0]); if (threadCount <= 0) { System.err.println("线程数量必须是正整数哦"); System.exit(1); } } catch (NumberFormatException e) { System.err.println("参数必须是有效的整数,请检查输入"); System.exit(1); } // 初始化阻塞消息队列 BlockingQueue<String> messageQueue = new LinkedBlockingQueue<>(); int totalMessages = threadCount * MESSAGE_COUNT; // 创建并启动N个写线程,每个写线程生产5条消息 for (int i = 1; i <= threadCount; i++) { Thread writerThread = new Thread(new MessageWriter(i, messageQueue), "Writer-" + i); writerThread.start(); } // 创建并启动N个读线程,读取所有生产的消息 for (int i = 1; i <= threadCount; i++) { Thread readerThread = new Thread(new MessageReader(i, messageQueue, totalMessages), "Reader-" + i); readerThread.start(); } } // 写线程任务:生产指定数量的消息 static class MessageWriter implements Runnable { private final int threadId; private final BlockingQueue<String> messageQueue; public MessageWriter(int threadId, BlockingQueue<String> messageQueue) { this.threadId = threadId; this.messageQueue = messageQueue; } @Override public void run() { try { for (int i = 1; i <= MESSAGE_COUNT; i++) { String message = String.format("【生产消息】线程%d-第%d条", threadId, i); // 阻塞式放入队列,队列满时自动等待 messageQueue.put(message); System.out.printf("[%s] 已生产:%s%n", Thread.currentThread().getName(), message); // 模拟生产耗时,比如业务处理时间 Thread.sleep(100); } } catch (InterruptedException e) { // 响应中断,优雅停止线程 Thread.currentThread().interrupt(); System.out.printf("[%s] 被中断,已停止生产%n", Thread.currentThread().getName()); } } } // 读线程任务:消费消息直到所有消息被读取 static class MessageReader implements Runnable { private final int threadId; private final BlockingQueue<String> messageQueue; private final int totalMessagesToRead; // 原子类统计已读消息数,保证多线程下计数安全 private static final AtomicInteger totalReadCount = new AtomicInteger(0); public MessageReader(int threadId, BlockingQueue<String> messageQueue, int totalMessagesToRead) { this.threadId = threadId; this.messageQueue = messageQueue; this.totalMessagesToRead = totalMessagesToRead; } @Override public void run() { try { while (totalReadCount.get() < totalMessagesToRead) { // 阻塞式获取消息,队列空时自动等待 String message = messageQueue.take(); totalReadCount.incrementAndGet(); System.out.printf("[%s] 已消费:%s%n", Thread.currentThread().getName(), message); // 模拟消费耗时 Thread.sleep(150); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.out.printf("[%s] 被中断,已停止消费%n", Thread.currentThread().getName()); } } } }
关键细节拆解
- 线程命名规范:创建线程时直接指定名称为
Writer-1、Reader-2这类格式,清晰区分线程角色和编号,完全满足“线程名有意义”的要求。 - 阻塞队列的优势:
LinkedBlockingQueue的put()和take()方法是线程安全的,自动处理了线程间的等待/唤醒逻辑,不需要我们手动写wait/notify,完美符合需求限制。 - 并发安全的计数:用
AtomicInteger统计已读消息总数,它是java.util.concurrent包提供的原子类,保证多个读线程同时更新计数时不会出现线程安全问题。 - 优雅的中断处理:在捕获
InterruptedException后,调用Thread.currentThread().interrupt()恢复中断状态,这是并发编程中良好的实践,保证线程能正确响应外部中断信号。
内容来源于stack exchange
相关产品推荐
相关产品推荐

