Java Internals & Concurrency/Fork/Join: Chia để trị song song với work-stealing
32/75
Bài 32 / 75~14 phútConcurrency cơ bảnMiễn phí lượt xem

Fork/Join: Chia để trị song song với work-stealing

Chia để trị song song với work-stealing: ForkJoinPool, RecursiveTask/RecursiveAction, ngưỡng sequential cutoff, và liên hệ với parallel streams.

TL;DR: Fork/Join là framework cho bài toán divide-and-conquer CPU-bound: task tự chẻ thành task con đệ quy rồi gộp kết quả. Sức mạnh nằm ở work-stealing — mỗi worker giữ một deque riêng, owner push/pop ở đầu (LIFO), thread đói việc trộm từ đuôi (task to nhất) — nên tải tự cân ra các core, không có hàng đợi trung tâm làm cổ chai. Viết task qua RecursiveTask/RecursiveAction với sequential cutoff: quá nhỏ thì chi phí điều phối nuốt hết lợi ích, quá lớn thì không phủ kín core. Parallel stream chạy trên chính commonPool này; tuyệt đối không block I/O trong compute() — buộc phải chặn thì dùng ManagedBlocker hoặc pool riêng.

1. Bài toán chia để trị, và vì sao thread pool thường không vừa

Lấy một việc cụ thể: tính tổng một mảng long rất lớn. Tuần tự chỉ là một vòng lặp, nhưng nếu mảng có hàng chục triệu phần tử và máy tám core, để bảy core ngồi không là lãng phí. Ý tưởng chia để trị rất tự nhiên: tổng cả mảng bằng tổng nửa trái cộng nửa phải, mỗi nửa chẻ tiếp tới khi đủ nhỏ để tính thẳng rồi gộp ngược lên - hai nửa độc lập hoàn toàn nên chạy song song mà không cần khóa nào.

Chạy song song bằng gì? Phản xạ đầu tiên là ném mỗi mảnh vào một ExecutorService rồi get - nhưng đó là một cái bẫy (mục 5): task cha sau khi giao hai con phải chờ, nếu cha con cùng một fixed pool thì cha chiếm hết thread ngồi chờ, con không còn thread chạy - pool deadlock do cạn thread. Fork/Join giải đúng nút thắt này: nó biết task có dạng cây, biết cha sẽ chờ con, và tổ chức lại cách thread tiêu thụ công việc sao cho "cha chờ con" không bao giờ làm chết pool.

2. ForkJoinPool và work-stealing

Một thread pool thông thường như Executors.newFixedThreadPool có một hàng đợi trung tâm: mọi thread đói việc đều lấy task từ đó — tốt cho task độc lập, kích thước tương đương. Nhưng với cây divide-and-conquer, hàng triệu task con bé tí chen vào một hàng đợi chung, nó thành điểm tranh chấp nóng.

ForkJoinPool làm khác: mỗi worker có một hàng đợi riêng — work-stealing deque, hàng đợi hai đầu. Khi một task fork ra con, con được đẩy vào đầu local của deque thread đó, và khi cần việc tiếp thread cũng lấy từ đầu local — gần như không tranh chấp. Khi làm hết việc, thay vì ngồi không nó trở thành kẻ trộm: nhìn sang deque thread khác và lấy một task từ đuôi, đầu đối diện.

Hình dung một bếp ăn tám đầu bếp, mỗi người một chồng phiếu order úp trước mặt: ai lo chồng nấy, lấy phiếu trên cùng mà làm; ai hết chồng thì rút một phiếu từ đáy chồng người đang bận nhất - phiếu đáy thường là order lớn chưa ai đụng. Chủ rút từ trên, kẻ giúp việc rút từ dưới, hiếm khi với cùng một tờ; tải tự cân.

Sơ đồ dưới chụp lại đúng khoảnh khắc đó — Worker 1 đang bận với deque đầy task, Worker 2 vừa hết việc và đi trộm:

