多队列场景下消费者线程的唤醒方案探讨
Great question! Let’s dive into some robust, low-latency-friendly solutions using Java’s built-in concurrent primitives that avoid the pitfalls you outlined:
1. Semaphore + Non-Blocking Queues
This is a straightforward, battle-tested approach:
- Producers: Every time you add a task to any of your queues, immediately call
semaphore.release()on a sharedSemaphore(initialized with 0 permits). - Consumer: In a loop, first call
semaphore.acquire()—this blocks until at least one task exists. Then iterate through all your queues, usingpoll()to drain tasks (no blocking here, since we know at least one queue has work). - Why this works: The semaphore guarantees you wake up the consumer the second a task is added, eliminating timeout-related delays. Draining all queues each pass ensures you don’t miss tasks added to other queues while processing one. For efficiency, you could track non-empty queues with an
AtomicReferenceArray, but even a full loop is trivial sinceConcurrentLinkedQueue.poll()is O(1).
2. Pre-Allocated Token Wakeup Queue
If you want to avoid even the tiny overhead of a semaphore’s permit counting, this approach uses a dedicated signal queue with pre-made tokens:
- Setup: Create one fixed, reusable object per task queue (e.g.,
private static final Object QUEUE_A_TOKEN = new Object();). Use aLinkedTransferQueuefor the wakeup signal queue. - Producers: After adding a task to queue X, offer
QUEUE_X_TOKENto the wakeup queue (duplicates are harmless, though you could use aConcurrentHashMapto track pending tokens if you want to optimize). - Consumer: Take a token from the wakeup queue (blocks until a signal arrives), then drain all tasks from the corresponding queue. After draining, check if the queue still has tasks (in case more were added mid-drain) and re-offer the token if needed.
- Why this works: No new objects are created (tokens are initialized once), and each queue’s producer only signals its own token—this splits contention across multiple task queues instead of funneling everything into one. The wakeup queue ensures instant notification when any queue has work.
3. AtomicBoolean + LockSupport
For an ultra-lightweight option, pair an atomic flag with LockSupport:
- Shared State: Use an
AtomicBoolean hasWork = new AtomicBoolean(false);visible to all producers and the consumer thread. - Producers: After adding a task, if
hasWork.compareAndSet(false, true)returnstrue, callLockSupport.unpark(consumerThread)—this only wakes the consumer if it was idle. - Consumer: Loop through these steps:
- Drain all tasks from every queue.
- If all queues are empty, set
hasWork.set(false), then callLockSupport.park()(blocks until unparked). - Before rechecking queues, set
hasWork.set(true)to cover edge cases where a task was added right as you started parking.
- Why this works: It’s about as lightweight as it gets—no queue overhead, no extra objects. The atomic flag prevents spurious wakeups and ensures you only wake the consumer when necessary.
Bonus: Phaser for Dynamic Queue Counts
If your number of task queues might change at runtime, a Phaser is a flexible choice:
- Setup: Call
phaser.register()every time you add a new queue to the system. - Producers: After adding a task, call
phaser.arrive()to signal work is available. - Consumer: Call
phaser.awaitAdvance(phaser.getPhase())to block until a producer signals. When removing a queue, callphaser.arriveAndDeregister()to clean up. - Why this works: It handles dynamic queue counts gracefully, and unlike wait/notify, there’s no risk of waiting through a full timeout when work is available.
All these approaches skip the issues you wanted to avoid: no wait/notify timeout gaps, no custom CountDownLatch implementations, and no unnecessary object creation or single-queue bottlenecks.
内容的提问来源于stack exchange,提问作者Ruslan Akhundov

