Java带消费超时的Producer and Consumer模式及CommunicationSupervisor实现咨询
Got it, let's break this down into two clear parts as you requested. I'll provide practical, tested Java implementations for both scenarios, with explanations to help you follow along.
1. Producer-Consumer Pattern with Consumer Timeout Mechanism
Core Approach
We'll leverage Java's java.util.concurrent.BlockingQueue (specifically LinkedBlockingQueue) for thread-safe message passing. The consumer will use the poll(long timeout, TimeUnit unit) method, which waits for an element for the specified duration and returns null if the timeout elapses—perfect for avoiding infinite blocking.
Implementation Code
import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; // Producer that generates messages with random delays class TimeoutProducer implements Runnable { private final BlockingQueue<String> messageQueue; private final int totalMessages; public TimeoutProducer(BlockingQueue<String> queue, int messageCount) { this.messageQueue = queue; this.totalMessages = messageCount; } @Override public void run() { try { for (int i = 1; i <= totalMessages; i++) { String msg = "Task-" + i; messageQueue.put(msg); System.out.println("[Producer] Generated: " + msg); Thread.sleep((long) (Math.random() * 1200)); // Simulate variable work time } // Send poison pill to signal consumer shutdown messageQueue.put("SHUTDOWN"); } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.out.println("[Producer] Interrupted mid-production"); } } } // Consumer with timeout logic for waiting on messages class TimeoutConsumer implements Runnable { private final BlockingQueue<String> messageQueue; private final long timeoutMs; public TimeoutConsumer(BlockingQueue<String> queue, long timeout) { this.messageQueue = queue; this.timeoutMs = timeout; } @Override public void run() { try { while (true) { // Wait for a message, or timeout after specified time String msg = messageQueue.poll(timeoutMs, TimeUnit.MILLISECONDS); if (msg == null) { System.out.println("[Consumer] Timed out waiting for new messages"); continue; // Or break here if you want the consumer to exit on timeout } if ("SHUTDOWN".equals(msg)) { System.out.println("[Consumer] Received shutdown signal, exiting"); break; } // Process the message System.out.println("[Consumer] Processing: " + msg); Thread.sleep(800); // Simulate message processing time } } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.out.println("[Consumer] Interrupted during processing"); } } } // Test the implementation public class TimeoutProducerConsumerDemo { public static void main(String[] args) { BlockingQueue<String> queue = new LinkedBlockingQueue<>(10); int messageCount = 6; long consumerTimeout = 2500; // 2.5 seconds Thread producer = new Thread(new TimeoutProducer(queue, messageCount)); Thread consumer = new Thread(new TimeoutConsumer(queue, consumerTimeout)); producer.start(); consumer.start(); try { producer.join(); consumer.join(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }
Key Notes
- The
pollmethod is critical here: it lets the consumer wait for messages without blocking forever. - We use a "poison pill" (
SHUTDOWN) to trigger a graceful consumer shutdown after all messages are processed. - All queue operations are inherently thread-safe, so no manual synchronization is needed between producers and consumers.
2. CommunicationSupervisor for Device Connection
Requirements Recap
- Outgoing messages stored in a queue
- Single consumer for sequential processing
- Consumer sends a message, waits for its response before moving to the next message
- New messages can be added while waiting for a response, but they’re queued and processed in order
Core Approach
- Use a
BlockingQueueto hold outgoing messages (thread-safe for producers to add messages at any time) - Use a single-threaded
ExecutorServiceto run the consumer logic—this guarantees strict sequential processing - Use
CompletableFutureto wait for responses, allowing us to unblock the consumer as soon as a response arrives
Implementation Code
import java.util.concurrent.*; // Represents an outgoing message to the device class DeviceOutgoingMessage { private final String messageId; private final String payload; public DeviceOutgoingMessage(String id, String content) { this.messageId = id; this.payload = content; } public String getMessageId() { return messageId; } public String getPayload() { return payload; } } // Represents an incoming response from the device class DeviceIncomingResponse { private final String correspondingMessageId; private final String responseContent; public DeviceIncomingResponse(String msgId, String content) { this.correspondingMessageId = msgId; this.responseContent = content; } public String getCorrespondingMessageId() { return correspondingMessageId; } public String getResponseContent() { return responseContent; } } class CommunicationSupervisor { private final BlockingQueue<DeviceOutgoingMessage> outgoingQueue = new LinkedBlockingQueue<>(); private final ExecutorService consumerService = Executors.newSingleThreadExecutor(); private final ConcurrentHashMap<String, CompletableFuture<DeviceIncomingResponse>> responseTracker = new ConcurrentHashMap<>(); public CommunicationSupervisor() { // Start the single consumer thread on initialization consumerService.submit(this::processOutgoingQueue); } // Method for producers to send messages (can be called from any thread) public void queueOutgoingMessage(DeviceOutgoingMessage message) { try { outgoingQueue.put(message); System.out.println("[Supervisor] Queued message: " + message.getMessageId()); } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.err.println("[Supervisor] Failed to queue message: " + message.getMessageId()); } } // Callback to handle incoming responses (triggered when device sends a response) public void handleIncomingResponse(DeviceIncomingResponse response) { CompletableFuture<DeviceIncomingResponse> pendingFuture = responseTracker.remove(response.getCorrespondingMessageId()); if (pendingFuture != null) { pendingFuture.complete(response); System.out.println("[Supervisor] Received response for: " + response.getCorrespondingMessageId()); } else { System.err.println("[Supervisor] Received response for unknown message: " + response.getCorrespondingMessageId()); } } // Core consumer logic: process messages in order, wait for responses private void processOutgoingQueue() { try { while (!Thread.currentThread().isInterrupted()) { // Block until a message is available in the queue DeviceOutgoingMessage msg = outgoingQueue.take(); System.out.println("[Supervisor] Processing message: " + msg.getMessageId()); // Create a future to track the pending response CompletableFuture<DeviceIncomingResponse> responseFuture = new CompletableFuture<>(); responseTracker.put(msg.getMessageId(), responseFuture); // Simulate sending the message to the device (replace with your actual protocol code) simulateMessageSend(msg); // Wait for the response (with a 15-second timeout) try { DeviceIncomingResponse response = responseFuture.get(15, TimeUnit.SECONDS); System.out.println("[Supervisor] Completed processing for " + msg.getMessageId() + ": " + response.getResponseContent()); } catch (TimeoutException e) { System.err.println("[Supervisor] Timeout waiting for response for: " + msg.getMessageId()); responseTracker.remove(msg.getMessageId()); } catch (ExecutionException e) { System.err.println("[Supervisor] Error processing response for: " + msg.getMessageId() + " - " + e.getCause()); responseTracker.remove(msg.getMessageId()); } } } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.out.println("[Supervisor] Consumer thread interrupted, shutting down"); } finally { consumerService.shutdown(); } } // Replace this with your actual device communication logic (TCP, MQTT, etc.) private void simulateMessageSend(DeviceOutgoingMessage message) { System.out.println("[Supervisor] Sent message to device: " + message.getMessageId()); } // Cleanup method to shut down the supervisor gracefully public void shutdown() { consumerService.shutdownNow(); } } // Test the CommunicationSupervisor public class CommunicationSupervisorDemo { public static void main(String[] args) throws InterruptedException { CommunicationSupervisor supervisor = new CommunicationSupervisor(); // Simulate sending multiple messages (some while waiting for responses) Thread sender = new Thread(() -> { supervisor.queueOutgoingMessage(new DeviceOutgoingMessage("DEV-REQ-001", "Check battery level")); try { Thread.sleep(2000); // Wait 2s before sending next message } catch (InterruptedException e) { Thread.currentThread().interrupt(); } supervisor.queueOutgoingMessage(new DeviceOutgoingMessage("DEV-REQ-002", "Activate sensor")); }); // Simulate device responses arriving asynchronously Thread responder = new Thread(() -> { try { // Wait 3s before sending first response Thread.sleep(3000); supervisor.handleIncomingResponse(new DeviceIncomingResponse("DEV-REQ-001", "Battery level: 85%")); // Wait 2s before sending second response Thread.sleep(2000); supervisor.handleIncomingResponse(new DeviceIncomingResponse("DEV-REQ-002", "Sensor activated")); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }); sender.start(); responder.start(); sender.join(); responder.join(); // Give time for final processing Thread.sleep(1000); supervisor.shutdown(); } }
Key Features Explained
- Single Consumer Guarantee:
Executors.newSingleThreadExecutor()ensures only one thread processes messages, so they’re handled in the exact order they’re queued. - Sequential Processing: The consumer blocks on
responseFuture.get()until the response arrives (or times out) before moving to the next message—even if new messages are added to the queue while waiting. - Thread-Safe Queue: Producers can call
queueOutgoingMessage()from any thread at any time, with no risk of race conditions. - Response Tracking:
ConcurrentHashMapmaps message IDs to their pending futures, allowing us to unblock the consumer as soon as a response is received. - Timeout Protection: The 15-second timeout ensures the consumer doesn’t block indefinitely if a device fails to respond.
内容的提问来源于stack exchange,提问作者Govan
相关产品推荐
相关产品推荐