Work-stealing deque trong Fork/JoinWorker 1 lấy và đẩy task ở đầu deque theo LIFO. Worker 2 khi hết việc sẽ trộm task ở đuôi deque của Worker 1 — task to nhất, fork sớm nhất — nên hai bên hiếm khi tranh nhau.Worker 1 · ownerWorker 2 · thiefĐẦULIFOĐUÔI(rỗng — hết việc)① nhỏ② vừa③ TO
owner · ĐẦUthief · ĐUÔI
Worker 1 đang bận với deque đầy 3 task (đầu = task nhỏ vừa fork, đuôi = task to fork sớm nhất). Worker 2 vừa hết việc. Bấm “Chạy”/“Bước”.

Hai chiều mũi tên là linh hồn thiết kế. Owner làm LIFO ở đầu deque vì task con vừa fork còn nóng trong cache, lấy làm ngay là rẻ nhất; thief trộm FIFO ở đuôi vì task cũ nhất là nhánh to nhất của cây, một lần trộm ôm cả cây con. Hai đầu đối diện nên hiếm khi giành cùng task, và không có hàng đợi trung tâm làm cổ chai — Fork/Join chịu được cây đệ quy hàng triệu node mà không nghẹt.

2.1 commonPool

Ta không nhất thiết phải tự tạo ForkJoinPool: JVM duy trì sẵn một pool dùng chung cho toàn ứng dụng, gọi là common pool, lấy qua ForkJoinPool.commonPool(). Đây cũng là pool mặc định của parallel stream và của CompletableFuture.supplyAsync khi không truyền executor.

Kích thước mặc định là số core trừ một (cộng chính thread đang gọi, nên hiệu quả bằng số core): với CPU-bound, nhiều thread hơn số core không nhanh hơn, chỉ thêm chi phí context switching. Trên JDK 25 vẫn chỉnh được qua property java.util.concurrent.ForkJoinPool.common.parallelism, nhưng mặc định thường đúng. Cái tiện của common pool cũng là cái bẫy (mục 5): vì dùng chung cho cả ứng dụng, một task cư xử xấu - đặc biệt task blocking - có thể bỏ đói mọi thứ khác đang dựa vào nó.

3. RecursiveTaskRecursiveAction

Để mô tả việc, ta kế thừa RecursiveTask<V> (task trả về kết quả) hoặc RecursiveAction (không trả gì). Cả hai bắt ta cài compute(), nơi chứa logic chẻ-và-gộp: mảnh đủ nhỏ thì làm trực tiếp, còn to thì chẻ đôi thành hai con rồi điều phối chúng chạy song song.

Viết bài toán tính tổng mảng long ở mục 1 thành RecursiveTask:

public final class SumTask extends RecursiveTask<Long> {

    private static final int THRESHOLD = 10_000;   // nguong sequential cutoff
    private final long[] data;
    private final int lo, hi;

    SumTask(long[] data, int lo, int hi) { this.data = data; this.lo = lo; this.hi = hi; }

    @Override
    protected Long compute() {
        int length = hi - lo;
        if (length <= THRESHOLD) {              // du nho: lam thang
            long sum = 0;
            for (int i = lo; i < hi; i++) sum += data[i];
            return sum;
        }
        int mid = lo + length / 2;              // con to: che doi
        SumTask left  = new SumTask(data, lo, mid);
        SumTask right = new SumTask(data, mid, hi);
        left.fork();                            // giao nua trai cho pool chay song song
        long rightSum = right.compute();        // tu minh lam nua phai
        long leftSum  = left.join();            // cho va lay ket qua nua trai
        return leftSum + rightSum;
    }

    public static long sum(long[] data) {
        return ForkJoinPool.commonPool().invoke(new SumTask(data, 0, data.length));
    }
}

Cây chia việc đoạn code dựng lên trông như sau — mảng 40 nghìn phần tử, THRESHOLD = 10_000: node trên ngưỡng chẻ đôi, lá dưới ngưỡng tính thẳng bằng vòng for:

