Java Internals & Concurrency/Structured Concurrency: vòng đời có kỷ luật cho nhóm task
35/75
Bài 35 / 75~14 phútConcurrency cơ bảnMiễn phí lượt xem

Structured Concurrency: vòng đời có kỷ luật cho nhóm task

StructuredTaskScope (JEP 505, preview ở Java 25) ràng vòng đời nhóm subtask vào một khối try — fork tỏa nhánh, join hợp lưu, close chống rò — với các Joiner fail-fast.

TL;DR: ExecutorService + Future để hở hai lỗ: subtask có thể sống lâu hơn method đã sinh ra nó (task leak), và lỗi một nhánh không tự lan sang các nhánh anh em. StructuredTaskScope (JEP 505, preview thứ năm ở Java 25) ràng vòng đời cả nhóm subtask vào một khối try: fork tỏa nhánh, join hợp lưu, close bảo đảm không nhánh nào sống sót ra ngoài block. Mặc định open()fail-fast (Joiner.awaitAllSuccessfulOrThrow()): một nhánh ném exception thì scope cancel ngay, các nhánh còn lại bị interrupt, join() ném FailedException. Chọn Joiner khác đổi được chính sách hợp lưu — chờ tất cả, hay lấy kết quả thành công đầu tiên — và áp được một deadline chung cho cả nhóm ngay tại điểm join.

1. Vì sao concurrency không cấu trúc lại rò rỉ?

Một handler cần dữ liệu user và tồn kho để dựng response — hai nguồn độc lập, nên ta chạy song song bằng ExecutorService:

ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor();
Future<User>      userF  = executor.submit(() -> findUser(userId));
Future<Inventory> invF   = executor.submit(() -> checkInventory(eventId));
User user      = userF.get();   // (1)
Inventory inv  = invF.get();    // (2)
return new Response(user, inv);

Đoạn này gọn nhưng rò rỉ ở chỗ khó thấy. Nếu findUser ném exception ở dòng (1), ta thoát method — nhưng checkInventory vẫn chạy tiếp: không ai gọi invF.cancel(...), nên nó còn chiếm một virtual thread, có thể giữ một database connection, rồi kết quả bị vứt đi. Đó là task leak. Chiều ngược lại cũng tệ: checkInventory treo vì service dưới chết thì dòng (2) block mãi cùng cả thread xử lý request — caller không áp được deadline lên cả nhóm, chỉ thấy hai Future rời rạc.

Gốc rễ ở chữ "rời rạc". Mỗi submit ném một task vào pool như thả thư vào hòm thư công cộng: quan hệ cha–con giữa code đang gọi và task sinh ra không được mã hóa ở đâu cả, nên task con sống lâu hơn method sinh ra nó, lỗi không tự lan về cha, việc hủy phải làm thủ công từng cái. Code tuần tự thì trái ngược: hàm con luôn trả về trước hàm cha, exception lan theo call stack, block kết thúc thì mọi thứ nó mở ra đã đóng. Structured concurrency mang chính cấu trúc khối ấy trở lại cho code song song: ràng vòng đời subtask vào một block cú pháp, mở scope đầu block thì chắc chắn mọi subtask kết thúc trước khi ra khỏi block.

2. StructuredTaskScope: fan-out, rồi fan-in

Đến Java 25, StructuredTaskScope (trong java.util.concurrent) đang ở preview lần thứ năm (JEP 505): API ổn định hình dạng nhưng vẫn cần --enable-preview để biên dịch. Hình dạng dưới đây bám JDK 25; các bản preview trước có API khác đáng kể.

Ý tưởng cốt lõi gói trong ba động tác: mở scope, fork các subtask, rồi join chờ cả nhóm. Viết lại ví dụ ở phần 1:

Response handle(String userId, String eventId) throws InterruptedException {
    try (var scope = StructuredTaskScope.open()) {          // mo pham vi
        Subtask<User>      user = scope.fork(() -> findUser(userId));
        Subtask<Inventory> inv  = scope.fork(() -> checkInventory(eventId));
        scope.join();                                       // diem hop luu duy nhat
        return new Response(user.get(), inv.get());
    }
}

