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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.08 12:49:37