并发编程需求:双文件并发逐行读取时匹配行号查看内容
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
lineCountvariable, so they have no way to know when the other has reached the same line number. - Duplicate
sc.nextLine()calls: You callsc.nextLine()once for printing and once forcounterSync, which skips every other line (you read line N for print, then immediately read line N+1 for the sync method). - Unrelated synchronization: The
counterSyncmethod 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:
- It checks if
currentLinestill equalsglobalLine(the value we read earlier) - If yes, it updates it to
globalLine + 1atomically - 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
currentLinevariable 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

