CodeGym /課程 /JAVA 25 SELF /檔案處理管線:Producer–Consumer

檔案處理管線:Producer–Consumer

JAVA 25 SELF
等級 59 , 課堂 4
開放

1. 在處理檔案時如何協調多個執行緒

當你處理大型檔案或進行複雜的資料處理時,常會遇到這樣的需求:將工作分配給多個執行緒。例如,一個執行緒從檔案讀取每一行,其他執行緒則處理這些行(搜尋單字、計算統計等)。

問題在哪裡?

  • 如果所有執行緒都操作同一個資源(例如檔案或共享集合),很容易出錯:有些執行緒可能「超前」其他執行緒,導致競態、資料遺失、記憶體過載。
  • 如果執行緒很多而缺乏協調——有些會閒置,有些會過載。

我們需要一個方案,它:

  • 允許執行緒之間安全地交換資料。
  • 避免記憶體被塞爆(當處理速度比讀取慢時)。
  • 在資料耗盡時,能輕鬆地讓所有執行緒結束。

2. 「Producer–Consumer」樣式(生產者–消費者)

Producer–Consumer 是一個經典的設計樣式,可讓多個執行緒協同工作、互不干擾。

它包含兩個角色:Producer(生產者)產生資料——例如從檔案讀取每一行或從網路接收訊息——並把資料放入共享佇列。Consumer(消費者)從佇列取出資料並處理:計算單字數、儲存到資料庫或寫到另一個檔案。

核心概念是這兩類執行緒彼此獨立運作。生產者可能比消費者處理得更快,反之亦然——但彼此不會相互阻塞。它們之間有一個緩衝區——用來平衡節奏的佇列。

視覺化示意

[檔案] --(讀取)--> [Producer] --(放入佇列)--> [BlockingQueue] --(取出)--> [Consumer] --(處理)

3. 以 BlockingQueue 實作

在 Java 中,要實作這樣的交換,最適合使用介面 BlockingQueue(例如其實作 ArrayBlockingQueue)。

什麼是 BlockingQueue?

BlockingQueue 是具有限制大小、且具備執行緒安全的佇列,會自行處理同步。如果生產者嘗試加入元素,而佇列已滿,該執行緒會自動被阻塞並等待直到有空位。同樣地,如果消費者嘗試取出元素而佇列為空,它就會等待直到有人把資料放入佇列。

這種機制自動解決了經典的「執行緒爭用」與記憶體溢位問題:生產者不會把佇列塞爆,消費者也不會對「空無一物」下手。所有執行緒都能平穩協調地完成各自的工作。

建立佇列的範例

import java.util.concurrent.*;

BlockingQueue<String> queue = new ArrayBlockingQueue<>(100); // 有 100 個元素的緩衝區

100 是佇列可容納的最大元素數量。如果 producer 跑得更快,它會等待,直到 consumer 處理掉一部分資料。

4. 執行緒協調:backpressure 與終止

限制佇列大小(backpressure)

Backpressure 是一種機制:當 consumer 來不及處理時,它可避免 producer 把記憶體「灌爆」。

  • 如果佇列已滿,producer 自動「減速」(put() 會被阻塞)。
  • 如果佇列為空,consumer 會等待(take() 會被阻塞)。

這讓系統即使在執行緒速度不同的情況下也能穩定運作。

結束流程:「poison pill」

當 producer 讀完檔案後,需要告訴 consumer 們資料已經沒有了,該收工了。

解法:

  • 往佇列放入一個特殊物件——「poison pill」(例如字串 "__END__" 或其他特殊值,但不要用 null)。
  • Consumer 收到「poison pill」後,就知道是時候結束了。

如果有多個 consumer,必須放入與 consumer 數量相同的「poison pill」!

5. 管線範例:讀取檔案並處理每行字串

我們來實作一個簡單的管線:

  • 一個執行緒從檔案讀取每一行並放入佇列。
  • 多個執行緒從佇列取出每一行,計算單字數並輸出結果。

步驟 1. Producer —— 檔案讀取器

import java.io.*;
import java.util.concurrent.*;

public class FileProducer implements Runnable {
    private final BlockingQueue<String> queue;
    private final File file;
    private final int consumerCount;
    private final String POISON_PILL = "__END__";

    public FileProducer(BlockingQueue<String> queue, File file, int consumerCount) {
        this.queue = queue;
        this.file = file;
        this.consumerCount = consumerCount;
    }