Cây chia việc fork/join ba tầng với bốn lá tính thẳng dưới ngưỡng

Kết quả chảy ngược từ lá lên gốc qua các lần join: mỗi node cộng hai con, gốc trả tổng cả mảng. Bốn lá độc lập nên bốn core cùng tính một lúc.

Thử đoán trước khi đọc tiếp

Thứ tự trong compute()fork nửa trái → compute nửa phải → join nửa trái. Nếu ta đảo lại thành fork nửa trái → join nửa trái ngay → rồi mới compute nửa phải, chương trình vẫn cho đúng kết quả. Vậy nó hỏng ở đâu? Viết ra câu trả lời trước khi đọc đoạn dưới.

Câu trả lời: fork đẩy nửa trái vào deque cho thread khác trộm, thread hiện tại tự chạy nửa phải bằng compute(); tới join, nếu nửa trái đã bị trộm thì chỉ lấy kết quả, nếu chưa thì tự thực thi nó ngay. Đảo thành "join ngay sau fork" thì join chặn tới khi nửa trái xong mới đụng nửa phải — hai con tuần tự hóa, mất sạch song song mà kết quả vẫn đúng nên khó phát hiện. Quy tắc: luôn join ngược thứ tự fork.

Vì quy tắc này dễ trượt tay, JDK cho sẵn idiom gọn hơn — ForkJoinTask.invokeAll tự lo đúng vũ đạo (fork task sau, chạy task đầu ngay trên thread hiện tại, chờ tất cả hoàn tất):

// Trong compute(), thay cho cap fork/compute/join thu cong (base case giu nguyen):
ForkJoinTask.invokeAll(left, right);    // fork right, compute left, wait for both
return left.join() + right.join();      // both done here -- join returns immediately

Sau khi invokeAll trả về, cả hai con đã xong nên hai join phía sau chỉ lấy kết quả, không chặn. Phiên bản này không thể đảo nhầm thứ tự và đọc rõ ý đồ hơn — nên làm lựa chọn mặc định.

compute() không cần khóa vì mỗi SumTask chỉ đọc đoạn [lo, hi) của riêng nó, không ghi vào data — task con không chia sẻ mutable state là điều kiện tiên quyết để Fork/Join phát huy. Nếu các nhánh tranh nhau ghi một biến đếm chung, ta rơi vào shared mutable state của bài Thread Safety, song song bị nuốt bởi contention.

4. Sequential cutoff nên đặt bao nhiêu?

Con số THRESHOLD = 10_000sequential cutoff - ranh giới mà dưới đó ta thôi chẻ và làm thẳng. Đặt sai, framework vẫn chạy đúng nhưng có thể chậm hơn cả vòng lặp tuần tự.

Mỗi lần chẻ tốn tạo hai object task, đẩy vào deque, có thể bị trộm, rồi join. Ngưỡng quá nhỏ - chẻ tới vài chục phần tử mỗi mảnh - tạo hàng triệu task để cộng vài con số mỗi cái; chi phí quản lý lấn át công việc thực, chậm hơn một vòng for. Ngưỡng quá lớn thì không đủ mảnh phủ core: mảng mười triệu với ngưỡng năm triệu chỉ tạo hai mảnh cho tám core, sáu core còn lại đứng nhìn. Điểm cân bằng là tạo nhiều mảnh hơn số core một chút (để còn task cho core rảnh trộm) nhưng mỗi mảnh vẫn đủ to so với chi phí fork/join - không có con số vàng, phải đo: đặt tổng số task vài lần tới vài chục lần số core rồi benchmark, luôn so với vòng for tuần tự (trên mảng vài nghìn phần tử nó thường thắng cả Fork/Join lẫn parallel stream).

Parallel stream chạy trên chính pool này