Khác biệt so với ExecutorService không nằm ở số dòng — gần như bằng nhau — mà ở những đảm bảo cấu trúc khối áp đặt. scope.fork(...) khởi chạy mỗi subtask trên một virtual thread riêng và trả về một Subtask — handle nhẹ cho kết quả tương lai. scope.join() chặn tới khi mọi subtask kết thúc, hoặc tới khi đủ điều kiện dừng sớm theo chính sách. Mấu chốt ở try-with-resources: khi luồng điều khiển rời khối try — dù return bình thường hay vì exception — scope.close() được gọi tự động và bảo đảm không subtask nào của scope còn sống, không có đường nào để một nhánh vượt ra ngoài block đã sinh ra nó.

Subtask.get() chỉ được gọi sau join() — khác cố ý với Future.get(). Future.get() tự nó là lệnh block, mời ta chờ từng task một theo thứ tự xen kẽ khó lường; Subtask.get() không block: tại lúc gọi nó, join() đã bảo đảm subtask kết thúc, nên nó chỉ đọc kết quả có sẵn. Việc chờ dồn về đúng một chỗ là join() — tỏa nhiều nhánh rồi gom tại một điểm hợp lưu duy nhất, đúng tinh thần "fan-out rồi fan-in".

Hãy hình dung scope như quản đốc giao việc cho tổ thợ: phát phiếu (fork), rồi đứng cửa xưởng đợi (join). Quy tắc bất di bất dịch: quản đốc không rời xưởng chừng nào còn một thợ chưa về; buộc phải đi (hết giờ, cháy xưởng) thì gọi tất cả về trước. Cái xưởng có tường có cửa ấy chính là khối try.

3. Khi nào scope dừng cả nhóm task?

Vậy scope.open() không tham số áp chính sách gì? Mặc định (theo JEP 505) là Joiner.awaitAllSuccessfulOrThrow()fail-fast, không phải "chờ tất cả rồi tính sau". Mọi nhánh thành công thì nó chờ đủ cả nhóm rồi cho join() trả về null (kiểu Void), đọc kết quả qua từng Subtask.get(). Nhưng chỉ cần một subtask ném exception, scope bị cancel ngay: mọi nhánh chưa xong bị interrupt, join() ném StructuredTaskScope.FailedException bọc nguyên nhân gốc, kích hoạt close() dọn sạch (đoạn code phần 2 đã mang sẵn hành vi này).

Vẽ chuỗi sự kiện đó thành sơ đồ — một subtask fail kéo cả cây xuống có trật tự:

Một subtask ném exception làm scope bị cancel, hai nhánh còn lại bị interrupt

Chính sách kết thúc do Joiner truyền vào open(...) quyết định. (Bản preview cũ dùng ShutdownOnFailure/ShutdownOnSuccess cho vai trò này — đã thay bằng các Joiner tương ứng.)

3.1 Ba Joiner cho kết quả "tất cả phải thành công"

Họ Joiner có ba thành viên hay gặp cho bài toán chờ nhiều nhánh, khác nhau ở hai câu hỏi: có fail-fast khôngjoin() trả về gì.

JoinerKhi một subtask failjoin() trả vềHợp với
awaitAllSuccessfulOrThrow() (default của open())Cancel cả scope ngay, join() ném FailedExceptionnull (kiểu Void) — đọc kết quả qua từng Subtask.get()Ít nhánh, mỗi nhánh một kiểu khác nhau, đã giữ sẵn Subtask ref
allSuccessfulOrThrow()Cancel cả scope ngay, join() ném FailedExceptionMột Stream các Subtask — duyệt kết quả như streamNhiều nhánh đồng kiểu, không muốn giữ từng ref
awaitAll()Không cancel — chờ mọi nhánh kết thúc bất kể thành bạinull — tự kiểm tra Subtask.state() từng cáiCần đủ kết quả lẫn lỗi của tất cả các nhánh (batch, báo cáo)

