1. Cách phối hợp nhiều luồng khi xử lý tệp
Khi bạn làm việc với các tệp lớn hoặc xử lý dữ liệu phức tạp, thường xuất hiện bài toán: chia nhỏ công việc giữa nhiều luồng. Ví dụ, một luồng đọc các dòng từ tệp, còn các luồng khác – xử lý những dòng đó (tìm từ, đếm thống kê, v.v.).
Vấn đề là gì?
- Nếu tất cả luồng cùng làm việc với một tài nguyên (ví dụ, tệp hoặc bộ sưu tập dùng chung), rất dễ phát sinh lỗi: một số luồng có thể “vượt trước” các luồng khác, xảy ra race condition, mất dữ liệu, tràn bộ nhớ.
- Nếu có nhiều luồng mà không có phối hợp – có luồng sẽ rỗi, có luồng sẽ quá tải.
Cần một giải pháp mà:
- Cho phép trao đổi dữ liệu giữa các luồng một cách an toàn.
- Không để tràn bộ nhớ (nếu xử lý chậm hơn đọc).
- Cho phép dễ dàng kết thúc tất cả luồng khi dữ liệu đã hết.
2. Mẫu “Producer–Consumer” (Nhà sản xuất – Người tiêu thụ)
Producer–Consumer là một mẫu kinh điển giúp các luồng làm việc nhịp nhàng mà không cản trở nhau.
Có hai vai trò: Producer (nhà sản xuất) tạo dữ liệu – ví dụ, đọc các dòng từ tệp hoặc nhận thông điệp từ mạng – và đặt chúng vào một hàng đợi chung. Consumer (người tiêu thụ) lấy dữ liệu từ hàng đợi và xử lý chúng: đếm từ, lưu vào cơ sở dữ liệu hoặc ghi sang tệp khác.
Ý tưởng chính là hai loại luồng này hoạt động độc lập. Producer có thể đọc nhanh hơn tốc độ consumer xử lý, hoặc ngược lại – và không ai bị tắc nghẽn. Ở giữa là bộ đệm – hàng đợi dùng để cân bằng nhịp độ công việc.
Sơ đồ trực quan
[Tệp] --(đọc)--> [Producer] --(đưa vào hàng đợi)--> [BlockingQueue] --(lấy)--> [Consumer] --(xử lý)
3. Triển khai với BlockingQueue
Trong Java, để tổ chức luồng trao đổi như vậy, giao diện BlockingQueue (ví dụ, triển khai ArrayBlockingQueue) là lựa chọn lý tưởng.
BlockingQueue là gì?
BlockingQueue là một hàng đợi an toàn cho nhiều luồng với kích thước giới hạn, tự đảm nhiệm đồng bộ hóa. Nếu producer cố thêm phần tử khi hàng đợi đã đầy, luồng sẽ tự động bị chặn và chờ cho tới khi có chỗ. Tương tự, nếu consumer cố lấy phần tử mà hàng đợi trống, nó sẽ chờ cho đến khi có thứ gì đó được đặt vào hàng đợi.
Cơ chế này tự động giải quyết vấn đề “vượt trước” giữa các luồng và tràn bộ nhớ: producer không nhồi thêm phần tử dư thừa vào hàng đợi, còn consumer không cố làm việc với rỗng. Tất cả luồng làm việc ổn định và nhịp nhàng.
Ví dụ tạo hàng đợi
import java.util.concurrent.*;
BlockingQueue<String> queue = new ArrayBlockingQueue<>(100); // bộ đệm 100 phần tử
100 – số phần tử tối đa trong hàng đợi. Nếu producer chạy nhanh hơn, nó sẽ chờ đến khi consumer xử lý bớt dữ liệu.
4. Phối hợp luồng: backpressure và kết thúc công việc
Giới hạn kích thước hàng đợi (backpressure)
Backpressure là cơ chế không cho producer “đổ” đầy bộ nhớ nếu consumer không kịp xử lý dữ liệu.
- Nếu hàng đợi đã đầy, producer tự động “hãm lại” (phương thức put() bị chặn).
- Nếu hàng đợi trống, consumer sẽ chờ (phương thức take() bị chặn).
Điều này cho phép hệ thống hoạt động ổn định ngay cả khi tốc độ giữa các luồng khác nhau.
Kết thúc: “poison pill” (viên thuốc độc)
Khi producer đọc xong tệp, cần thông báo cho các consumer rằng sẽ không còn dữ liệu nữa và đã đến lúc kết thúc.
Giải pháp:
- Đặt vào hàng đợi một đối tượng đặc biệt – “poison pill” (ví dụ, chuỗi "__END__" hoặc giá trị đặc biệt khác, nhưng không phải null).
- Consumer, khi nhận “poison pill”, sẽ hiểu rằng cần kết thúc.
Nếu có nhiều consumer, cần đặt số “poison pill” tương ứng với số consumer!
5. Ví dụ pipeline: đọc tệp và xử lý dòng
Hãy triển khai một pipeline đơn giản:
- Một luồng đọc các dòng từ tệp và đưa chúng vào hàng đợi.
- Nhiều luồng lấy các dòng từ hàng đợi, đếm số từ và in kết quả.
Bước 1. Producer – đọc tệp
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); // nếu hàng đợi đầy – chờ
}
} catch (IOException | InterruptedException e) {
e.printStackTrace();
} finally {
// Đặt poison pill cho mỗi consumer
try {
for (int i = 0; i < consumerCount; i++) {
queue.put(POISON_PILL);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}
}
Bước 2. Consumer – bộ xử lý dòng
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(); // nếu hàng đợi trống – chờ
if (POISON_PILL.equals(line)) {
break; // kết thúc
}
int wordCount = line.trim().isEmpty() ? 0 : line.trim().split("\\s+").length;
System.out.println(Thread.currentThread().getName() + ": " + wordCount + " từ trong dòng: " + line);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}
Bước 3. Khởi chạy pipeline qua 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"); // tệp của bạn
// Khởi chạy producer
Thread producer = new Thread(new FileProducer(queue, file, consumerCount));
producer.start();
// Khởi chạy các consumer qua ExecutorService
ExecutorService consumers = Executors.newFixedThreadPool(consumerCount);
for (int i = 0; i < consumerCount; i++) {
consumers.submit(new LineConsumer(queue));
}
// Chờ producer kết thúc
producer.join();
// Chờ các consumer kết thúc
consumers.shutdown();
consumers.awaitTermination(1, TimeUnit.MINUTES);
System.out.println("Xử lý hoàn tất!");
}
}
6. Trực quan hóa: sơ đồ hoạt động của pipeline
flowchart LR
A[Producer: đọc tệp] -- put() --> Q[BlockingQueue]
Q -- take() --> C1[Consumer 1: đếm từ]
Q -- take() --> C2[Consumer 2: đếm từ]
Q -- take() --> C3[Consumer 3: đếm từ]
A -.->|poison pill| Q
Q -.->|poison pill| C1
Q -.->|poison pill| C2
Q -.->|poison pill| C3
7. Những lưu ý hữu ích và best practices
- Giới hạn kích thước hàng đợi – giúp tránh tràn bộ nhớ.
- Sử dụng “poison pill” để kết thúc – nếu không các consumer có thể treo mãi.
- Không dùng null làm “poison pill” nếu trong hàng đợi có thể có các giá trị null thực sự. Tốt hơn là một chuỗi hoặc đối tượng đặc biệt.
- Xử lý InterruptedException – điều này quan trọng để kết thúc luồng một cách đúng đắn.
- Sử dụng ExecutorService để quản lý pool luồng – dễ và an toàn hơn so với tự tạo luồng bằng tay.
8. Những lỗi thường gặp khi triển khai pipeline
Lỗi số 1: Hàng đợi không giới hạn → OutOfMemoryError. Nếu dùng LinkedBlockingQueue mà không giới hạn kích thước, producer có thể “đổ” đầy bộ nhớ nếu consumer không kịp xử lý.
Lỗi số 2: Không đặt “poison pill” → các consumer bị treo. Nếu quên đặt “viên thuốc độc”, các luồng consumer sẽ chờ dữ liệu mới vô thời hạn.
Lỗi số 3: Chỉ đặt một “poison pill” trong khi có nhiều consumer. Mỗi luồng consumer cần một “viên thuốc” riêng – nếu không sẽ không phải tất cả đều kết thúc.
Lỗi số 4: Không xử lý InterruptedException. Nếu luồng bị ngắt, cần kết thúc công việc đúng cách (khôi phục cờ ngắt), nếu không có thể dẫn đến luồng “bị treo”.
Lỗi số 5: Dùng biến chung giữa các luồng mà không đồng bộ. Đừng “tự phát minh lại bánh xe” – hãy dùng BlockingQueue, không phải danh sách hay mảng thông thường.
GO TO FULL VERSION