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

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 poll method 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 BlockingQueue to hold outgoing messages (thread-safe for producers to add messages at any time)
  • Use a single-threaded ExecutorService to run the consumer logic—this guarantees strict sequential processing
  • Use CompletableFuture to 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: ConcurrentHashMap maps 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 03:54:18