    @Override
    public void run() {
        try (BufferedReader reader = new BufferedReader(new FileReader(file))) {
            String line;
            while ((line = reader.readLine()) != null) {
                queue.put(line); // 若佇列已滿——就等待
            }
        } catch (IOException | InterruptedException e) {
            e.printStackTrace();
        } finally {
            // 為每個 consumer 放入一顆 poison pill
            try {
                for (int i = 0; i < consumerCount; i++) {
                    queue.put(POISON_PILL);
                }
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        }
    }
}

步驟 2. Consumer —— 行處理器

import java.util.concurrent.BlockingQueue;

public class LineConsumer implements Runnable {
    private final BlockingQueue<String> queue;
    private final String POISON_PILL = "__END__";

    public LineConsumer(BlockingQueue<String> queue) {
        this.queue = queue;
    }

    @Override
    public void run() {
        try {
            while (true) {
                String line = queue.take(); // 若佇列為空——就等待
                if (POISON_PILL.equals(line)) {
                    break; // 結束工作
                }
                int wordCount = line.trim().isEmpty() ? 0 : line.trim().split("\\s+").length;
                System.out.println(Thread.currentThread().getName() + ": " + wordCount + " 個單字於該行: " + line);
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

步驟 3. 透過 Executors 啟動管線

import java.io.File;
import java.util.concurrent.*;

public class PipelineDemo {
    public static void main(String[] args) throws Exception {
        int consumerCount = 3;
        BlockingQueue<String> queue = new ArrayBlockingQueue<>(100);

        File file = new File("input.txt"); // 你的檔案

        // 啟動 producer
        Thread producer = new Thread(new FileProducer(queue, file, consumerCount));
        producer.start();

        // 透過 ExecutorService 啟動 consumers
        ExecutorService consumers = Executors.newFixedThreadPool(consumerCount);
        for (int i = 0; i < consumerCount; i++) {
            consumers.submit(new LineConsumer(queue));
        }

        // 等待 producer 結束
        producer.join();

        // 等待 consumers 結束
        consumers.shutdown();
        consumers.awaitTermination(1, TimeUnit.MINUTES);

        System.out.println("處理完成!");
    }
}

6. 視覺化:管線運作示意

flowchart LR
    A[Producer: 讀取檔案] -- put() --> Q[BlockingQueue]
    Q -- take() --> C1[Consumer 1: 計算單字數]
    Q -- take() --> C2[Consumer 2: 計算單字數]
    Q -- take() --> C3[Consumer 3: 計算單字數]
    A -.->|poison pill| Q
    Q -.->|poison pill| C1
    Q -.->|poison pill| C2
    Q -.->|poison pill| C3

7. 實用細節與最佳實務

  • 限制佇列大小——可避免記憶體溢位。
  • 使用「poison pill」來結束——否則 consumer 可能會永遠卡住。
  • 不要使用 null 當作「poison pill」,如果佇列中可能出現真實的 null 值,請改用特殊字串或物件。
  • 處理 InterruptedException——這對於正確地終止執行緒很重要。
  • 使用 ExecutorService 管理執行緒池——比手動建立執行緒更簡單且更安全。

8. 實作管線時的常見錯誤

錯誤 №1: 無上限佇列 → OutOfMemoryError。 如果使用未限制大小的 LinkedBlockingQueue,當 consumer 來不及處理時,producer 可能會把記憶體「灌爆」。

錯誤 №2: 沒有放入「poison pill」 → consumer 卡住。 如果忘了放入「毒藥丸」,消費者執行緒會一直等待新資料。

錯誤 №3: 只放了一顆「poison pill」,但有多個 consumer。 每個消費者執行緒都需要自己的「藥丸」——否則不會全部結束。

錯誤 №4: 沒有處理 InterruptedException 當執行緒被中斷時,必須正確地結束工作(恢復中斷旗標),否則可能導致執行緒「卡住」。

錯誤 №5: 未經同步就跨執行緒共用變數。 不要嘗試「重新發明輪子」——請使用 BlockingQueue,而不是一般的清單或陣列。

1
問卷/小測驗
平行處理檔案,等級 59,課堂 4
未開放
平行處理檔案
平行處理檔案
留言
TO VIEW ALL COMMENTS OR TO MAKE A COMMENT,
GO TO FULL VERSION