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

Java中以线程运行指定方法:多Reader连接与读取方案咨询

Managing Reader Connections & Data Collection with Thread Safety

Hey there! Let's break down your problem and work through a cleaner, more maintainable approach for managing your Reader instances alongside your Client logic.

First: What's Off with Your Current Implementation?

Your current setup has a few pain points that make it feel "messy":

  • Tight coupling between Reader and Client: Using doneSignal and readSignal directly ties the Reader's data flow to the Client's signaling logic. This makes both classes harder to test and modify independently.
  • State management risks: The isConnected and isRunning flags are shared across threads without proper synchronization, which can lead to race conditions (e.g., a thread checking isConnected right as another thread updates it).
  • Manual thread management: Your startReaders() method runs only one reader at a time, and handling thread lifecycle manually is error-prone.

A Better Approach: Use Executor Framework + Blocking Queues

Let's refactor this to separate concerns, use Java's ExecutorService for thread management, and use BlockingQueue to safely pass data between Readers and the Client. Here's how to do it:

Step 1: Refactor the Reader Class

Make Reader focus solely on maintaining its connection and pushing data to a queue. This removes the need for cross-class signaling:

import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;

public class Reader implements Runnable {
    private final int id;
    private final int port;
    private volatile boolean isRunning = true;
    private volatile boolean isConnected = false;
    private final BlockingQueue<String> dataQueue = new LinkedBlockingQueue<>();

    public Reader(int id, int port) {
        this.id = id;
        this.port = port;
    }

    @Override
    public void run() {
        while (isRunning) {
            if (!isConnected) {
                attemptReconnection();
            } else {
                try {
                    // Read data and add to queue (blocks if queue is full, adjust capacity if needed)
                    String readings = readData();
                    if (readings != null && !readings.isEmpty()) {
                        dataQueue.put(readings);
                    }
                    // Add a small delay to avoid spamming the reader
                    Thread.sleep(100);
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                    isRunning = false;
                } catch (Exception e) {
                    // If reading fails, mark as disconnected to trigger reconnection
                    isConnected = false;
                    System.err.printf("Reader %d failed to read data: %s%n", id, e.getMessage());
                }
            }
        }
        // Cleanup connection when stopping
        disconnect();
    }

    private void attemptReconnection() {
        try {
            System.out.printf("Reader %d attempting to connect on port %d...%n", id, port);
            connect(); // Your existing connect logic here
            isConnected = true;
            System.out.printf("Reader %d connected successfully!%n", id);
        } catch (Exception e) {
            System.err.printf("Reader %d failed to connect: %s%n", id, e.getMessage());
            // Wait before retrying to avoid flooding
            try {
                Thread.sleep(5000);
            } catch (InterruptedException ie) {
                Thread.currentThread().interrupt();
                isRunning = false;
            }
        }
    }

    // Your existing connect logic
    private void connect() throws Exception {
        // Implement your socket/device connection here
    }

    // Your existing read logic
    private String readData() throws Exception {
        // Implement your data reading here, return null if no data available
        return "Sample data from reader " + id;
    }

    private void disconnect() {
        // Implement cleanup for your connection here
        isConnected = false;
        System.out.printf("Reader %d disconnected.%n", id);
    }

    public void stop() {
        isRunning = false;
    }

    public BlockingQueue<String> getDataQueue() {
        return dataQueue;
    }

    public int getId() {
        return id;
    }
}

Step 2: Update the Client Class

Use ExecutorService to manage all Reader threads, handle server instructions, and collect data from Readers to send to the server:

import java.io.*;
import java.net.Socket;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class Client {
    private final Reader[] readers;
    private final String serverAddress;
    private ExecutorService readerExecutor;

    public Client(int numReaders, String serverAddress) {
        this.serverAddress = serverAddress;
        this.readers = new Reader[numReaders];
        // Initialize readers (you can adjust this to match your port assignment logic)
        for (int i = 0; i < numReaders; i++) {
            readers[i] = new Reader(i, 10000 + i); // Example port assignment
        }
    }

    public void run() {
        try {
            // Wait for server signal to start
            Map<Integer, Integer> readerPortMap = readStartSignal();
            assignPorts(readerPortMap);

            // Start all reader threads using ExecutorService
            startReaders();

            // Start data collection thread to send data to server
            startDataCollectionToServer();

            // Keep client running until server sends terminate
            waitForTerminationSignal();
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            // Cleanup: stop all readers and shutdown executor
            stopReaders();
        }
    }

    private Map<Integer, Integer> readStartSignal() throws IOException {
        // Implement logic to receive start signal and port assignments from server
        // Example mock return:
        Map<Integer, Integer> portMap = new HashMap<>();
        for (int i = 0; i < readers.length; i++) {
            portMap.put(i, 10000 + i);
        }
        System.out.println("Received start signal from server with port assignments.");
        return portMap;
    }

    private void assignPorts(Map<Integer, Integer> readerPortMap) {
        // Update each reader's port based on server's assignment
        for (Map.Entry<Integer, Integer> entry : readerPortMap.entrySet()) {
            int readerId = entry.getKey();
            int port = entry.getValue();
            // Note: If you need to modify the reader's port, add a setter in Reader class
            System.out.printf("Assigned port %d to Reader %d%n", port, readerId);
        }
    }

    private void startReaders() {
        // Create a fixed thread pool matching the number of readers
        readerExecutor = Executors.newFixedThreadPool(readers.length);
        for (Reader reader : readers) {
            readerExecutor.submit(reader);
        }
        System.out.println("All reader threads started.");
    }

    private void startDataCollectionToServer() {
        // Start a separate thread to collect data from all readers and send to server
        new Thread(() -> {
            try (Socket socket = new Socket(serverAddress, 9000);
                 ObjectOutputStream out = new ObjectOutputStream(socket.getOutputStream());
                 BufferedReader in = new BufferedReader(new InputStreamReader(socket.getInputStream()))) {

                String serverCommand;
                while (!"terminate".equals(serverCommand = in.readLine())) {
                    // Collect data from all connected readers
                    for (Reader reader : readers) {
                        // Poll queue to avoid blocking if no data is available
                        String data = reader.getDataQueue().poll();
                        if (data != null) {
                            out.writeObject(new String[]{String.valueOf(reader.getId()), data});
                            out.flush();
                        }
                    }
                }
                System.out.println("Received terminate command from server.");
            } catch (IOException e) {
                e.printStackTrace();
            }
        }).start();
    }

    private void waitForTerminationSignal() throws IOException {
        // If you need to wait for a direct termination signal, implement here
        // Alternatively, the startDataCollectionToServer thread handles termination
        Thread.currentThread().join();
    }

    private void stopReaders() {
        if (readerExecutor != null) {
            for (Reader reader : readers) {
                reader.stop();
            }
            readerExecutor.shutdown();
            System.out.println("Reader executor shutdown complete.");
        }
    }
}

Key Improvements in This Approach

  • Decoupled components: Readers push data to their own queues, and the Client pulls from those queues. No cross-class signaling needed.
  • Thread safety: BlockingQueue handles thread-safe data transfer, and volatile flags ensure state changes are visible across threads.
  • Simplified thread management: ExecutorService handles thread creation, lifecycle, and pooling, so you don't have to manage threads manually.
  • Automatic reconnection: Each Reader handles its own connection status and retries independently, without blocking other Readers.

Final Notes

  • Adjust the BlockingQueue capacity if you expect high data volume (use LinkedBlockingQueue(int capacity) to prevent unbounded growth).
  • Add proper error handling for socket connections in both Reader and Client (e.g., handling connection timeouts).
  • Consider using ScheduledExecutorService if you need more control over reconnection intervals.

内容的提问来源于stack exchange,提问作者A. Clark

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 04:24:17