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

并发编程需求:双文件并发逐行读取时匹配行号查看内容

解决双文件逐行并发读取的行号对齐问题

Hey there! Let's break down your problem step by step. First, let's look at the issues in your current code that are keeping you from achieving your goal, then we'll build a solution that aligns the line counts between your two threads—and I'll explain how patterns like Compare-and-Swap (CAS) and avoiding Slipped Condition fit in along the way.

现有代码的核心问题

Your current setup has a few key flaws that prevent line count alignment:

  • No shared line state: Each thread maintains its own lineCount variable, so they have no way to know when the other has reached the same line number.
  • Duplicate sc.nextLine() calls: You call sc.nextLine() once for printing and once for counterSync, which skips every other line (you read line N for print, then immediately read line N+1 for the sync method).
  • Unrelated synchronization: The counterSync method only synchronizes print output—it doesn't coordinate line number progress between threads.

基础解决方案:基于锁的行协调器

Let's start with a lock-based implementation (easiest for beginners to grasp) that uses a shared coordinator to sync line progress. This avoids Slipped Condition by keeping condition checks and actions atomic.

Step 1: Create a Shared Coordinator Class

This class will track line progress and wait for both threads to reach the same line before allowing you to inspect the content:

static class LineCoordinator {
    private int currentLine = 0;
    private final Map<String, String> lineData = new HashMap<>();
    private final Object lock = new Object();

    // Called by each thread to submit their line and wait for a match
    public void waitForMatch(String threadId, int lineNum, String lineContent) throws InterruptedException {
        synchronized (lock) {
            // Wait until this thread's line number catches up to the global current line
            while (lineNum > currentLine) {
                lock.wait();
            }

            // Store this thread's line content
            lineData.put(threadId, lineContent);

            // If both threads have submitted the current line, process it
            if (lineData.size() == 2) {
                // This is where you inspect the matching lines!
                System.out.println("\n=== Matching Line " + currentLine + " ===");
                System.out.println("Source Thread: " + lineData.get("source"));
                System.out.println("Target Thread: " + lineData.get("target"));

                // Advance to the next line, clear data, and wake waiting threads
                currentLine++;
                lineData.clear();
                lock.notifyAll();
            } else {
                // Wait for the other thread to submit its line
                lock.wait();
            }
        }
    }
}

Step 2: Update Thread and File Reading Logic

Modify your thread and file reading code to use the coordinator:

public class ConcurrencyTest {
    public static void main(String[] args) throws IOException {
        String filePath1 = "path1.txt";
        String filePath2 = "path2.txt";
        LineCoordinator coordinator = new LineCoordinator();

        // Give each thread a unique ID for the coordinator to distinguish them
        MyThread source = new MyThread(coordinator, filePath1, "source");
        MyThread target = new MyThread(coordinator, filePath2, "target");

        source.start();
        target.start();
    }

    static class MyThread extends Thread {
        private final LineCoordinator coordinator;
        private final String filePath;
        private final String threadId;

        MyThread(LineCoordinator coordinator, String filePath, String threadId) {
            this.coordinator = coordinator;
            this.filePath = filePath;
            this.threadId = threadId;
        }

        @Override
        public void run() {
            // Use try-with-resources to auto-close streams
            try (FileInputStream inputStream = new FileInputStream(filePath);
                 Scanner sc = new Scanner(inputStream, "UTF-8")) {
                int lineCount = 0;
                while (sc.hasNextLine()) {
                    lineCount++;
                    String line = sc.nextLine();
                    // Optional: Print individual thread progress for debugging
                    System.out.println(threadId + " - Line " + lineCount + ": " + line);
                    // Wait for the other thread to reach the same line
                    coordinator.waitForMatch(threadId, lineCount, line);
                }

                // Handle end-of-file: Notify the other thread we're done
                synchronized (coordinator.lock) {
                    lineData.put(threadId, "[END OF FILE]");
                    coordinator.lock.notifyAll();
                }
            } catch (IOException | InterruptedException e) {
                e.printStackTrace();
            }
        }
    }
}

Avoiding Slipped Condition Here

Slipped Condition happens when a thread checks a condition, then another thread modifies the state before the first thread acts on the condition. By wrapping all condition checks (lineNum > currentLine, lineData.size() == 2) and subsequent actions (wait, store data, advance line) inside the synchronized block, we ensure these operations are atomic—no other thread can modify the state while we're working with it.

进阶:无锁CAS实现

If you want to use Compare-and-Swap (CAS) instead of explicit locks (good for high-concurrency scenarios), you can leverage Java's AtomicInteger for atomic line number updates. CAS works by atomically checking if a value is still what you expect, then updating it if so.

CAS-Based Coordinator Example

static class CASLineCoordinator {
    private final AtomicInteger currentLine = new AtomicInteger(0);
    private final ConcurrentHashMap<String, String> lineData = new ConcurrentHashMap<>();

    public void waitForMatch(String threadId, int lineNum, String lineContent) throws InterruptedException {
        while (true) {
            int globalLine = currentLine.get();

            // Wait until this thread's line catches up to the global line
            if (lineNum > globalLine) {
                Thread.sleep(10); // Simple spin wait; use LockSupport.park() for better efficiency
                continue;
            }

            // Store the current line content
            lineData.put(threadId, lineContent);

            // Check if both threads have submitted their lines
            if (lineData.size() == 2) {
                // Use CAS to atomically advance the global line (only if it hasn't changed)
                if (currentLine.compareAndSet(globalLine, globalLine + 1)) {
                    // Inspect the matching lines
                    System.out.println("\n=== Matching Line " + globalLine + " ===");
                    System.out.println("Source Thread: " + lineData.get("source"));
                    System.out.println("Target Thread: " + lineData.get("target"));

                    lineData.clear();
                    break;
                } else {
                    // CAS failed (another thread advanced the line), reset and retry
                    lineData.remove(threadId);
                }
            } else {
                // Wait for the other thread, then check if the line advanced
                Thread.sleep(10);
                if (currentLine.get() > globalLine) {
                    lineData.remove(threadId);
                }
            }
        }
    }
}

How CAS Helps

The currentLine.compareAndSet(globalLine, globalLine + 1) call is the CAS operation:

  1. It checks if currentLine still equals globalLine (the value we read earlier)
  2. If yes, it updates it to globalLine + 1 atomically
  3. If no (another thread already advanced the line), it returns false, and we retry

This avoids the need for explicit locks, using hardware-level atomic operations to ensure thread safety.

关键总结

  • Shared state is mandatory: You need a single source of truth for line progress (like the currentLine variable in the coordinator) to align threads.
  • Lock-based sync is beginner-friendly: It eliminates Slipped Condition by keeping condition checks and actions atomic.
  • CAS is for high concurrency: It's more efficient in busy systems but requires handling retry logic, making it trickier for new developers.

内容的提问来源于stack exchange,提问作者Timothy Cumberland

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.06 22:58:17