Khi viết stream.parallel() hay list.parallelStream(), cái máy bên dưới chính là ForkJoinPool.commonPool(): stream tự chẻ nguồn qua Spliterator, gói thành task Fork/Join, gộp kết quả - nên bài tính tổng gọn còn một dòng Arrays.stream(data).parallel().sum(). Đây gần như luôn là cách nên thử trước; tự viết RecursiveTask chỉ đáng khi cấu trúc không khớp mô hình stream. Và vì chia chung common pool, nhét blocking vào một map cũng ghim thread dùng chung — đúng cạm bẫy mục tiếp theo.

5. Cạm bẫy: blocking trong Fork/Join và chuyện chia sẻ common pool

Cạm bẫy lớn nhất lặp lại cảnh báo từ bài CompletableFuture pipeline nhưng nghiêm trọng hơn. ForkJoinPool kích thước theo số core với giả định ngầm: mỗi worker gần như lúc nào cũng đang tính, không ngồi chờ. Nhét một thao tác blocking - đọc file, gọi HTTP, chờ khóa, sleep - vào compute() phá vỡ giả định đó: thread chờ I/O vẫn chiếm một slot, vài task blocking là đủ bỏ đói cả pool dù CPU rảnh.

Khi bắt buộc phải chặn, framework cho lối thoát có kiểm soát: ForkJoinPool.ManagedBlocker. Gói thao tác blocking vào một ManagedBlocker (cài block() chạy lời gọi chặn, isReleasable() báo đã xong) rồi chờ qua ForkJoinPool.managedBlock(...); pool nhận biết một thread sắp chặn và tạm bù thêm worker để giữ mức song song. Đây là van an toàn cho một thao tác chặn bất khả kháng, không phải giấy phép chạy I/O ồ ạt — contract ở javadoc ManagedBlocker.

Cạm bẫy thứ hai: mọi parallel stream, CompletableFuture.supplyAsync mặc định và commonPool().invoke đổ chung một pool, nên một tác vụ tham lam ở góc này âm thầm làm chậm góc khác. Cách phòng là cô lập: workload đáng kể - nhất là có nguy cơ chặn hoặc chạy lâu - nên tạo một ForkJoinPool riêng để invoke/submit:

try (ForkJoinPool pool = new ForkJoinPool(4)) {   // pool rieng, co lap khoi common pool
    long total = pool.invoke(new SumTask(data, 0, data.length));
    // ... dung total
}                                                 // try-with-resources tu shutdown pool

Từ Java 19, ForkJoinPool cài AutoCloseable nên đặt được trong try-with-resources cho tự đóng. Một pool riêng tách bạch sự cố: nếu nó bị bỏ đói hay chạy chậm, thiệt hại khoanh trong workload đó, không lan ra phần còn lại.

6. Tie-in capstone: đối soát báo cáo bán vé theo divide-and-conquer

Fork/Join xuất hiện trong TicketFlow v3 đúng chỗ: một tác vụ tổng hợp CPU-bound ngoài đường đi của request. Cuối ngày đối soát - quét mảng rất lớn bản ghi giao dịch trong bộ nhớ, cộng doanh thu và đếm vé theo sự kiện: divide-and-conquer thuần khiết (không I/O, phép gộp có tính kết hợp), cấu trúc giống hệt SumTask với mỗi lá gộp một map doanh thu-theo-sự-kiện.

Ngược lại, nếu đối soát phải đọc từng dòng từ database thì không còn CPU-bound - truy vấn blocking trong compute() vấp đúng cạm bẫy mục 5. Ranh giới: I/O thuộc về Executor/CompletableFuture, chỉ tính toán thuần trong bộ nhớ mới giao cho Fork/Join trên pool riêng.

7. 📚 Deep Dive Oracle

📚 Deep Dive Oracle