Ví dụ user + inventory ở phần 2 rơi vào ô thứ nhất. Khi các nhánh đồng kiểu, muốn gom kết quả như dòng chảy, allSuccessfulOrThrow() gọn hơn vì join() trả thẳng stream:

List<Quote> fetchQuotes(List<String> symbols) throws InterruptedException {
    try (var scope = StructuredTaskScope.open(Joiner.<Quote>allSuccessfulOrThrow())) {
        symbols.forEach(s -> scope.fork(() -> fetchQuote(s)));
        return scope.join()                  // Stream cac Subtask da thanh cong
                    .map(Subtask::get)
                    .toList();
    }
}

Còn awaitAll() (ô thứ ba) cho khi fail-fast không phải điều ta muốn — chạy mười phép kiểm tra rồi cần báo cáo đầy đủ cái nào đậu cái nào rớt: scope chờ hết mọi nhánh, ta tự duyệt Subtask.state() từng cái.

Việc hủy dựa trên interrupt của Java, nên chỉ cắt được subtask tôn trọng interrupt: phần lớn blocking I/O trên virtual thread đều tôn trọng, nhưng vòng lặp thuần CPU không kiểm tra Thread.interrupted() thì không bị cắt. Đây là cancellation hợp tác — cơ chế đã học ở Thread API và vòng đời, gặp lại ở Executor (hủy qua Future.cancel).

3.2 Lấy kết quả thành công đầu tiên

Mặt đối xứng của fail-fast là success-fast. Query cùng thông tin từ ba replica, chỉ cần một câu trả lời nhanh nhất, các cái còn lại đều dư. Joiner.anySuccessfulResultOrThrow() phục vụ đúng kiểu này:

String fetchFastest(List<String> replicas) throws InterruptedException {
    try (var scope = StructuredTaskScope.open(
             Joiner.<String>anySuccessfulResultOrThrow())) {
        for (String r : replicas) {
            scope.fork(() -> queryReplica(r));
        }
        return scope.join();                // tra ket qua thanh cong DAU TIEN
    }
}

Subtask thành công đầu tiên kích hoạt cancel scope; các nhánh còn lại bị interrupt. scope.join() trả thẳng giá trị đó nên không cần giữ Subtask ref; nếu mọi nhánh đều thất bại, join() ném exception gom các nguyên nhân. Đây là pattern hedged request kinh điển — đánh đổi ít tài nguyên dư để cắt đuôi độ trễ — gói trong vài dòng lifecycle đảm bảo, thay vì CompletableFuture.anyOf (bài CompletableFuture pipeline) cộng logic hủy thủ công.

3.3 Deadline cho cả nhóm

Vì cả nhóm subtask sống trong một scope, ta áp được deadline lên toàn bộ nhóm tại điểm join — thứ ba Future rời rạc không cho làm gọn. Timeout là tham số thứ hai của open(...): open(joiner, cf -> cf.withTimeout(Duration.ofMillis(500))) (phải viết joiner tường minh vì cần config phía sau). Quá hạn 500ms, scope bị cancel như khi một subtask fail: mọi nhánh chưa xong bị interrupt, join() ném StructuredTaskScope.TimeoutException (unchecked, lồng trong API mới — không phải java.util.concurrent.TimeoutException cũ). Deadline thành thuộc tính của cả phép tính concurrent, không phải gánh thủ công lên từng nhánh; pattern đầy đủ ở capstone của bài 22b.

4. So với ExecutorService: vì sao có kỷ luật hơn

Đặt hai mô hình cạnh nhau, khác biệt là hình dạng vòng đời. ExecutorService tách rời ba việc — submit, chờ kết quả, dọn dẹp — nên quan hệ cha–con chỉ tồn tại trong đầu lập trình viên; quên hủy hay quên chờ thì trình biên dịch không cản. StructuredTaskScope buộc cả ba về cùng một block try, khiến ba thuộc tính khó-đảm-bảo thành mặc định miễn phí: không task leak (close() luôn chờ hết), lỗi không bị nuốt (lan lên cha), cancellation không bị quên (tự động theo chính sách).

