多线程读取java.util.Queue时每个线程获全部元素,有无标准/第三方实现?
Great question! Let's break this down clearly—because what you're asking for isn't a typical queue consumption pattern (where each item is processed once by one thread), but rather a broadcast/publish-subscribe pattern where every message gets delivered to every processing thread, regardless of their processing speed.
为什么普通java.util.Queue不行?
First off: standard concurrent queues like ConcurrentLinkedQueue or LinkedBlockingQueue are designed for competitive consumption. When a thread takes an item from the queue, it's removed—so only one thread ever processes that item. That's the opposite of what you need here.
核心思路
Your goal requires a system where every new message added to the "source" is copied/distributed to all listening processing threads. Think of it as a news feed: every subscriber gets every new post, not just one person.
实现方案
1. 基于Java标准库的自定义实现
If you want to stick to core Java without third-party dependencies, you can build a simple broadcast system using thread-safe collections and executors:
步骤:
- Create a central broadcaster that holds a list of registered handlers (your processing threads/logic).
- Use
CopyOnWriteArrayListto store handlers—it's thread-safe and ideal for read-heavy scenarios (since you'll add handlers rarely but broadcast often). - When a new message arrives, iterate over all handlers and submit the message processing task to an executor (so each handler runs in its own thread).
代码示例:
// 定义消息类 public class QueueMessage { private final String content; public QueueMessage(String content) { this.content = content; } public String getContent() { return content; } } // 消息广播器 import java.util.List; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; public class MessageBroadcaster { private final List<MessageHandler> handlers = new CopyOnWriteArrayList<>(); private final ExecutorService executor = Executors.newCachedThreadPool(); // 注册处理逻辑 public void registerHandler(MessageHandler handler) { handlers.add(handler); } // 广播消息给所有处理线程 public void broadcastMessage(QueueMessage message) { for (MessageHandler handler : handlers) { executor.submit(() -> handler.process(message)); } } // 处理逻辑接口 public interface MessageHandler { void process(QueueMessage message); } // 关闭资源 public void shutdown() { executor.shutdown(); } } // 使用示例 public class Main { public static void main(String[] args) throws InterruptedException { MessageBroadcaster broadcaster = new MessageBroadcaster(); // 注册慢处理线程 broadcaster.registerHandler(message -> { try { Thread.sleep(1000); // 模拟1秒的慢处理 System.out.printf("[Slow Thread] Processed: %s%n", message.getContent()); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }); // 注册快处理线程 broadcaster.registerHandler(message -> { System.out.printf("[Fast Thread] Processed: %s%n", message.getContent()); }); // 发布消息 broadcaster.broadcastMessage(new QueueMessage("First incoming message")); Thread.sleep(500); // 模拟消息间隔 broadcaster.broadcastMessage(new QueueMessage("Second incoming message")); // 等待所有任务完成后关闭 Thread.sleep(2000); broadcaster.shutdown(); } }
2. 第三方库的现成实现
If you want to avoid reinventing the wheel, there are mature libraries that handle this pattern out of the box:
Guava EventBus
Guava's EventBus (and AsyncEventBus for async processing) is built exactly for this publish-subscribe use case. Every registered subscriber gets every event posted to the bus.
import com.google.common.eventbus.AsyncEventBus; import com.google.common.eventbus.Subscribe; import java.util.concurrent.Executors; public class EventBusDemo { public static void main(String[] args) throws InterruptedException { // 异步EventBus,每个订阅者在独立线程处理 AsyncEventBus eventBus = new AsyncEventBus(Executors.newCachedThreadPool()); // 注册慢订阅者 eventBus.register(new Object() { @Subscribe public void handleSlow(QueueMessage message) { try { Thread.sleep(1000); System.out.printf("[Slow Subscriber] Got message: %s%n", message.getContent()); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }); // 注册快订阅者 eventBus.register(new Object() { @Subscribe public void handleFast(QueueMessage message) { System.out.printf("[Fast Subscriber] Got message: %s%n", message.getContent()); } }); // 发布消息 eventBus.post(new QueueMessage("Hello from EventBus!")); Thread.sleep(500); eventBus.post(new QueueMessage("Another message for everyone")); // 等待处理完成 Thread.sleep(2000); } }
RxJava
RxJava's Observable is another powerful option. You can create an observable stream of messages, and each observer (processing thread) will receive every item in the stream. It's great if you need more complex thread scheduling or transformation logic.
关键注意事项
- Thread Safety: Always use thread-safe collections for storing handlers/subscribers (like
CopyOnWriteArrayList). - Resource Management: Don't forget to shut down executors/buses to avoid leaking threads.
- Backpressure: If slow handlers can't keep up with message volume, you might need to add backpressure logic (e.g., buffering messages per handler) to prevent out-of-memory issues. RxJava handles this well with built-in operators.
内容的提问来源于stack exchange,提问作者Maruthi