Spec / reference chính thức:

  • Doug Lea — A Java Fork/Join Framework (2000) — paper gốc của tác giả framework; mô tả thiết kế work-stealing deque và các phép đo hiệu năng đầu tiên, ngắn và rất dễ đọc.
  • ForkJoinPool javadoc (Java 21) — quy định kích thước mặc định của commonPool, property chỉnh parallelism, và contract của ManagedBlocker.
  • ForkJoinTask javadoc — ghi rõ ngữ nghĩa fork/join/invokeAll và cảnh báo chính thức về task blocking.

Ghi chú: Fork/Join vào JDK 7 qua JSR 166y; cùng dòng công trình của Doug Lea sau này thành nền cho parallel stream (JDK 8). Đọc paper trước, javadoc sau — paper cho bức tranh, javadoc cho contract.

8. Tóm tắt

  • Fork/Join dành cho divide-and-conquer trên dữ liệu CPU-bound trong bộ nhớ. Sức mạnh là work-stealing: mỗi worker một deque riêng nên gần như không tranh chấp, tự trộm việc từ đuôi deque kẻ khác khi đói — tải tự cân, không cần hàng đợi trung tâm.
  • RecursiveTask/RecursiveAction mô tả việc qua compute(): nhỏ thì làm thẳng, to thì chẻ; fork một con, tự compute con kia, rồi join ngược thứ tự fork — hoặc dùng invokeAll cho chắc. Thứ tự này quyết định hiệu năng.
  • Sequential cutoff sống còn: quá nhỏ thì chi phí fork/join nuốt hết lợi ích, quá lớn thì không đủ mảnh phủ core. Tạo nhiều mảnh hơn số core một chút, luôn benchmark so với tuần tự.
  • Parallel stream chạy trên commonPool nên là cách nên thử trước; tự viết RecursiveTask chỉ đáng khi cấu trúc không khớp mô hình stream.
  • Đừng chặn trong Fork/Join — sizing theo số core giả định mọi thread đều đang tính, vài task blocking là đủ bỏ đói cả pool. Buộc phải chặn thì dùng ManagedBlocker; workload nặng thì tạo pool riêng.

Cả Phần B — Executor, Future/CompletableFuture, Fork/Join — đứng trên một giả định: mỗi Thread của Java là tài nguyên đắt gắn với một OS thread suốt vòng đời, nên phải gói vào pool và tránh để nó ngồi chờ. Phần C lật ngược giả định ấy — virtual thread, bài tiếp theo.

9. Tự kiểm tra

Tự kiểm tra
0/5 câu đã trả lời
  1. Q1
    Vì sao trong work-stealing deque, owner lấy việc từ đầu (LIFO) còn thief trộm từ đuôi (FIFO)?
  2. Q2
    Sequential cutoff đặt quá nhỏ thì chuyện gì xảy ra? Quá lớn thì sao?
  3. Q3
    Vì sao không được block I/O bên trong compute()? ManagedBlocker giúp gì khi buộc phải chặn?
  4. Q4
    So sánh ba cách viết: (a) fork cả hai con rồi join cả hai; (b) fork một con, compute con kia, rồi join; (c) fork một con, join nó ngay, rồi compute con kia. Cách nào đúng, cách nào sai, vì sao?
  5. Q5
    Vì sao ném task chia-để-trị vào một fixed thread pool thường gây deadlock cạn thread, còn ForkJoinPool thì không?

Bài tiếp theo: Virtual Threads: Thread-per-request trở lại

Bài này đáng gửi cho bạn học cùng?

Copy link đã gắn nguồn — dán group, chat, hoặc LinkedIn.

Bài này có giúp bạn hiểu bản chất không?

Hỏi đáp về bài này

Chưa có câu hỏi

Đặt câu hỏi

Có gì chưa rõ trong bài? Đặt câu hỏi đầu tiên — câu trả lời từ cộng đồng giúp bạn (và người sau).

Đặt câu hỏi đầu tiên

Bài tiếp theo

Virtual Threads: vì sao rẻ hơn platform thread nhiều bậc