Điều này không có nghĩa ExecutorService bị khai tử. Một pool dài hạn nhận việc từ nhiều nguồn không liên quan (ví dụ background worker xử lý hàng đợi sự kiện) vẫn là việc của nó — ở đó vốn dĩ không có quan hệ cha–con để cấu trúc hóa. StructuredTaskScope dành cho tình huống ngược lại: một đơn vị công việc tỏa nhiều nhánh rồi gom lại, sống và chết trong phạm vi một lời gọi.

Còn một mảnh ghép: khi các nhánh fork ra cần chung ngữ cảnh (user đăng nhập, traceId), truyền tay qua từng chữ ký hàm là cực hình. ScopedValue giải đúng bài toán đó — và đó là nội dung bài kế tiếp, nơi ta ráp StructuredTaskScope với ScopedValue thành capstone khép lại module.

5. 📚 Deep Dive Oracle

📚 Deep Dive Oracle

Spec / reference chính thức:

Ghi chú: structured concurrency vẫn là preview — hình dạng API có thể còn chỉnh ở các JDK sau, nên khi đọc tài liệu trên mạng hãy kiểm tra nó viết cho bản preview nào; các lớp ShutdownOnFailure/ShutdownOnSuccess bạn gặp trong bài viết cũ thuộc API đời trước, đã bị thay bằng Joiner.

6. Liên hệ các bài khác

  • Thread API và vòng đời — interrupt và cooperative cancellation là cơ chế bên dưới việc scope hủy subtask; không hiểu interrupt thì không giải thích được vì sao vòng lặp CPU thuần "không chịu chết".
  • Executor và thread pool — đối trọng trực tiếp: hủy thủ công qua Future.cancel so với cancel tự động theo chính sách Joiner.
  • CompletableFuture pipelineanySuccessfulResultOrThrow() thay cho pattern CompletableFuture.anyOf + hủy tay; nhìn lại để thấy structured concurrency rút gọn được gì.
  • Virtual Threads — tiền đề vật chất của bài này: fork tạo một virtual thread cho mỗi subtask, chỉ khả thi vì thread đã rẻ.
  • ScopedValue ở quy mô virtual thread — bài kế tiếp: truyền ngữ cảnh bất biến xuống các nhánh fork, và capstone ráp cả hai cơ chế.

7. Tóm tắt

  • Structured concurrency mang cấu trúc khối của code tuần tự trở lại code song song: fork tỏa nhánh, join hợp lưu ở đúng một điểm, close (qua try-with-resources) bảo đảm không subtask nào sống sót ra ngoài block — chặn task leak tận gốc. Subtask.get() chỉ hợp lệ sau join() và không block.
  • Default của open() là fail-fast (awaitAllSuccessfulOrThrow()): một nhánh ném exception thì cancel cả scope, các nhánh còn lại bị interrupt, join() ném FailedException.
  • Đổi Joiner là đổi chính sách hợp lưu: allSuccessfulOrThrow() trả stream, awaitAll() không cancel (báo cáo đủ thành/bại), anySuccessfulResultOrThrow() lấy kết quả nhanh nhất (hedged request).
  • Vì cả nhóm sống trong một scope, áp được một deadline chung tại join — hủy dựa trên interrupt nên chỉ cắt được nhánh tôn trọng interrupt.

8. Tự kiểm tra

Tự kiểm tra
0/5 câu đã trả lời
  1. Q1
    Structured concurrency giải quyết loại leak nào mà ExecutorService + Future để hở?
  2. Q2
    Default joiner của StructuredTaskScope.open() làm gì khi một nhánh fail?
  3. Q3
    awaitAllSuccessfulOrThrow() và allSuccessfulOrThrow() khác nhau ở đâu, chọn cái nào khi nào?
  4. Q4
    Vì sao Subtask.get() không block còn Future.get() thì block?
  5. Q5
    Việc hủy subtask trong scope dựa trên cơ chế gì, và khi nào nó không cắt được?

Bài tiếp theo: ScopedValue ở quy mô virtual thread — truyền context không leak

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

ScopedValue trên virtual thread: truyền context không leak