Concurrency · همزمانی سنیورSenior ~61 دقیقه مطالعه~50 min read
CompletableFuture، ForkJoin و موازیسازیCompletableFuture, ForkJoin & Parallelism
در این درس قدمبهقدم یاد میگیری CompletableFuture چطور کارها را مثل یک آشپزخانهی هماهنگ به هم میچسباند، ForkJoinPool چطور با «سرقت کار» همهی هستهها را مشغول نگه میدارد، چرا مسدود شدن روی common pool کل برنامه را زمین میزند، و virtual threadها در جاوا ۲۱ چطور بازی را عوض میکنند.A step-by-step lesson on how CompletableFuture wires async work together like a well-run kitchen, how ForkJoinPool keeps every core busy through "work-stealing," why blocking on the shared common pool can freeze your whole app, and how Java 21 virtual threads change the game.
پیشنیاز:Prerequisites: نخها، Runnable/Callable و ExecutorهاThreads, Runnable/Callable & Executors
خب، بیا رک باشیم: همروندی (concurrency) در جاوا جایی است که خیلی از مهندسهای باتجربه هم گیر میکنند — نه چون سخت است، بلکه چون دو ابزار کاملاً متفاوت را با هم قاطی میکنند و بعد تعجب میکنند چرا برنامهی تولیدیشان نصفهشب قفل میشود. این درس دقیقاً میخواهد این گره را باز کند. قرار است از صفر بسازیم، هر واژهی فنی را همان لحظه که میبینی معنایش کنیم، و در پایان تو باید بتوانی هر سؤال مصاحبهی این حوزه را با اعتماد جواب بدهی.
در این فصل این مسیر را طی میکنیم: (۱) یک مدل ذهنی میسازیم تا ForkJoinPool را از CompletableFuture تشخیص بدهی؛ (۲) CompletableFuture را از ساختن تا ترکیب و مدیریت خطا کامل یاد میگیری؛ (۳) قانون طلایی «کدام نخ چه چیزی را اجرا میکند» را میفهمی — همان چیزی که بیشترین باگ تولیدی از آن میآید؛ (۴) ForkJoinPool و مکانیزم سرقت کار را میشکافیم؛ (۵) درون parallelStream() را میبینیم؛ (۶) با reactive مقایسه میکنیم؛ و (۷) میبینیم virtual threadها در جاوا ۲۱ چطور کل معادله را عوض میکنند. آخرش هم یک بخش کامل سؤالات مصاحبه.
بخش ۰ — واژههایی که باید بشناسی
قبل از هر چیز، چند کلمه که مدام تکرار میشوند را با یک تصویر ساده جا بیندازیم تا بعداً سردرگم نشوی.
- همروندی (concurrency): یعنی چند کار را طوری مدیریت کنی که انگار همزمان پیش میروند. مثل یک آشپز که همزمان چند قابلمه را روی گاز دارد و بینشان میچرخد.
- مسدودکننده (blocking): کاری که نخ (thread) را میخواباند تا منتظر یک اتفاق بماند — مثلاً منتظر جواب دیتابیس یا شبکه. نخِ مسدودشده هیچ کار دیگری نمیکند، فقط منتظر است.
- وابسته به CPU (CPU-bound): کاری که واقعاً پردازنده را مشغول میکند، مثل محاسبهی ریاضی سنگین. اینجا نخ نمیخوابد، دارد عرق میریزد.
- نخ (thread): یک خط اجرای مستقل. تصورش کن یک کارگر که یک لیست از دستورها را پشت سر هم انجام میدهد.
- executor / استخر نخ (thread pool): یک جعبهی از پیش آماده از کارگرها. بهجای اینکه هر بار کارگر جدید استخدام کنی (گران است)، از استخر یکی را برمیداری، کار را میدهی، و برمیگردانیاش.
دو نوع کار داریم و دو نوع ابزار. کارِ CPU-bound (محاسبه) با ForkJoin؛ کارِ مسدودکننده (I/O مثل شبکه و دیتابیس) با executor اختصاصی یا virtual thread. قاطی کردنِ این دو، ریشهی تقریباً همهی فاجعههاست.
مدل ذهنی — دو ابزاری که همه اشتباه میگیرند
جاوا دو جعبهابزار همروندی دارد که ظاهرشان شبیه است ولی برای دو کار کاملاً مختلف ساخته شدهاند.
تصور کن یک رستوران داری. ForkJoinPool مثل آشپزخانه است: یک مشت آشپز حرفهای که کارشان بریدن، سرخ کردن و پختن است — کارهای کوتاه، سریع، و بدون انتظار. هیچ آشپزی وسط کار نمیایستد پشت تلفن منتظر تحویل بار بماند؛ همه دارند دستهایشان را تکان میدهند.
اما CompletableFuture آشپز نیست؛ مثل دستور کارِ روی کاغذ است: «اول سالاد را آماده کن، بعد وقتی استیک حاضر شد کنار هم بچین، اگر ماهی نبود جایگزینش را بگذار.» این کاغذ خودش هیچ کاری نمیکند و هیچ آشپزی استخدام نمیکند — فقط میگوید چه چیزی بعد از چه چیزی و کجا انجام شود. و بهصورت پیشفرض، این کاغذ کارها را میسپارد به همان آشپزخانهی ForkJoin.
پس دقیقتر:
ForkJoinPoolیک موتور کارِ وابسته به CPU (CPU-bound) است. کل طراحیاش — صفهای دوسر (deque) با سرقت کار (work-stealing)، و زمانبندی محلی LIFO — بر این فرض بنا شده که وظایف کوتاه، غیرمسدودکننده و بهصورت بازگشتی قابل تقسیماند. این استخر ستون فقراتparallelStream()و executor پیشفرضCompletableFutureاست.CompletableFutureیک گراف ترکیب ناهمگام (asynchronous composition graph) است؛ یعنی یک promise/future با ترکیبکنندهها (combinatorها). این کلاس خودش نخ نمیسازد و مسدود شدن را مدیریت نمیکند؛ فقط تعیین میکند callbackهای تو کجا اجرا شوند — و بهصورت پیشفرض این «کجا» همان common pool مربوط به ForkJoin است.
CompletableFuture یک گراف وابستگی از callbackهاست، و common pool مربوط به ForkJoin مکان پیشفرض اجرای آن callbackهاست. تقریباً همهی باگهای تولیدی این حوزه از این میآیند که کارِ مسدودکننده را روی آن استخرِ غیرمسدودکننده اجرا کنیم.
هر جا در ادامه گیج شدی، برگرد به همین جمله. کل درس دارد این یک خط را باز میکند.
بخش ۱ — CompletableFuture
ساختن future
قبل از اینکه بتوانی کارها را به هم بچسبانی، باید یک future بسازی. چند راه داریم:
// supplier را روی common pool اجرا میکند و بلافاصله برمیگردد
CompletableFuture<String> f1 = CompletableFuture.supplyAsync(() -> fetchUser());
// بهجای آن روی executor دادهشده اجرا میشود
ExecutorService io = Executors.newFixedThreadPool(16);
CompletableFuture<String> f2 = CompletableFuture.supplyAsync(() -> fetchUser(), io);
// مقادیر از پیش تکمیلشده (هیچ کار ناهمگامی ندارد)
CompletableFuture<Integer> done = CompletableFuture.completedFuture(42);
// futureای که خودت دستی کاملش میکنی — الگوی «promise»
CompletableFuture<String> promise = new CompletableFuture<>();
// ... بعداً، از هر نخی:
promise.complete("value"); // یا promise.completeExceptionally(ex);
بیا اینها را باز کنیم:
supplyAsyncیعنی «این کار را برو یک جای دیگر انجام بده و به من یک قبضِ رسید بده که بعداً جوابش را بگیرم.» آن «جای دیگر» بهصورت پیشفرض common pool است، مگر اینکه مثلf2یک executor خودت بدهی.completedFutureیعنی «جواب همین الان آماده است، اصلاً کاری برای انجام نیست.» مثل یک قبض که رویش از قبل جوابش نوشته شده.- الگوی promise: یک
CompletableFutureخالی میسازی و بعداً از هر نخ دلخواهی باcomplete(...)پُرش میکنی. این وقتی بهدرد میخورد که مثلاً منتظر یک callback از یک کتابخانهی قدیمی هستی.
runAsync هم خواهرِ بدونخروجیِ supplyAsync است (یک Runnable میگیرد و CompletableFuture<Void> میدهد) — وقتی فقط میخواهی کاری انجام شود و جوابی لازم نداری.
supplyAsync مشتاق (eager) است: کار همین لحظه که صدایش میزنی شروع میشود، نه وقتی جوابش را میخواهی. این با بعضی کتابخانههای reactive که «تنبل» هستند فرق دارد — و بعداً در همین درس رویش برمیگردیم چون یک سؤال مصاحبهی رایج است.
سه خانوادهی ترکیبکنندهها
حالا که یک future داری، میخواهی رویش کار کنی: تبدیلش کنی، به کارِ بعدی بچسبانی، یا با یک future دیگر ترکیبش کنی. این متدها را combinator (ترکیبکننده) میگویند. در سه محور فکر کن:
| هدف | متد | نوع تابع |
|---|---|---|
| تبدیل مقدار | thenApply |
T -> U |
| زنجیر کردن یک مرحلهی ناهمگام دیگر (flatten) | thenCompose |
T -> CompletableFuture<U> |
| مصرف مقدار، بدون خروجی | thenAccept / thenRun |
T -> void / () -> void |
| ترکیب دو future مستقل | thenCombine |
(T, U) -> V |
| انتظار برای چندتا | allOf / anyOf |
— |
مهمترین جفتِ این جدول، thenApply در برابر thenCompose است.
thenApply مثل این است که داخل یک جعبهی هدیه دست ببری، شیء را برداری، رنگش کنی و بگذاری سرِ جایش. یک جعبه، یک شیء. اما بعضی وقتها کاری که میکنی خودش یک جعبهی دیگر برمیگرداند (چون آن کار هم ناهمگام است). اگر باز از thenApply استفاده کنی، به یک جعبه داخل جعبه میرسی — CompletableFuture<CompletableFuture<U>>. thenCompose این پوستهی اضافه را باز میکند و فقط یک جعبه به تو میدهد.
به زبان برنامهنویسی تابعی: thenApply دقیقاً map است و thenCompose دقیقاً flatMap. قانون ساده: اگر تابعت یک مقدار ساده برمیگرداند از thenApply استفاده کن؛ اگر یک future دیگر برمیگرداند از thenCompose.
CompletableFuture<Order> pipeline =
CompletableFuture.supplyAsync(() -> loadCart(userId)) // CF<Cart>
.thenApply(cart -> priceCart(cart)) // CF<PricedCart> (تبدیل همگام)
.thenCompose(priced -> chargeAsync(priced)) // CF<Payment> (future برمیگرداند -> flatten)
.thenApply(payment -> new Order(payment)); // CF<Order>
ببین priceCart یک مقدار عادی میدهد پس thenApply است؛ اما chargeAsync خودش یک CompletableFuture میدهد پس thenCompose است تا جعبهی تودرتو نگیریم.
thenCombine و allOf / anyOf
تا اینجا یک زنجیرهی خطی داشتیم. اما گاهی میخواهی دو کارِ مستقل را همزمان اجرا کنی و بعد جوابهایشان را کنار هم بگذاری. اینجا thenCombine وارد میشود:
CompletableFuture<Integer> price = CompletableFuture.supplyAsync(this::fetchPrice);
CompletableFuture<Integer> stock = CompletableFuture.supplyAsync(this::fetchStock);
// هر دو همزمان اجرا میشوند؛ وقتی هر دو تمام شدند ترکیب میشوند
CompletableFuture<Quote> quote =
price.thenCombine(stock, (p, s) -> new Quote(p, s));
نکتهی طلایی: چون price و stock هر دو با supplyAsync ساخته شدهاند، همزمان شروع میشوند. thenCombine فقط منتظر میماند تا هر دو تمام شوند و بعد تابع را صدا میزند. یعنی زمان کل تقریباً برابر با کندترینشان است، نه جمعشان.
وقتی نه دو تا بلکه یک لیست از futureها داری، از allOf استفاده میکنی:
List<CompletableFuture<Product>> futures = ids.stream()
.map(id -> CompletableFuture.supplyAsync(() -> fetch(id), io))
.toList();
CompletableFuture<List<Product>> all =
CompletableFuture.allOf(futures.toArray(CompletableFuture[]::new))
.thenApply(v -> futures.stream()
.map(CompletableFuture::join) // امن: اینجا همه از قبل تکمیل شدهاند
.toList());
allOf یک CompletableFuture<Void> برمیگرداند — یعنی فقط میگوید «همه تمام شدند» ولی خودِ مقادیر را دور میریزد. برای همین بعد از اینکه تمام شد، باید دوباره روی هر future جداگانه join() بزنی تا جوابش را بگیری. اینجا join() امن است چون میدانیم همه از قبل کامل شدهاند و اصلاً مسدود نمیشود.
anyOf برعکسِ allOf است: با نتیجهی اولین futureای که تمام شود کامل میشود. برای درخواستهای hedged (یک درخواست را به دو سرور بفرست، هرکدام زودتر جواب داد را بردار) یا timeout مفید است. فقط یک زخم دارد: چون نمیداند بین چند نوع مختلف چه نوع مشترکی هست، CompletableFuture<Object> برمیگرداند و باید خودت cast کنی.
انتشار خطا: exceptionally، handle، whenComplete
اینجاست که کاندیداهای ارشد از بقیه جدا میشوند. اول یک نکتهی کلیدی که خیلیها را گیج میکند:
وقتی داخل یک مرحله (مثلاً یک thenApply) یک استثنا رخ میدهد، برخلاف کد معمولی، آن استثنا در همانجا که lambda را نوشتی پرتاب نمیشود. بهجایش، آن مرحله را «بهصورت استثنایی کامل» میکند و استثنا در طول زنجیره به پایین سُر میخورد تا به یک هندلر برسد — و در راه، پیچیده میشود درون یک CompletionException.
CompletableFuture<String> result =
CompletableFuture.supplyAsync(() -> { throw new IllegalStateException("boom"); })
.thenApply(s -> s.toUpperCase()) // رد میشود — بالادست شکست خورد
.exceptionally(ex -> "fallback"); // بازیابی: ex یک CompletionException روی ISE است
ببین thenApply اصلاً اجرا نشد چون مرحلهی قبلش شکست خورده بود؛ استثنا از رویش پرید و مستقیم رفت سراغ exceptionally. سه ابزار برای مدیریت خطا داریم و فرقشان دقیقاً همان چیزی است که ازت در مصاحبه میپرسند:
| متد | مقدار را میبیند | استثنا را میبیند | میتواند مقدار را تبدیل کند | میتواند از خطا بازیابی کند | برمیگرداند |
|---|---|---|---|---|---|
exceptionally |
خیر | بله | خیر | بله | CompletableFuture<T> |
handle |
بله | بله | بله | بله | CompletableFuture<U> (نوع جدید) |
whenComplete |
بله | بله | خیر (فقط side-effect) | خیر (استثنای اصلی را دوباره پرتاب میکند) | CompletableFuture<T> |
تصور کن یک محصول روی تسمهنقاله میآید و ممکن است سالم باشد یا خراب. exceptionally مثل کارگری است که فقط وقتی محصول خراب است دست به کار میشود و یک جایگزین میگذارد — اگر سالم بود اصلاً بیدار نمیشود. handle مثل بازرسی است که هر محصول را میبیند (چه سالم چه خراب)، میتواند سالم را هم عوض کند و خراب را هم تعمیر کند، و حتی نوع محصول را تغییر دهد. whenComplete مثل دوربین مداربسته است: هر چیزی را میبیند و ثبت میکند (log، پاکسازی)، ولی هیچچیزی را عوض نمیکند — محصول خراب همانطور خراب میرود پایین.
تلههای کلیدی که باید حفظ باشی:
- استثنایی که میگیری معمولاً یک
CompletionExceptionاست که علت واقعی را میپیچد. برای دیدن علت واقعی حتماًex.getCause()را صدا بزن. whenCompleteاستثنا را نمیبلعد — future بازگشتی همچنان شکست میخورد. این برای log و پاکسازی است، نه بازیابی.handleروی هر دو مسیر اجرا میشود؛ تنها combinatorی است که میگذارد شکست را به موفقیت تبدیل کنی و نوع را عوض کنی.- هر combinator یک نسخهی
*Asyncهم دارد (exceptionallyAsync،handleAsync) اگر لازم باشد بازیابی به نخ دیگری بپرد.
CompletableFuture<Response> safe =
callServiceAsync()
.orTimeout(2, TimeUnit.SECONDS) // جاوا ۹+: بعد از ۲ ثانیه با TimeoutException شکست
.handle((resp, ex) -> {
if (ex != null) {
log.warn("degraded", ex.getCause()); // بازکردن CompletionException
return Response.degraded();
}
return resp;
});
orTimeout(2, SECONDS) میگوید «اگر تا ۲ ثانیه جواب نیامد، با TimeoutException شکست بخور.» و completeOnTimeout(defaultValue, ...) میگوید «اگر دیر شد، بهجای شکست، این مقدار پیشفرض را بگذار.» هر دو از جاوا ۹ اضافه شدند و کار قشنگشان این است که بدون یک scheduler یا نخِ زمانسنجِ جداگانه، یک future را زمانبندی محدود میکنند.
کدام executor چه چیزی را اجرا میکند — قانونی که همه را زمین میزند
حالا رسیدیم به قلبِ همهچیز. اگر فقط یک بخش از این درس را قرار است حفظ کنی، همین است. دو قانون ساده که پشتِ درصد بزرگی از باگهای همروندی جاوا هستند:
قانون ۱ — combinatorهای غیر-async (thenApply، thenCombine، handle، …) روی همان نخی که future بالادست را کامل کرده اجرا میشوند. اگر بالادست هنگام اتصال callback از قبل کامل شده باشد، روی نخ فراخوانِ خودت، همگام (synchronous)، اجرا میشود.
قانون ۲ — combinatorهای *Async بدون executor روی common pool مربوط به ForkJoin اجرا میشوند. اگر آرگومان executor بدهی، آنجا اجرا میشوند.
CompletableFuture.supplyAsync(() -> load(), io) // روی 'io' اجرا میشود
.thenApply(x -> step1(x)) // روی 'io' (نخ تکمیلکننده)
.thenApplyAsync(x -> step2(x)) // به ForkJoinPool.commonPool() میپرد
.thenApplyAsync(x -> step3(x), io); // روی 'io' اجرا میشود
خطبهخط بخوانش: load() روی استخر io اجرا میشود. حالا step1 یک thenApply غیر-async است، پس طبق قانون ۱ روی همان نخِ io اجرا میشود که کار قبلی را تمام کرد. بعد step2 یک thenApplyAsync بدون executor است، پس طبق قانون ۲ میپرد روی common pool. و step3 یک thenApplyAsync با executor io است، پس برمیگردد روی io.
تصور کن یک پرونده روی میز کارمندِ بخش IO تمام میشود. حالا مرحلهی بعدی یک کار «غیر-async» است. بهجای اینکه پرونده برود توی صف مرکزی، همان کارمندِ IO بلافاصله مرحلهی بعد را هم انجام میدهد. اگر مرحلهی بعد کوتاه باشد اشکالی ندارد؛ اما اگر یک زنجیرهی طولانی از این کارها باشد، این کارمندِ بختبرگشتهی IO آنقدر درگیر کارهای دیگران میشود که برای کارِ اصلی خودش (پاسخ به شبکه) وقت نمیماند — و استخرِ IO تو قحطیزده میشود.
پیامدهایی که در تولید گاز میگیرند:
- اگر
supplyAsync(..., io)کامل شود، آنگاهstep1هم یک نخ ازioرا اشغال میکند — یک زنجیرهی طولانیthenApplyمیتواند استخر IO تو را قحطیزده کند، حتی اگر فقط یک مرحله «async» بوده باشد. - اگر future از قبل تمام شده باشد (مثلاً
completedFuture)،thenApplyتو بهصورت inline روی نخِ فراخوان اجرا میشود — گاهی در تستها غافلگیرکننده است، چون فکر میکنی کار در پسزمینه میرود ولی همانجا سرِ جای خودت اجرا میشود. - چون پیشفرض common pool است، هر فراخوان مسدودکننده درون یک combinator پیشفرض، یک worker از common pool را مسدود میکند، و بهصورت پیشفرض تنها
Runtime.availableProcessors() - 1تای آنها وجود دارد. اگر بهاندازهی کافی مسدود کنی، کل ماشینِ موازیِ JVM (streamها، futureهای دیگر) متوقف میشود. همیشه برای کار مسدودکننده یک executor صریح پاس بده.
یک مثال واقعگرایانه از ترکیب
بیا همهچیز را کنار هم بگذاریم: سه سرویس را همزمان صدا بزن، جوابهایشان را ترکیب کن، هرکدام timeout خودش را داشته باشد، و اگر چیزی خراب شد به یک نسخهی تنزلیافته برگرد.
// فراخوان همزمان ۳ سرویس، ترکیب، با timeout هر فراخوان و یک fallback سراسری.
CompletableFuture<Profile> profile = supplyAsync(() -> userSvc.get(id), io)
.orTimeout(500, MILLISECONDS);
CompletableFuture<List<Order>> orders = supplyAsync(() -> orderSvc.recent(id), io)
.exceptionally(ex -> List.of()); // سفارشها اختیاریاند
CompletableFuture<Score> score = supplyAsync(() -> riskSvc.score(id), io)
.orTimeout(300, MILLISECONDS);
CompletableFuture<Dashboard> dashboard =
profile.thenCombine(orders, Dashboard::new) // (Profile, List<Order>) -> Dashboard
.thenCombine(score, Dashboard::withScore) // افزودن نتیجهی سوم
.exceptionally(ex -> Dashboard.EMPTY); // اگر profile/score شکست خورد، تنزل بده
Dashboard d = dashboard.join(); // یک بار در لبه، فراخوان را مسدود کن
دقت کن به فلسفهاش: هر سرویس روی io (استخر اختصاصی، نه common pool) اجرا میشود، هرکدام timeout دارد، orders که «اختیاری» است اگر شکست بخورد به لیست خالی برمیگردد، و در انتها فقط یک بار، در لبهی گراف، با join() مسدود میشویم. درونِ گراف کاملاً غیرمسدودکننده است.
join() یک CompletionException بدونبررسی (unchecked) پرتاب میکند؛ get() استثناهای بررسیشدهی InterruptedException و ExecutionException پرتاب میکند. چون در lambdaهای stream نمیتوانی استثنای checked پرتاب کنی، آنجا join() را ترجیح بده — تمیزتر است.
بخش ۲ — ForkJoinPool و سرقت کار (work-stealing)
حالا برویم سراغ آن آشپزخانه. ForkJoinPool برای یک نوع خاص از موازیسازی ساخته شده: تقسیموحل (divide-and-conquer).
مدل
ایده ساده است: یک وظیفهی بزرگ را به دو نصفه تقسیم کن (fork)، هر نصفه را جداگانه حل کن، بعد جوابها را با هم ادغام کن (join). و اگر نصفه هنوز بزرگ بود، خودش را باز نصف کن — تا وقتی به تکههای کوچکِ قابلحل برسیم.
قلبِ فنیِ ForkJoin در نحوهی توزیع کار بین نخهاست. هر نخِ worker یک صف دوسر (deque) دارد — صفی که هم از سر و هم از دُم میشود کار برداشت.
هر آشپز یک دسته کار روی میزِ خودش دارد. وقتی خودش کاری میسازد، آن را میگذارد روی همانهایی که تازه گذاشته و از همانجا (سرِ میز) هم برمیدارد — چون کارِ تازه «داغ» است، مواد و ابزارش هنوز دمِ دست است. این میشود LIFO (آخرین ورودی، اولین خروجی) و باعث میشود کش پردازنده بهینه بماند.
اما وقتی یک آشپزِ دیگر بیکار شد، میرود سراغ میزِ یک همکارِ شلوغ و از ته میزِ او (دُم) یک کار برمیدارد — یعنی قدیمیترین کار، که معمولاً هنوز نصف نشده و بزرگترین تکه است، و همکار هم فعلاً به آن نیاز فوری ندارد. این میشود سرقت کار (work-stealing) بهصورت FIFO.
پس مالک از سرِ deque بهصورت LIFO push/pop میکند (محلیّت کش خوب، کارهای اخیر داغاند)، و workerهای بیکار از دُمِ deque یک قربانی بهصورت FIFO سرقت میکنند (قدیمیترین، بزرگترین، کماحتمالترین برای نیاز فوری). این عدمتقارنِ LIFO-محلی + FIFO-سرقت همهی هستهها را بدون رقابتِ مرکزی مشغول نگه میدارد — نکتهی کلیدیِ کل طراحی همین است.
مثال RecursiveTask
بیا این را با کد ببینیم. میخواهیم جمعِ یک آرایهی بزرگ را موازی محاسبه کنیم. RecursiveTask<V> را extend میکنیم چون جواب برمیگرداند (Long):
class SumTask extends RecursiveTask<Long> {
private static final int THRESHOLD = 10_000;
private final long[] arr; private final int lo, hi;
SumTask(long[] arr, int lo, int hi) { this.arr = arr; this.lo = lo; this.hi = hi; }
@Override protected Long compute() {
if (hi - lo <= THRESHOLD) { // حالت پایه: مستقیم محاسبه کن
long s = 0;
for (int i = lo; i < hi; i++) s += arr[i];
return s;
}
int mid = (lo + hi) >>> 1;
SumTask left = new SumTask(arr, lo, mid);
SumTask right = new SumTask(arr, mid, hi);
left.fork(); // چپ را بهصورت ناهمگام زمانبندی کن
long rightResult = right.compute(); // راست را در همین نخ محاسبه کن (هر دو را fork نکن!)
long leftResult = left.join(); // منتظر چپِ دزدیده/صفشده بمان
return leftResult + rightResult;
}
}
long total = ForkJoinPool.commonPool().invoke(new SumTask(data, 0, data.length));
یک شاخه را fork کن، شاخهی دیگر را inline محاسبه کن، سپس join کن. توجه کن که left.fork() میکنیم ولی right.compute() را همانجا در نخ فعلی صدا میزنیم و آخرش left.join(). اگر هر دو را fork میکردیم و بعد هر دو را join، یک جای صف و یک سرقت را الکی هدر میدادیم؛ محاسبهی یکی بهصورت inline نخ فعلی را بهرهور نگه میدارد و شاخهی دیگر آماده است تا اگر لازم شد دزدیده شود.
RecursiveAction هم دقیقاً همین است ولی برای وقتی که جوابی برنمیگردانی (نسخهی Void).
انتخابِ THRESHOLD هنر است: اگر خیلی بزرگش کنی، موازیسازی کافی نمیشود و هستهها بیکار میمانند؛ اگر خیلی کوچکش کنی، سربارِ fork/join خودش از کارِ واقعی بیشتر میشود. آن را طوری بگذار که هر برگ بهاندازهی کافی بزرگ باشد که سربار را ناچیز کند، اما بهاندازهی کافی کوچک باشد که بار خوب متعادل شود.
common pool و درجهی موازیاش
حالا یک نکتهی بسیار مهم: ForkJoinPool.commonPool() یک تکنمونهی (singleton) در سطح کل JVM است. یعنی یک استخرِ مشترک برای کل برنامه.
درجهی موازیاش (parallelism) بهصورت پیشفرض Runtime.getRuntime().availableProcessors() - 1 است. یعنی روی یک ماشین ۸ هستهای، ۷ worker داری (چرا منهای یک؟ چون نخِ ثبتکنندهی کار خودش هم میتواند دست به کار شود و کمک کند). قابل تنظیم است:
-Djava.util.concurrent.ForkJoinPool.common.parallelism=32
چون این استخر بین parallelStream()، CompletableFuture پیشفرض، و هر کتابخانهای که سراغش برود به اشتراک گذاشته میشود، یک منبع سراسری است که تو مالکش نیستی. اندازهگذاریاش برای کارِ CPU (تقریباً برابر تعداد هسته) و بعد انجام I/O مسدودکننده روی همان، همان باگ کلاسیک قحطی است که کل برنامه را زمین میزند.
مسدود شدنِ درست: ManagedBlocker
فرض کن واقعاً چارهای نداری جز اینکه درون یک وظیفهی ForkJoin مسدود شوی. راهحلِ رسمی این است که به استخر خبر بدهی، تا بتواند یک نخ جبرانی (compensation thread) بالا بیاورد و درجهی موازی را حفظ کند — یعنی جای آن نخی که خوابیده، یکی دیگر بیدار میکند تا هستهها بیکار نمانند:
ForkJoinPool.managedBlock(new ForkJoinPool.ManagedBlocker() {
private String result;
@Override public boolean block() throws InterruptedException {
result = blockingRemoteCall(); // فراخوان مسدودکنندهی واقعی
return true; // true = تمام شد، نیازی به مسدود شدن دوباره نیست
}
@Override public boolean isReleasable() {
return result != null;
}
});
تصور کن یک آشپزِ تو مجبور شد برود پشت تلفن و منتظر بماند. اگر خبر ندهی، یک هسته عملاً بیکار میماند. ManagedBlocker مثل زنگی است که میزنی و میگویی «من رفتم پشت تلفن، یک آشپزِ موقت بیار.» استخر یک نخِ جبرانی اضافه میکند تا درجهی موازیِ واقعی افت نکند. وقتی برگشتی، اوضاع به حالت عادی برمیگردد.
parallelStream() روی مسیرهای پر از synchronized/Collections و فراخوانهای JDBC دقیقاً همان مواردی هستند که مسدود شدنِ سادهلوحانه بیسروصدا توان عملیاتی (throughput) را افت میدهد. ManagedBlocker دریچهی فرارِ مجاز است — اما صادقانه بگویم، روی JDKهای مدرن پاسخِ تمیزتر یک executor اختصاصی یا virtual threadهاست.
بخش ۳ — جزئیات درونی parallel stream
حالا که ForkJoin را میشناسی، parallelStream() دیگر جادو نیست. وقتی collection.parallelStream() (یا stream().parallel()) میزنی، خط لوله روی همان common pool اجرا میشود و منبع را از طریق یک شیء بهنام Spliterator تقسیم میکند.
اسمش از Split + Iterator میآید. یک Iterator عناصر را یکییکی میدهد؛ یک Spliterator علاوه بر آن میتواند خودش را نصف کند و بگوید «این نصفِ عناصر مالِ تو، آن نصف مالِ من.» اینکه یک منبع چقدر خوب و ارزان نصف میشود، مستقیماً تعیین میکند parallel stream رویش سریع است یا کند.
قواعد سرانگشتی که مستقیم از این جزئیات درونی میآیند:
- موازیسازی فقط وقتی کمک میکند که: منبع یکنواخت و ارزان تقسیم شود (آرایهها،
ArrayList،IntStream.range— عالی؛ چون میشود وسطشان را با یک محاسبه پیدا کرد. اماLinkedListو منابع مبتنی برIterator— ضعیف، چون برای رسیدن به وسط باید از اول قدم بزنی)، کارِ هر عنصر سنگین از نظر CPU باشد، و N بزرگ باشد. یک اکتشافِ تقریبی: N × هزینههرعنصر باید خیلی بیشتر از حدود ۱۰۰ میکروثانیه باشد تا سربارِ موازیسازی صرف کند. - برای reduction، خط لوله باید بدونحالت (stateless)، غیرمداخلهگر (non-interfering) و شرکتپذیر (associative) باشد. یعنی: تابعت نباید حالتِ بیرونی نگه دارد یا منبع را وسطِ کار عوض کند، و ترتیبِ گروهبندی نباید جواب را عوض کند.
reduceباید شرکتپذیر باشد؛collectاز یکcombinerبرای چسباندن نتایج جزئی استفاده میکند. - هرگز در یک parallel stream کار I/O انجام نده — workerهای common pool را مسدود و بقیه را قحطیزده میکنی. (همان جملهی کلیدی فصل، بار دیگر.)
forEachهیچ ترتیبی تضمین نمیکند؛ اگر به ترتیبِ برخورد (encounter order) نیاز داری ازforEachOrderedاستفاده کن — که البته هزینه دارد چون باید هماهنگی برقرار کند.
یعنی (a op b) op c باید همان جوابِ a op (b op c) را بدهد. جمع و ضرب شرکتپذیرند، اما تفریق نه ((۵-۳)-۲ ≠ ۵-(۳-۲)). چون parallel stream تکهها را در ترتیبِ نامعلوم گروهبندی و ادغام میکند، اگر عملیاتت شرکتپذیر نباشد، جوابِ موازی با جوابِ ترتیبی فرق میکند — و این یکی از بدترین باگهاست چون بیسروصداست.
// وابسته به CPU، منبع قابلتقسیم، reduction شرکتپذیر — انتخاب خوب
long primes = LongStream.rangeClosed(2, 10_000_000)
.parallel()
.filter(this::isPrime)
.count();
این یک نمونهی عالی است: منبع (rangeClosed) کاملاً ارزان نصف میشود، کارِ هر عنصر (isPrime) واقعاً CPU-bound است، و N هم بزرگ است.
و یک ترفندِ شناختهشده: اگر بخواهی یک parallel streamِ خاص را از common pool فراری بدهی (مثلاً چون نمیخواهی بقیهی برنامه را تحتتأثیر بگذارد)، آن را درون ForkJoinPool خودت submit کن:
ForkJoinPool custom = new ForkJoinPool(4);
long r = custom.submit(() ->
list.parallelStream().filter(this::heavy).count()
).join(); // این parallel stream حالا روی 'custom' اجرا میشود، نه common pool
چون parallel stream روی همان نخی که آن را اجرا میکند خطش را میبندد، و وقتی داخل custom.submit(...) هستی، آن نخ یکی از workerهای custom است. پس تقسیمها هم روی custom انجام میشوند نه common pool. یک ترفند رسمی نیست ولی رفتارِ قابلاتکایی است.
بخش ۴ — future در برابر reactive (بهاختصار)
گاهی در مصاحبه میپرسند «فرق CompletableFuture با reactive چیست؟» جوابش کوتاه اما مهم است.
CompletableFuture یک مقدارِ ناهمگامِ تکی را با ترکیب مدلسازی میکند — «یک جواب، بعداً میآید.» اما Reactive Streams (Mono/Flux در Project Reactor، یا RxJava) جریانهایی از ۰ تا N مقدار با backpressure را مدل میکنند و مجموعهی عملگرِ بسیار غنیتری دارند (retry، buffer، window، merge، کنترل جریان).
CompletableFuture مثل انتظار برای یک بستهی پستی است: میآید یا نمیآید، تمام. Flux مثل ایستادن کنار یک نوار نقاله است که ممکن است هیچ بستهای، یا ده، یا دههزار بسته رویش بیاید، و تو میتوانی بگویی «آهستهتر بفرست، دستم پُر است» — این همان backpressure است، یعنی مصرفکننده به تولیدکننده فشار برمیگرداند که سرعتش را کم کند.
قاعدهی انتخاب: از CF برای «این سرویسها را صدا بزن، ترکیب کن، یکبار برگردان» استفاده کن؛ از reactive وقتی جریان، backpressure داری یا به ترکیب در طول زمان نیاز داری. Mono در Reactor تقریباً یعنی «CF با صدها عملگر و subscription تنبل.»
CompletableFuture مشتاق (eager) است — کار در لحظهی ساخت شروع میشود. اما Mono/Flux تنبل (lazy) است — تا وقتی subscribe() نکنی، هیچ چیزی اجرا نمیشود. این تفاوت روی retry، caching و اینکه «آیا ساختنِ یک future بدونِ مصرفش کاری انجام میدهد یا نه» اثر مستقیم دارد.
بخش ۵ — virtual threadها معادله را عوض میکنند (جاوا ۲۱)
حالا بزرگترین تحول. تا اینجا کلِ فلسفهی CompletableFuture این بود که «نخهای سیستمعامل گراناند، پس روی مسدود شدن هدرشان نده — بهجایش callback بنویس.» اما virtual threadها (نخهای مجازی، JEP 444، پایدار در جاوا ۲۱) این فرض را زیر و رو میکنند: مسدود شدن را دوباره ارزان میکنند.
نخِ سیستمعامل مثل استخدامِ یک کارمندِ تماموقتِ گران است؛ نمیتوانی هزاران تا داشته باشی، پس نباید بگذاری بیکار پشت تلفن بماند. اما virtual thread مثل یک کارگرِ فوقارزان است که میلیونها تای آن را میتوانی داشته باشی. حالا دیگر مهم نیست یکیشان پشت تلفن منتظر بماند — JVM آن را از روی هسته برمیدارد، هسته را به یکی دیگر میدهد، و وقتی جواب آمد دوباره برش میگرداند. پس دیگر لازم نیست کدِ پیچیدهی callback بنویسی.
با virtual threadها میتوانی کدِ مسدودکنندهی خطی و ساده بنویسی و هزاران تای آن را همزمان اجرا کنی:
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
var f1 = executor.submit(() -> userSvc.get(id)); // فراخوان مسدودکننده، اما روی یک vthread
var f2 = executor.submit(() -> orderSvc.recent(id));
return new Dashboard(f1.get(), f2.get()); // ساختارمند، خوانا، بدون زنجیره
}
ببین چقدر خوانا است — بدون هیچ thenCombine و زنجیرهای، فقط «این دو کار را همزمان انجام بده و جوابشان را بگیر.»
virtual threadها توسط ForkJoinPool مخصوص خودشان در حالت FIFO زمانبندی میشوند (بهصورت پیشفرض یک نخِ carrier بهازای هر هسته، قابلتنظیم با jdk.virtualThreadScheduler.parallelism) — این استخر جدا از common pool است. پس مسدود شدن روی virtual thread، common pool را که مالِ parallel streamها و CF پیشفرض است، درگیر نمیکند.
پس اینها رقیب هم نیستند، مکملاند:
- از
CompletableFutureبرای عملگرهای ترکیبش استفاده کن، وقتی واقعاً به یک گراف نیاز داری (thenCombine،allOf، timeoutها). - به virtual threadها (یا Structured Concurrency با
StructuredTaskScope) روی بیاور وقتی منطق واقعاً «اینها را موازی انجام بده و join کن» است. - ForkJoin را برای موازیسازیِ خالصِ CPU-bound نگه دار.
تلهها و نکات ظریف
بیا همهی تلههایی که در راه دیدیم را یکجا جمع کنیم — اینها همانهاییاند که در تولید گاز میگیرند:
- مسدود شدن روی common pool.
supplyAsync/parallelStreamپیشفرض + JDBC/HTTP = قحطی. یکExecutorاختصاصی پاس بده. - پیچیده شدن در
CompletionException.catch/exceptionallyتو wrapper را میبیند؛getCause()را صدا بزن. whenCompleteبازیابی نیست. دوباره پرتاب میکند؛ فقطhandle/exceptionallyبازیابی میکنند.allOfمقادیر را از دست میدهد. بعد از تکمیلِallOf، هر future را دوباره join کن.- combinatorهای غیر-async نخِ تکمیلکننده را میدزدند. یک زنجیرهی طولانی
thenApplyروی هر نخی که مرحلهی قبل را تمام کرده اجرا میشود — شاید استخر IO تو، شاید یک worker از common pool. - فراموش کردن اینکه future مشتاق است.
supplyAsyncهمین حالا شروع میشود؛ هیچ subscribe تنبلی وجود ندارد. - fork کردن هر دو شاخه. یکی را inline محاسبه کن؛ دیگری را fork کن.
- fan-out بیکران. نگاشتِ ۱۰۰ هزار آیتم به
supplyAsync(..., commonPool)آن را اشباع میکند — همروندی را با یک executor اندازهگذاریشده یا یک semaphore محدود کن.
بهترین روشها
۱. مالکِ executorهای خودت باش برای هر چیزی که مسدود میشود. هرگز common pool را برای I/O باور نکن.
۲. نخها را نامگذاری کن (ThreadFactory) تا stack dumpها قابلدیباگ باشند.
۳. join() را در لبهی گراف ترجیح بده؛ درونِ گراف را غیرمسدودکننده نگه دار.
۴. به هر فراخوان خارجی orTimeout/completeOnTimeout اضافه کن.
۵. خطاها را با handle/exceptionally مدیریت کن؛ همیشه getCause().
۶. RecursiveTask را فقط برای تقسیموحلِ وابسته به CPU با آستانهی تنظیمشده بهکار ببر.
۷. روی جاوا ۲۱+، برای fan-out مربوط به I/O، virtual threadها / structured concurrency را ترجیح بده؛ CompletableFuture را برای گرافهای ترکیبِ واقعی نگه دار؛ ForkJoin را برای موازیسازیِ وابسته به CPU.
سؤالات مصاحبه
حالا بخشی که برایش آمدهای. هر سؤال را جوری جواب میدهیم که هم درست باشد هم نشاندهندهی عمق.
روی نخی که future بالادست را کامل کرده. اگر بالادست هنگام اتصالِ callback از قبل کامل شده باشد، بهصورت همگام روی نخِ فراخوان اجرا میشود. فقط نسخههای *Async پرش به یک executor را تضمین میکنند (اگر هیچ executorی داده نشود، common pool).
thenApply همان map است (T -> U)؛ thenCompose همان flatMap (T -> CompletableFuture<U>). وقتی تابعت خودش ناهمگام است و یک future برمیگرداند، از thenCompose استفاده کن تا از CompletableFuture<CompletableFuture<U>> تودرتو پرهیز شود.
exceptionally فقط روی شکست اجرا میشود و میتواند بازیابی کند (همان نوع). handle روی هر دو مسیر اجرا میشود، میتواند مقدار را تبدیل و بازیابی کند و نوعِ نتیجه را عوض کند. whenComplete روی هر دو اجرا میشود اما فقط side-effect است — نمیتواند بازیابی کند و استثنای اصلی را به پاییندست دوباره پرتاب میکند.
چون join() یک CompletionException پرتاب میکند که علت واقعی را میپیچد. باید CompletionException را بگیری و getCause() را بررسی کنی، وگرنه نوعِ اصلی مطابقت نخواهد کرد.
parallel streamها روی common pool در سطح کل JVM اجرا میشوند (تقریباً هسته - 1 worker). یک فراخوان مسدودکنندهی JDBC یک worker را پارک میکند؛ بهاندازهی کافی از اینها استخر را قحطیزده میکنند و در همین فرآیند هر parallel stream و CompletableFuture پیشفرض را متوقف میکنند. راهحل: executor اختصاصی، ManagedBlocker، یا virtual threadها.
هر worker یک deque دارد؛ وظایف خودش را LIFO از سر push/pop میکند (کشداغ). workerهای بیکار از دُمِ deque یک worker دیگر FIFO سرقت میکنند (قدیمیترین/بزرگترین وظیفه). این کار بار را با حداقل رقابت متعادل میکند.
fork کردنِ هر دو و سپس join هر دو، یک جای صف را هدر میدهد و یک سرقت را اجباری میکند؛ محاسبهی یک شاخه بهصورت inline نخِ فعلی را مشغولِ کارِ مفید نگه میدارد در حالی که شاخهی دیگر برای سرقت در دسترس است. این فرمِ متعارف و کارآمدترین است.
availableProcessors() - 1. با -Djava.util.concurrent.ForkJoinPool.common.parallelism=N تغییرش بده. توجه کن که یک singleton سراسری است که در کل JVM به اشتراک گذاشته میشود.
var cf = CompletableFuture.supplyAsync(() -> 1)
.thenApply(x -> x / 0)
.thenApply(x -> x + 1)
.exceptionally(ex -> -1);
System.out.println(cf.join());
پاسخ: -1. عبارت / 0 یک ArithmeticException پرتاب میکند و آن مرحله را بهصورت استثنایی کامل میکند؛ thenApply بعدی رد میشود؛ exceptionally به -1 بازیابی میکند.
List<CompletableFuture<Integer>> fs = ...;
CompletableFuture.allOf(fs.toArray(new CompletableFuture[0]));
List<Integer> results = fs.stream().map(CompletableFuture::join).toList();
باگ: allOf(...) ساخته میشود اما هرگز join/await نمیشود، پس نتیجهاش دور ریخته میشود و هیچ تضمینی نیست که مراحلِ زمانبندیشده بعد از آن دچار race نشوند. مهمتر اینکه اگر هر futureای شکست بخورد، join() وسطِ جریان پرتاب میکند و بقیه را از دست میدهی. راهحل: .thenApply(...) را از futureِ مربوط به allOf زنجیر کن، یا با مدیریت استثنا برای هر future جداگانه جمعآوری کن.
مشتاق — supplyAsync بلافاصله شروع به اجرا میکند. Mono/Flux تنبلاند: تا subscribe() هیچ چیز اجرا نمیشود. این روی retry، caching و اینکه آیا ساختن-اما-مصرفنکردن کاری انجام میدهد اثر میگذارد.
خیر. virtual threadها روی یک ForkJoinPool جداگانه در حالت FIFO اجرا میشوند، بهصورت پیشفرض یک نخِ carrier بهازای هر هسته (jdk.virtualThreadScheduler.parallelism). common pool (LIFO) توسط parallel streamها و CompletableFuture پیشفرض استفاده میشود.
همهی ۱۰۰ هزار callback روی common pool فرود میآیند؛ اگر هر کدام مسدود شود یا حجم از هسته - 1 جای بهرهور فراتر رود، صفبندی/قحطی میگیری. یک executor با اندازهی مناسب پاس بده، یا fan-out را محدود کن.
وقتی به معناشناسیِ ترکیب آن نیاز داری: ترکیبِ نتایجِ ناهمگون (thenCombine)، انتظار برای چندتا (allOf/anyOf)، hedging (anyOf)، timeoutهای اعلانی (orTimeout)، و ساختنِ گرافهای ناهمگامِ قابلاستفادهی مجدد. برای «N فراخوان مسدودکننده را اجرا کن و join کن»، virtual threadها / StructuredTaskScope سادهترند.
چون combinatorهای غیر-async روی نخِ تکمیلکننده اجرا میشوند. اگر هر مرحله روی نخِ همان executor کامل شود (مثلاً supplyAsync(..., io) سپس thenCompose غیر-async که futureِ درونیاش هم روی io کامل میشود)، کلِ زنجیره میتواند روی workerهای io سریالی شود — یک منبعِ ظریفِ قحطیِ استخر. برای پرشِ عمدی، *Async(..., otherPool) قرار بده.
نکاتِ سنیور و موارد پیشرفته
تا اینجا فصل، ساز و کارِ CompletableFuture، ForkJoin و virtual thread را کامل یادت داد. اما آنچه یک سنیور را از یک میدلِ خوب جدا میکند، نه دانستنِ این APIها، بلکه دانستنِ جاهایی است که این APIها بیسر و صدا دروغ میگویند — کارهایی که فکر میکنی انجام میدهند ولی نمیدهند، و درست همانها هستند که ساعت ۳ بامداد صفحهی pager تو را روشن میکنند. این بخش دقیقاً همان لایه است.
(۱) چرا کنسل کردن یک CompletableFuture یک دروغ است و orTimeout نشتی (leak) میسازد؛ (۲) اینکه anyOf/allOf برادرها را کنسل نمیکنند؛ (۳) چطور یک future را در برابر دستکاریِ بیرونی محافظت کنی؛ (۴) رفتارِ common pool داخل کانتینر و Kubernetes؛ (۵) مسئلهی pinning در virtual threadها و رفعش در جاوا ۲۴؛ (۶) ThreadLocal، MDC و ScopedValue؛ (۷) Structured Concurrency بهعنوان جانشینِ واقعی؛ (۸) retry با backoff و exceptionallyCompose؛ (۹) تضمینهای حافظه (happens-before). آخرش هم چند سؤال مصاحبهی سنگین.
۱. کنسل کردنِ یک future یک دروغ است — و orTimeout نشتی میسازد
این نکتهای است که خودِ Javadoc هم صریح نمیگویدش و در تولید خون به پا میکند. تصور رایج این است: «اگر روی future صبرم تمام شد، cancel(true) میزنم و کارِ پشتش قطع میشود.» این غلط است.
CompletableFuture<String> f = CompletableFuture.supplyAsync(() -> slowDbCall(), io);
boolean cancelled = f.cancel(true); // true = mayInterruptIfRunning
پارامتر mayInterruptIfRunning در CompletableFuture کاملاً نادیده گرفته میشود. برخلاف FutureTask، یک CompletableFuture هیچ ارجاعی به نخی که supplyAsync را اجرا میکند نگه نمیدارد؛ پس نمیتواند آن نخ را interrupt کند. تنها کاری که cancel میکند این است که خودِ future را به حالتِ استثنایی (CancellationException) میبرد تا مرحلههای پاییندستِ گراف رها شوند. اما slowDbCall() تا آخر اجرا میشود و کانکشنِ دیتابیس تا لحظهی برگشتش checked-out میماند.
همین منطق برای orTimeout هم برقرار است:
CompletableFuture<Response> r = callServiceAsync().orTimeout(2, SECONDS);
orTimeout(2, SECONDS) بعد از ۲ ثانیه futureِ تو را با TimeoutException میبندد، ولی درخواستِ HTTP/JDBC پشتش همچنان زنده است و منبع (کانکشن، thread) را نگه داشته. زیرِ بار، همین میشود نشتیِ کانکشنپول: تایماوتها سریع رخ میدهند، ولی کارهای واقعی روی هم تلنبار میشوند تا pool خفه شود. تایماوتِ واقعی باید در خودِ کلاینت باشد (مثلاً HttpClient با .timeout(...))، نه در لایهی CompletableFuture.
نتیجهی سنیوری: کنسل در دنیای CompletableFuture تعاونی (cooperative) نیست مگر اینکه تو دستی سیمش را بکشی. یا کلاینتی به کار بگیر که واقعاً کنسل را propagate کند (مثل HttpClient.sendAsync که futureِ برگشتیاش به کنسلِ واقعی سیمکشی شده)، یا از StructuredTaskScope استفاده کن که هنگام شکستِ scope، subtaskها را واقعاً interrupt میکند.
۲. anyOf و allOf برادرها را نمیکشند
فصل گفت allOf مقادیر را دور میریزد و anyOf اولین جواب را برمیگرداند. اما دو ظرافتِ حیاتی جا ماند:
anyOfبا «اولین چیزی که کامل شود» کامل میشود — حتی اگر آن کاملشدن یک شکست باشد. یعنی یک شکستِ سریع بر یک موفقیتِ کند برنده میشود. اگر داری باanyOfدرخواستِ hedged میزنی (همان درخواست را به دو سرور بفرست، هرکدام زودتر جواب داد بردار)، یک خطای فوریِ سرورِ اول کلِ نتیجه را شکست میدهد در حالی که سرورِ دوم داشت درست جواب میداد.- هیچکدام از دو تا برادرهای بازنده را کنسل نمیکنند. بعد از اینکه
anyOfبرنده را برگرداند، همهی futureهای دیگر تا آخر اجرا میشوند و منبع میسوزانند.allOfهم بهمحضِ اولین شکست، بهصورت استثنایی کامل میشود ولی بقیهی futureها بیمراقب ادامه میدهند.
اگر واقعاً hedging میخواهی، بعد از اینکه anyOf برنده را داد، باید دستی روی بازندهها منطقِ کنسل/رهاسازیِ منبع بزنی، یا futureهایی بهکار ببری که کنسلشان واقعاً منبع را آزاد میکند. anyOf بهتنهایی فقط «اولین سیگنال» است، نه «مدیریتِ چرخهی عمرِ رقبا».
۳. یک future را به بیرون نده مگر محافظتشده
وقتی متدِ کتابخانهات CompletableFuture<T> برمیگرداند، هر فراخوانی میتواند روی آن complete(...), obtrudeValue(...) یا cancel() صدا بزند و کلِ گرافِ پاییندستِ تو را بدزدد — انگار یک متغیرِ private را public کرده باشی.
// بد: هر کسی میتواند نتیجهی تو را جعل کند
public CompletableFuture<Price> quote() { return internalFuture; }
// خوب: نسخهای که متدهای completion رویش استثنا میاندازند
public CompletionStage<Price> quote() { return internalFuture.minimalCompletionStage(); }
minimalCompletionStage() (جاوا ۹) یک CompletionStage میدهد که فقط combinatorها را دارد و صدا زدنِ complete/cancel/get رویش UnsupportedOperationException میدهد. copy() هم یک کپیِ وابسته میدهد که کاملکردنش روی اصل اثر ندارد. و obtrudeValue/obtrudeException را فقط در تست بهکار ببر — اینها یک futureِ ازقبلکاملشده را با زور بازنویسی میکنند و همهی invariantها را میشکنند.
۴. common pool داخل کانتینر — دامِ Kubernetes
فصل گفت common pool پیشفرض availableProcessors() - 1 نخ دارد. اما در دنیای کانتینر این عدد گاز میگیرد:
- از JDK 10 (و 8u191) به بعد
availableProcessors()نسبت به cgroup آگاه است و به CPU limit احترام میگذارد. - اگر pod تو یک vCPU limit دارد،
cores - 1میشود صفر و common pool عملاً به بیموازیسازی فرو میریزد:parallelStream()روی نخِ صداکننده سریال اجرا میشود. آن سرعتی که در لپتاپِ ۸ هستهای دیدی، در prod ناپدید میشود. - برعکس، اگر limit نگذاری، JVM ممکن است همهی هستههای نودِ میزبان را ببیند و poolای بسیار بزرگتر از سهمِ واقعیات بسازد.
به سایزِ خودکار در کانتینر اعتماد نکن. یا -XX:ActiveProcessorCount=N بگذار، یا مستقیم -Djava.util.concurrent.ForkJoinPool.common.parallelism=N. ضمناً یادت باشد نخهای common pool daemon هستند؛ هنگامِ خروجِ JVM کارهای ناتمامشان بدون تضمین رها میشوند — پس کارِ حیاتیِ «باید تمام شود» را به common pool نسپار.
۵. Pinning در virtual thread — دامِ مدرن و رفعش در جاوا ۲۴
فصل گفت virtual thread بلاکشدن را ارزان میکند، ولی یک ستارهی مهم دارد که تا جاوا ۲۳ خونریزی میکرد: pinning.
وقتی یک virtual thread وارد یک بلاکِ synchronized میشد، در نسخههای تا جاوا ۲۳ monitorِ آن قفل روی نخِ حاملِ (carrier) پلتفرمی ثبت میشد، نه روی خودِ virtual thread. نتیجه: تا وقتی داخلِ آن بلاک بودی، virtual thread به carrier میخکوب (pinned) بود و اگر همانجا بلاک میشدی (مثلاً I/O داخلِ synchronized)، carrier نمیتوانست آزاد شود و به virtual threadِ دیگری برسد. اگر همهی carrierها اینطور میخکوب میشدند، کلِ scheduler خفه میشد — و در بدترین حالت، deadlock.
JEP 491 در جاوا ۲۴ کاری کرد که monitor به خودِ virtual thread گره بخورد؛ حالا synchronized دیگر carrier را میخکوب نمیکند و میتوانی بیترس از آن استفاده کنی. اما هنوز دو چیز pin میکنند: فریمهای native (JNI) و بلاکشدن داخلِ یک متدِ native. پس pinning نمرده، فقط بسیار نادر شده.
در هاتپثهای I/O، synchronized را با ReentrantLock جایگزین کن — چون قفلهای java.util.concurrent پین نمیکنند. برای دیباگ، -Djdk.tracePinnedThreads (که در ۲۴ deprecated شد) یا رویدادِ JFRِ jdk.VirtualThreadPinned را روشن کن. و دو قانونِ همیشگی: virtual thread را pool نکن (خودشان ارزاناند)، و برای کارِ CPU-bound از virtual thread استفاده نکن — آنها برای کارِ بلاکشوندهاند نه محاسبهی سنگین.
۶. ThreadLocal، MDC و ScopedValue
با virtual threadها دو مشکلِ زمینهای (context) بزرگ میشود:
- حافظه: اگر میلیونها virtual thread داشته باشی و هرکدام یک
ThreadLocalسنگین بگیرد، مصرفِ حافظه منفجر میشود.InheritableThreadLocalهم بدتر است چون به هر child کپی میشود. - گمشدنِ context در پرشهای async: مقادیری مثل MDC (همان
traceIdدر لاگها) وSecurityContextرویThreadLocalسوارند. وقتی یکthenApplyAsyncبه نخِ دیگری میپرد، آنThreadLocalهمراه نمیآید — و یکدفعه لاگهایتtraceIdرا گم میکنند و امنیتت خالی میشود. همین مشکل سرِ مرزِexecutor.submit(...)هم هست.
ScopedValue (که در جاوا ۲۵ با JEP 506 نهایی شد) یک زمینهی تغییرناپذیر و محدود به یک scope میدهد: با ScopedValue.where(KEY, val).run(...) مقدار را فقط برای طولِ آن اجرا و برای همهی subtaskهای فرزند در StructuredTaskScope میبندی. برخلاف ThreadLocal، نه نشت میکند، نه باید دستی پاکش کنی، و برای دنیای میلیوننخی ساخته شده. برای propagate کردنِ MDC در فریمورکها هم Spring از TaskDecorator و کتابخانهی context-propagation میکرومتر استفاده میکند.
۷. Structured Concurrency — جانشینِ واقعیِ گرافهای شکننده
مشکلی که CompletableFuture بد حلش میکند این است: اگر یکی از N تا subtask شکست بخورد، بقیه بهصورت یتیم ادامه میدهند و منبع میسوزانند؛ هیچ کسی مالکِ چرخهی عمرِ همهشان نیست. Structured Concurrency دقیقاً همین را حل میکند: subtaskها فرزندِ یک scope میشوند، و مثل یک بلوکِ try با کروشه، یا همه با هم موفق میشوند یا با شکستِ یکی، بقیه کنسل و interrupt میشوند.
نمودارِ درختِ scope (هر subtask فرزندِ scope است و در همان بلوک join میشود):
flowchart TD
Scope[StructuredTaskScope] --> A[fork: userSvc.get]
Scope --> B[fork: orderSvc.recent]
Scope --> C[fork: riskSvc.score]
Scope --> J{join}
J -->|any fails| Cancel[cancel & interrupt siblings]
J -->|all ok| Combine[combine results]
// نسخهی جاوا ۲۵ (هنوز preview): factory استاتیک open() + یک Joiner برای سیاست
try (var scope = StructuredTaskScope.open(Joiner.<Object>allSuccessfulOrThrow())) {
var user = scope.fork(() -> userSvc.get(id)); // فرزندِ scope
var order = scope.fork(() -> orderSvc.recent(id));
scope.join(); // منتظرِ همه؛ اگر یکی بترکد، بقیه interrupt
return new Dashboard(user.get(), order.get());
} // بستنِ scope تضمین میکند هیچ subtaskِ یتیمی جا نمیماند
Structured Concurrency هنوز preview است و امضایش بین نسخهها عوض شده: در جاوا ۲۵ سازندههای عمومی با فکتوریِ استاتیکِ StructuredTaskScope.open(...) و اینترفیسِ Joiner (با سیاستهایی مثل allSuccessfulOrThrow, anySuccessfulResultOrThrow) جایگزین شدند. اگر روی جاوای قدیمیتری، هنوز الگوی قدیمیِ ShutdownOnFailure/ShutdownOnSuccess را میبینی — مفهوم یکی است، امضا فرق دارد.
۸. Retry با backoff — جایی که exceptionally کم میآورد
exceptionally فقط میتواند به یک مقدارِ ساده برگردد؛ نمیتواند یک تلاشِ دوبارهی async راه بیندازد. برای این، جاوا ۱۲ متدِ exceptionallyCompose را داد که خطا را به یک future دیگر میبندد. همراهش delayedExecutor (جاوا ۹) میآید که بدونِ Timer جداگانه، یک اجرا را عقب میاندازد:
CompletableFuture<Resp> withOneRetry() {
return callAsync()
.exceptionallyCompose(ex -> // خطا؟ یک future تازه بساز (نه یک مقدار)
CompletableFuture.supplyAsync(this::callAsyncOnce,
CompletableFuture.delayedExecutor(200, MILLISECONDS))); // بعد از ۲۰۰ms دوباره
}
اگر داخلِ خودِ اکشنِ whenComplete یک استثنا پرتاب شود و مرحلهی بالادست هم قبلاً شکست خورده بوده، استثنای اصلی گم میشود و استثنای جدیدِ تو پایین میرود. پس در whenComplete هرگز کاری که ممکن است بترکد نگذار، مگر خودت آن را try/catch کنی — وگرنه علتِ ریشهایِ باگت ناپدید میشود.
۹. تضمینِ حافظه (happens-before) و پرشِ بیفایده
دو نکتهی ریزِ سنیوری برای بستنِ بحث:
CompletableFutureیک رابطهی happens-before برقرار میکند: هر کاری که در نخِ کاملکننده قبل ازcomplete(v)انجام شده، برای مرحلههای وابسته قابلمشاهده است. یعنی برای خواندنِ مقدارِ تولیدشده نیازی بهvolatileیا قفلِ اضافی نداری. (اما این تضمین فقط شاملِ خودِ آن مقدار است، نه هر state جانبیِ mutable که بین مرحلهها به اشتراک گذاشتهای.)thenApplyAsync(fn, sameExecutor)باز هم میپرد: حتی اگر همین الان روی همان executor باشی، نسخهی*Asyncیک تسکِ جدید submit میکند — یعنی صفبندی و context switchِ اضافه. اگر به پرش نیاز نداری، combinatorِ غیر-async ارزانتر است. پرشِ عمدی فقط وقتی ارزش دارد که بخواهی از یک نخِ گران (مثلاً نخِ event-loop یا IO pool) فرار کنی.
سؤالات مصاحبهی سنیور (سطحِ سخت)
نه. پارامترِ mayInterruptIfRunning در CompletableFuture نادیده گرفته میشود؛ چون CF هیچ ارجاعی به آن نخ ندارد. cancel فقط خودِ future را استثنایی میکند تا مرحلههای پاییندست رها شوند، ولی کارِ زیرین تا آخر اجرا میشود. کنسلِ واقعی باید تعاونی باشد یا از کلاینتی بیاید که کنسل را propagate میکند (یا از StructuredTaskScope).
هیچ. future با TimeoutException بسته میشود اما کوئری/کانکشن همچنان زنده است و checked-out میماند. زیرِ بار، این میشود نشتیِ کانکشنپول: تایماوتها سریع رخ میدهند ولی کارِ واقعی تلنبار میشود تا pool خفه شود. تایماوتِ واقعی باید در خودِ کلاینت (JDBC/HTTP) تنظیم شود، نه فقط در لایهی CompletableFuture.
شکستِ سریع. anyOf با اولین چیزی که کامل شود کامل میشود — و کاملشدنِ استثنایی هم کاملشدن است. پس اگر سرورِ اول فوری خطا بدهد، کلِ نتیجه شکست میخورد حتی اگر سرورِ دوم داشت درست جواب میداد. ضمناً anyOf بازندهها را کنسل نمیکند، پس منبع میسوزانند. hedgingِ درست نیاز به منطقِ کنسلِ دستیِ رقبا دارد.
چون هر فراخوانی میتواند complete(...), cancel() یا obtrudeValue(...) صدا بزند و گرافِ پاییندستِ تو را بدزدد — مثلِ public کردنِ یک فیلدِ private. راهحل: minimalCompletionStage() (جاوا ۹) که یک CompletionStage بدونِ متدهای completion میدهد و صدا زدنشان UnsupportedOperationException میدهد؛ یا copy() برای یک کپیِ محافظتشده.
تا جاوا ۲۳، ورود به synchronized باعثِ pinning میشد: monitor روی carrierِ پلتفرمی ثبت میشد و virtual thread به carrier میخکوب میماند. اگر داخلِ آن بلاک بلاک میشدی (I/O)، carrier آزاد نمیشد؛ اگر همهی carrierها اینطور میشدند، موازیسازی به صفر میرسید و میشد deadlock. JEP 491 در جاوا ۲۴ monitor را به خودِ virtual thread گره زد؛ حالا synchronized پین نمیکند. فقط فریمهای native هنوز پین میکنند. راهحلِ پیش از ۲۴: ReentrantLock بهجای synchronized.
در allOf، بهمحضِ اولین شکست futureِ ترکیبی استثنایی میشود ولی بقیهی futureها یتیم و بیمراقب ادامه میدهند و منبع میسوزانند — چون هیچکس مالکِ چرخهی عمرِ آنها نیست. در StructuredTaskScope، subtaskها فرزندِ scopeاند؛ با شکستِ یکی، بقیه کنسل و interrupt میشوند و بستنِ scope تضمین میکند هیچ نخِ یتیمی جا نمیماند. این همان «ساختاریافته» بودن است.
از JDK 10 به بعد availableProcessors() نسبت به cgroup آگاه است و ۱ برمیگرداند؛ پس parallelismِ common pool میشود 1 - 1 = 0 و pool عملاً فرو میریزد: parallel stream روی نخِ صداکننده سریال اجرا میشود. آن speedupای که در لپتاپِ چندهستهای دیدی در prod ناپدید میشود. درمان: parallelism را صریح با -Djava.util.concurrent.ForkJoinPool.common.parallelism=N یا -XX:ActiveProcessorCount=N تنظیم کن، یا برای این stream یک ForkJoinPool اختصاصی بساز.
چون exceptionally فقط میتواند خطا را به یک مقدارِ ساده برگرداند؛ نمیتواند یک عملیاتِ asyncِ دوباره را راه بیندازد. exceptionallyCompose (جاوا ۱۲) خطا را به یک future دیگر میبندد، پس میتوانی داخلش یک supplyAsync جدید (تلاشِ دوباره) برگردانی. همراهش delayedExecutor(200, MILLISECONDS) را بگذار تا بدونِ Timer جداگانه، تلاشِ بعدی با تأخیر اجرا شود.
- کنسل و
orTimeoutکارِ زیرین را متوقف نمیکنند — منبع نشت میکند مگر کنسل را واقعاً propagate کنی. anyOf/allOfبرادرها را نمیکشند؛anyOfبا شکستِ سریع هم کامل میشود.- future را با
minimalCompletionStage()محافظتشده بیرون بده. - در کانتینر، سایزِ common pool را صریح بگذار؛ روی ۱ vCPU، parallel stream سریال میشود.
- pinning: پیش از جاوا ۲۴
synchronizedcarrier را میخکوب میکرد؛JEP 491رفعش کرد. virtual thread را pool نکن و برای CPU به کار نبر. - context (MDC/security) در پرشهای async گم میشود؛
ScopedValue(جاوا ۲۵) جانشینِ ساختاریافتهیThreadLocalاست. - Structured Concurrency جانشینِ واقعیِ گرافهای شکننده است: یا همه موفق، یا بقیه interrupt.
- دو ابزار، دو کار:
ForkJoinPoolموتورِ کارِ CPU-bound با سرقت کار است؛CompletableFutureیک گراف ترکیبِ ناهمگام از callbackهاست که خودش نخ نمیسازد. - common pool مالِ تو نیست: یک singletonِ سراسری با
هسته - 1worker که بین parallel stream و CF پیشفرض مشترک است. کارِ مسدودکننده روی آن = قحطی = ریشهی بیشترِ باگها. - قانونِ نخ: combinatorهای غیر-async روی نخِ تکمیلکننده اجرا میشوند؛
*Asyncبدون executor روی common pool. برای کارِ مسدودکننده همیشه executor اختصاصی بده. - map در برابر flatMap:
thenApplyبرای مقدارِ ساده،thenComposeبرای futureِ تودرتو. - خطا سُر میخورد پایین و در
CompletionExceptionپیچیده میشود —getCause()را صدا بزن.exceptionally/handleبازیابی میکنند،whenCompleteفقط ثبت میکند. - fork/join: یکی را fork کن، دیگری را inline محاسبه کن، بعد join؛
THRESHOLDرا تنظیم کن. - virtual threadها (جاوا ۲۱، JEP 444) مسدود شدن را ارزان میکنند و استخرِ FIFOِ جدای خودشان را دارند؛ برای fan-out مربوط به I/O ترجیحشان بده،
CompletableFutureرا برای گرافهای ترکیبِ واقعی نگه دار.
Let's be honest: concurrency in Java is where even experienced engineers stumble — not because it's hard, but because they mix up two completely different tools and then wonder why their production service locks up at midnight. This lesson exists to untangle exactly that. We'll build from zero, define every piece of jargon the moment you meet it, and by the end you should be able to answer any interview question in this area with confidence.
Here's our path: (1) build a mental model so you can tell ForkJoinPool apart from CompletableFuture; (2) learn CompletableFuture fully — from creation to composition to error handling; (3) master the golden rule of "which thread runs what" — the source of most production bugs here; (4) crack open ForkJoinPool and its work-stealing; (5) look inside parallelStream(); (6) compare against reactive; and (7) see how virtual threads in Java 21 change the whole calculus. Then a full interview-questions section.
Part 0 — Words you must know
Before anything, let's lock in a few words that keep recurring, each with a simple picture, so you never get lost later.
- Concurrency: managing several tasks so they appear to progress at once. Like a cook juggling several pots on the stove, rotating between them.
- Blocking: work that puts a thread to sleep waiting for something — a database or network reply, say. A blocked thread does nothing else; it just waits.
- CPU-bound: work that actually keeps the processor busy, like heavy math. Here the thread isn't sleeping — it's sweating.
- Thread: an independent line of execution. Picture a worker running down a list of instructions one by one.
- Executor / thread pool: a pre-built box of workers. Instead of hiring a new worker each time (expensive), you borrow one from the pool, hand it work, and return it.
Two kinds of work, two kinds of tools. CPU-bound work (computation) goes on ForkJoin; blocking work (I/O like network and database) goes on a dedicated executor or a virtual thread. Mixing the two is the root of nearly every disaster.
Mental model — the two tools everyone confuses
Java gives you two concurrency toolkits that look similar but are built for two very different jobs.
Imagine you run a restaurant. ForkJoinPool is like the kitchen: a bunch of professional cooks whose job is chopping, frying, cooking — short, fast tasks with no waiting. No cook stops mid-task to sit by the phone waiting for a delivery; everyone's hands are always moving.
But CompletableFuture isn't a cook; it's the recipe card: "First make the salad, then plate it once the steak is done, and if there's no fish, substitute it." The card itself does no work and hires no cooks — it only says what follows what and where each step should run. And by default, this card hands its steps to that same ForkJoin kitchen.
More precisely:
ForkJoinPoolis a CPU-bound work engine. Its whole design — double-ended queues (deques) with work-stealing, and LIFO local scheduling — assumes tasks are short, non-blocking, and recursively divisible. It's the backbone ofparallelStream()and the default executor forCompletableFuture.CompletableFutureis an asynchronous composition graph; a promise/future with combinators. It does not create threads by itself and does not manage blocking; it only decides where your callbacks run — and by default that "where" is the ForkJoin common pool.
CompletableFuture is a dependency graph of callbacks, and ForkJoin's common pool is the default place those callbacks execute. Almost every production bug in this area comes from running blocking work on that non-blocking pool.
Whenever you get confused below, come back to this sentence. The whole lesson is just unpacking this one line.
Part 1 — CompletableFuture
Creating futures
Before you can wire work together, you need a future. There are several ways to make one:
// Runs the supplier on the ForkJoin common pool, returns immediately
CompletableFuture<String> f1 = CompletableFuture.supplyAsync(() -> fetchUser());
// Runs on a supplied executor instead
ExecutorService io = Executors.newFixedThreadPool(16);
CompletableFuture<String> f2 = CompletableFuture.supplyAsync(() -> fetchUser(), io);
// Already-completed values (no async work at all)
CompletableFuture<Integer> done = CompletableFuture.completedFuture(42);
// A future you complete manually — the "promise" pattern
CompletableFuture<String> promise = new CompletableFuture<>();
// ... later, from any thread:
promise.complete("value"); // or promise.completeExceptionally(ex);
Let's unpack these:
supplyAsyncmeans "go run this work somewhere else and hand me a receipt I can redeem later." That "somewhere else" is the common pool by default, unless you pass your own executor like inf2.completedFuturemeans "the answer is ready right now, there's nothing to do." Like a receipt that already has the answer written on it.- The promise pattern: you create an empty
CompletableFutureand later fill it from any thread withcomplete(...). Handy when you're waiting on a callback from some older library.
runAsync is the void-returning sibling of supplyAsync (takes a Runnable, yields CompletableFuture<Void>) — for when you just want work done and don't need a result.
supplyAsync is eager: the work starts the moment you call it, not when you ask for the answer. That's different from some reactive libraries which are "lazy" — and we'll come back to this later, because it's a common interview question.
The three families of combinators
Now that you have a future, you want to do things to it: transform it, chain the next step, or combine it with another future. These methods are called combinators. Think along three axes:
| Purpose | Method | Function type |
|---|---|---|
| Transform the value | thenApply |
T -> U |
| Chain another async stage (flatten) | thenCompose |
T -> CompletableFuture<U> |
| Consume the value, no result | thenAccept / thenRun |
T -> void / () -> void |
| Combine two independent futures | thenCombine |
(T, U) -> V |
| Wait for many | allOf / anyOf |
— |
The most important pair here is thenApply vs thenCompose.
thenApply is like reaching into a gift box, taking the item, painting it, and putting it back. One box, one item. But sometimes the work you do itself returns another box (because that work is also async). If you use thenApply for that, you end up with a box inside a box — CompletableFuture<CompletableFuture<U>>. thenCompose peels off that extra wrapper and hands you just one box.
In functional terms: thenApply is exactly map, and thenCompose is exactly flatMap. Simple rule: if your function returns a plain value, use thenApply; if it returns another future, use thenCompose.
CompletableFuture<Order> pipeline =
CompletableFuture.supplyAsync(() -> loadCart(userId)) // CF<Cart>
.thenApply(cart -> priceCart(cart)) // CF<PricedCart> (sync transform)
.thenCompose(priced -> chargeAsync(priced)) // CF<Payment> (returns a future -> flatten)
.thenApply(payment -> new Order(payment)); // CF<Order>
Notice priceCart returns a plain value so it's thenApply; but chargeAsync itself returns a CompletableFuture, so it's thenCompose to avoid the nested box.
thenCombine and allOf / anyOf
So far we've had a linear chain. But sometimes you want two independent pieces of work to run concurrently and then join their answers. That's where thenCombine comes in:
CompletableFuture<Integer> price = CompletableFuture.supplyAsync(this::fetchPrice);
CompletableFuture<Integer> stock = CompletableFuture.supplyAsync(this::fetchStock);
// Both run concurrently; combine when BOTH are done
CompletableFuture<Quote> quote =
price.thenCombine(stock, (p, s) -> new Quote(p, s));
The key win: because price and stock are both created with supplyAsync, they start concurrently. thenCombine just waits until both finish, then calls the function. So total time is roughly the slower of the two, not their sum.
When you have not two but a whole list of futures, use allOf:
List<CompletableFuture<Product>> futures = ids.stream()
.map(id -> CompletableFuture.supplyAsync(() -> fetch(id), io))
.toList();
CompletableFuture<List<Product>> all =
CompletableFuture.allOf(futures.toArray(CompletableFuture[]::new))
.thenApply(v -> futures.stream()
.map(CompletableFuture::join) // safe: all are already complete here
.toList());
allOf returns a CompletableFuture<Void> — meaning it only tells you "everyone's done" but throws away the values themselves. So after it completes, you must re-join() each future individually to get its answer. Here join() is safe because we know all are already complete, so it never actually blocks.
anyOf is the opposite of allOf: it completes with the result of the first future to finish. It's useful for hedged requests (fire the same request at two servers, take whichever answers first) or timeouts. It has one wart: because it can't infer a common type across differing futures, it returns CompletableFuture<Object>, so you cast.
Error propagation: exceptionally, handle, whenComplete
This is where senior candidates separate themselves. First, a key idea that trips many people up:
When an exception occurs inside a stage (say a thenApply), unlike ordinary code, it is not thrown where you wrote the lambda. Instead, it completes that stage exceptionally, and the exception slides down the chain until it reaches a handler — wrapped, along the way, in a CompletionException.
CompletableFuture<String> result =
CompletableFuture.supplyAsync(() -> { throw new IllegalStateException("boom"); })
.thenApply(s -> s.toUpperCase()) // SKIPPED — upstream failed
.exceptionally(ex -> "fallback"); // recovers: ex is CompletionException wrapping ISE
See how thenApply never ran because the stage before it failed; the exception jumped right over it straight to exceptionally. We have three tools for error handling, and their differences are exactly what interviewers ask about:
| Method | Sees value | Sees exception | Can transform value | Can recover from error | Returns |
|---|---|---|---|---|---|
exceptionally |
no | yes | no | yes | CompletableFuture<T> |
handle |
yes | yes | yes | yes | CompletableFuture<U> (new type) |
whenComplete |
yes | yes | no (side-effect only) | no (re-throws original) | CompletableFuture<T> |
Imagine a product coming down a conveyor belt, possibly good or defective. exceptionally is the worker who springs into action only when the product is defective and drops in a replacement — if it's good, they don't even wake up. handle is the inspector who sees every product (good or bad), can swap out the good one, repair the bad one, and even change the product type. whenComplete is the security camera: it sees and records everything (logging, cleanup) but changes nothing — a defective product goes down the line just as defective.
Key traps to memorize:
- The exception you catch is usually a
CompletionExceptionwrapping the real cause. Always callex.getCause()to inspect the real one. whenCompletedoes not swallow the exception — the returned future still fails. It's for logging/cleanup, not recovery.handleruns on both paths; it's the only combinator that lets you turn a failure into a success and change the type.- Each combinator has an
*Asyncvariant (exceptionallyAsync,handleAsync) if you need the recovery to hop to another thread.
CompletableFuture<Response> safe =
callServiceAsync()
.orTimeout(2, TimeUnit.SECONDS) // Java 9+: fail with TimeoutException after 2s
.handle((resp, ex) -> {
if (ex != null) {
log.warn("degraded", ex.getCause()); // unwrap CompletionException
return Response.degraded();
}
return resp;
});
orTimeout(2, SECONDS) says "if no answer within 2 seconds, fail with TimeoutException." And completeOnTimeout(defaultValue, ...) says "if it's late, instead of failing, drop in this default value." Both were added in Java 9, and their beauty is they bound a future without a separate scheduler or timer thread.
Which executor runs what — the rule that trips everyone
Now we reach the heart of it all. If you memorize only one part of this lesson, this is it. Two simple rules behind a huge fraction of Java concurrency bugs:
Rule 1 — non-async combinators (thenApply, thenCombine, handle, …) run on whichever thread completed the upstream future. If the upstream is already complete when you attach the callback, it runs on your calling thread, synchronously.
Rule 2 — *Async combinators with no executor run on the ForkJoin common pool. With an executor argument, they run there.
CompletableFuture.supplyAsync(() -> load(), io) // runs on 'io'
.thenApply(x -> step1(x)) // runs on 'io' (the completing thread)
.thenApplyAsync(x -> step2(x)) // hops to ForkJoinPool.commonPool()
.thenApplyAsync(x -> step3(x), io); // runs on 'io'
Read it line by line: load() runs on the io pool. Now step1 is a non-async thenApply, so by Rule 1 it runs on that same io thread that finished the previous work. Then step2 is a thenApplyAsync with no executor, so by Rule 2 it hops to the common pool. And step3 is a thenApplyAsync with the io executor, so it comes back to io.
Imagine a file finishes on the desk of a clerk in the IO department. Now the next step is a "non-async" task. Instead of the file going into the central queue, that same IO clerk immediately does the next step too. If the next step is short, no problem; but if it's a long chain of such tasks, that poor IO clerk gets so buried in other people's work that there's no time left for their actual job (answering the network) — and your IO pool gets starved.
Consequences that bite in production:
- If
supplyAsync(..., io)completes, thenstep1also occupies aniothread — a longthenApplychain can starve your IO pool even though only one stage was "async." - If the future is already done (e.g.
completedFuture), yourthenApplyruns inline on the caller — sometimes surprising in tests, because you think work goes to the background but it runs right where you are. - Because the default is the common pool, any blocking call inside a default combinator blocks a common-pool worker, and there are only
Runtime.availableProcessors() - 1of them by default. Block enough and the whole JVM's parallel machinery (streams, other futures) stalls. Always pass an explicit executor for blocking work.
A realistic composition example
Let's put it all together: call three services concurrently, combine their answers, each with its own timeout, and if anything fails, degrade to a fallback.
// Fan out to 3 services, combine, with per-call timeout and a global fallback.
CompletableFuture<Profile> profile = supplyAsync(() -> userSvc.get(id), io)
.orTimeout(500, MILLISECONDS);
CompletableFuture<List<Order>> orders = supplyAsync(() -> orderSvc.recent(id), io)
.exceptionally(ex -> List.of()); // orders are optional
CompletableFuture<Score> score = supplyAsync(() -> riskSvc.score(id), io)
.orTimeout(300, MILLISECONDS);
CompletableFuture<Dashboard> dashboard =
profile.thenCombine(orders, Dashboard::new) // (Profile, List<Order>) -> Dashboard
.thenCombine(score, Dashboard::withScore) // add third result
.exceptionally(ex -> Dashboard.EMPTY); // if profile/score failed, degrade
Dashboard d = dashboard.join(); // block the caller once, at the edge
Notice the philosophy: each service runs on io (a dedicated pool, not the common pool), each has a timeout, orders — which is "optional" — falls back to an empty list if it fails, and at the very end we block only once, at the edge of the graph with join(). The interior of the graph stays entirely non-blocking.
join() throws an unchecked CompletionException; get() throws the checked InterruptedException and ExecutionException. Since you can't throw checked exceptions inside stream lambdas, prefer join() there — it's cleaner.
Part 2 — ForkJoinPool & work-stealing
Now to that kitchen. ForkJoinPool is built for one specific kind of parallelism: divide-and-conquer.
The model
The idea is simple: split a big task into two halves (fork), solve each half separately, then merge the answers (join). And if a half is still big, split it again — until we reach small, directly-solvable chunks.
The technical heart of ForkJoin is how work is distributed among threads. Each worker thread owns a double-ended queue (deque) — a queue you can take work from at both the head and the tail.
Each cook has a stack of tasks on their own desk. When they create a task, they put it on top of the ones they just placed and pick from that same spot (the head) — because fresh work is "hot," its ingredients and tools are still at hand. That's LIFO (last in, first out) and it keeps the CPU cache warm.
But when another cook goes idle, they walk over to a busy coworker's desk and take a task from the bottom (the tail) — that is, the oldest task, usually not yet split and the largest chunk, one the coworker won't need urgently anyway. That's work-stealing, done FIFO.
So the owner pushes/pops from the head of its deque LIFO (good cache locality, recent tasks are hot), and idle workers steal from the tail of a victim's deque FIFO (oldest, largest, least-likely-to-be-needed-soon). This LIFO-local + FIFO-steal asymmetry keeps all cores busy with minimal central contention — it's the key idea of the whole design.
RecursiveTask example
Let's see this in code. We want to compute the sum of a large array in parallel. We extend RecursiveTask<V> because it returns a value (Long):
class SumTask extends RecursiveTask<Long> {
private static final int THRESHOLD = 10_000;
private final long[] arr; private final int lo, hi;
SumTask(long[] arr, int lo, int hi) { this.arr = arr; this.lo = lo; this.hi = hi; }
@Override protected Long compute() {
if (hi - lo <= THRESHOLD) { // base case: compute directly
long s = 0;
for (int i = lo; i < hi; i++) s += arr[i];
return s;
}
int mid = (lo + hi) >>> 1;
SumTask left = new SumTask(arr, lo, mid);
SumTask right = new SumTask(arr, mid, hi);
left.fork(); // schedule left asynchronously
long rightResult = right.compute(); // compute right in THIS thread (don't fork both!)
long leftResult = left.join(); // wait for the stolen/queued left
return leftResult + rightResult;
}
}
long total = ForkJoinPool.commonPool().invoke(new SumTask(data, 0, data.length));
Fork one branch, compute the other inline, then join. Notice we left.fork() but call right.compute() right here in the current thread, and finish with left.join(). If we forked both and then joined both, we'd waste a queue slot and force a steal for nothing; computing one inline keeps the current thread productive while the other branch stands ready to be stolen if needed.
RecursiveAction is exactly the same but for when you return nothing (the Void variant).
Choosing THRESHOLD is an art: too big and you don't get enough parallelism and cores sit idle; too small and the fork/join overhead itself exceeds the real work. Set it so each leaf is big enough to dwarf the overhead but small enough to load-balance well.
The common pool and its parallelism
Now a very important point: ForkJoinPool.commonPool() is a JVM-wide singleton. That is, one shared pool for your entire application.
Its parallelism defaults to Runtime.getRuntime().availableProcessors() - 1. So on an 8-core box you get 7 workers (why minus one? because the submitting thread itself can also pitch in and help). It's tunable:
-Djava.util.concurrent.ForkJoinPool.common.parallelism=32
Because it's shared by parallelStream(), default CompletableFuture, and any library that grabs it, it's a global resource you don't own. Sizing it for CPU work (≈ core count) and then doing blocking I/O on it is the classic starvation bug that takes down the whole app.
Blocking correctly: ManagedBlocker
Suppose you truly have no choice but to block inside a ForkJoin task. The sanctioned answer is to tell the pool, so it can spin up a compensation thread to keep parallelism up — that is, in place of the thread that went to sleep, it wakes another so cores don't sit idle:
ForkJoinPool.managedBlock(new ForkJoinPool.ManagedBlocker() {
private String result;
@Override public boolean block() throws InterruptedException {
result = blockingRemoteCall(); // the actual blocking call
return true; // true = done, no need to block again
}
@Override public boolean isReleasable() {
return result != null;
}
});
Imagine one of your cooks has to go to the phone and wait. If you don't say anything, a core effectively sits idle. ManagedBlocker is like ringing a bell to say "I'm on the phone, bring a temp cook." The pool adds a compensation thread so the effective parallelism doesn't drop. When you return, things go back to normal.
parallelStream() on synchronized/Collections-heavy paths and JDBC calls are exactly the cases where naive blocking silently degrades throughput. ManagedBlocker is the sanctioned escape hatch — but honestly, on modern JDKs the cleaner answer is a dedicated executor or virtual threads.
Part 3 — Parallel streams internals
Now that you know ForkJoin, parallelStream() is no longer magic. When you call collection.parallelStream() (or stream().parallel()), the pipeline runs on that same common pool, splitting the source via an object called a Spliterator.
Its name comes from Split + Iterator. An Iterator hands you elements one by one; a Spliterator can additionally split itself in half and say "this half of the elements is yours, that half is mine." How well and cheaply a source splits directly determines whether a parallel stream on it is fast or slow.
Rules of thumb that follow directly from these internals:
- Parallelism helps only when: the source splits evenly and cheaply (arrays,
ArrayList,IntStream.range— great, because you can find the midpoint with one calculation; butLinkedListandIterator-backed sources — poor, because to reach the middle you must walk from the start), the per-element work is CPU-heavy, and N is large. A rough heuristic: N × cost-per-element should be well above ~100μs for the parallelism overhead to be worth it. - For reductions, the pipeline must be stateless, non-interfering, and associative. That is: your function must not keep external state or mutate the source mid-flight, and the grouping order must not change the answer.
reducemust be associative;collectuses acombinerto stitch partial results together. - Never do I/O in a parallel stream — you'll block common-pool workers and starve everything else. (The chapter's key line, once more.)
forEachguarantees no ordering; if you need encounter order, useforEachOrdered— which costs, because it has to coordinate.
It means (a op b) op c must give the same answer as a op (b op c). Addition and multiplication are associative, but subtraction isn't ((5-3)-2 ≠ 5-(3-2)). Because a parallel stream groups and merges chunks in an unspecified order, if your operation isn't associative, the parallel answer differs from the sequential one — and this is one of the worst bugs because it's silent.
// CPU-bound, splittable source, associative reduction — a good fit
long primes = LongStream.rangeClosed(2, 10_000_000)
.parallel()
.filter(this::isPrime)
.count();
This is a great fit: the source (rangeClosed) splits perfectly cheaply, the per-element work (isPrime) is genuinely CPU-bound, and N is large.
And a known trick: if you want to escape the common pool for a specific parallel stream (say, so it doesn't affect the rest of the app), submit it into your own ForkJoinPool:
ForkJoinPool custom = new ForkJoinPool(4);
long r = custom.submit(() ->
list.parallelStream().filter(this::heavy).count()
).join(); // this parallel stream now runs on 'custom', not the common pool
Because a parallel stream binds to the same thread that runs it, and while you're inside custom.submit(...), that thread is one of custom's workers. So the splits happen on custom, not the common pool. It's not officially documented but it's reliable behavior.
Part 4 — Futures vs reactive (brief)
Interviewers sometimes ask "what's the difference between CompletableFuture and reactive?" The answer is short but important.
CompletableFuture models a single async value with composition — "one answer, coming later." But Reactive Streams (Mono/Flux in Project Reactor, or RxJava) model streams of 0..N values with backpressure, and have a far richer operator set (retry, buffer, window, merge, flow-control).
CompletableFuture is like waiting for one parcel: it arrives or it doesn't, done. Flux is like standing by a conveyor belt that might bring zero, ten, or ten thousand parcels, and you can say "slow down, my hands are full" — that's backpressure, the consumer pushing back on the producer to ease off.
The choice rule: use CF for "call these services, combine, return once"; use reactive when you have streams, backpressure, or need composition over time. Reactor's Mono is roughly "CF with hundreds of operators and lazy subscription."
CompletableFuture is eager — work starts at creation. But Mono/Flux is lazy — nothing runs until you subscribe(). This directly affects retries, caching, and whether "creating a future without consuming it" does any work at all.
Part 5 — Virtual threads change the calculus (Java 21)
Now the biggest shift. Until now, the whole philosophy of CompletableFuture was "OS threads are expensive, so don't waste them on blocking — write callbacks instead." But virtual threads (JEP 444, stable in Java 21) turn that assumption upside down: they make blocking cheap again.
An OS thread is like hiring an expensive full-time employee; you can't have thousands, so you mustn't let one sit idle by the phone. But a virtual thread is like an ultra-cheap worker you can have millions of. Now it no longer matters if one waits by the phone — the JVM lifts it off the core, gives the core to another, and puts it back when the answer arrives. So you no longer need complex callback code.
With virtual threads you can write simple, straight-line blocking code and run thousands of them concurrently:
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
var f1 = executor.submit(() -> userSvc.get(id)); // blocking call, but on a vthread
var f2 = executor.submit(() -> orderSvc.recent(id));
return new Dashboard(f1.get(), f2.get()); // structured, readable, no chains
}
See how readable that is — no thenCombine, no chains, just "do these two things concurrently and get their answers."
Virtual threads are scheduled by their own ForkJoinPool in FIFO mode (one carrier thread per core by default, tunable via jdk.virtualThreadScheduler.parallelism) — that pool is separate from the common pool. So blocking on a virtual thread doesn't touch the common pool that belongs to parallel streams and default CF.
So these aren't competitors — they're complementary:
- Use
CompletableFuturefor its composition operators, when you genuinely need a graph (thenCombine,allOf, timeouts). - Reach for virtual threads (or Structured Concurrency via
StructuredTaskScope) when the logic is really "do these in parallel and join." - Keep ForkJoin for pure CPU-bound parallelism.
Common pitfalls & gotchas
Let's gather all the traps we met along the way in one place — these are the ones that bite in production:
- Blocking on the common pool. Default
supplyAsync/parallelStream+ JDBC/HTTP = starvation. Pass a dedicatedExecutor. CompletionExceptionwrapping. Yourcatch/exceptionallysees the wrapper; callgetCause().whenCompleteis not recovery. It re-throws; onlyhandle/exceptionallyrecover.allOfloses values. Re-join each future afterallOfcompletes.- Non-async combinators steal the completing thread. A long
thenApplychain runs on whatever thread finished the previous stage — possibly your IO pool, possibly a common-pool worker. - Forgetting the future is eager.
supplyAsyncstarts now; there's no lazy subscribe. - Forking both branches. Compute one inline; fork the other.
- Unbounded fan-out. Mapping 100k items to
supplyAsync(..., commonPool)overwhelms it — bound concurrency with a sized executor or a semaphore.
Best practices
- Own your executors for anything that blocks. Never trust the common pool for I/O.
- Name your threads (
ThreadFactory) so stack dumps are debuggable. - Prefer
join()at the edge of the graph; keep the interior non-blocking. - Add
orTimeout/completeOnTimeoutto every external call. - Handle errors with
handle/exceptionally; alwaysgetCause(). - Use
RecursiveTaskonly for CPU-bound divide-and-conquer with a tuned threshold. - On Java 21+, prefer virtual threads / structured concurrency for I/O fan-out; reserve
CompletableFuturefor genuine composition graphs; reserve ForkJoin for CPU-bound parallelism.
Interview Questions
Now the part you came for. We'll answer each so it's both correct and shows depth.
The thread that completed the upstream future. If the upstream was already complete when you attached the callback, it runs synchronously on the calling thread. Only *Async variants guarantee a hop to an executor (the common pool if none given).
thenApply is map (T -> U); thenCompose is flatMap (T -> CompletableFuture<U>). Use thenCompose when your function is itself async and returns a future, to avoid a nested CompletableFuture<CompletableFuture<U>>.
exceptionally runs only on failure and can recover (same type). handle runs on both paths, can transform the value and recover, and can change the result type. whenComplete runs on both but is side-effect only — it cannot recover and re-throws the original exception downstream.
Because join() throws a CompletionException wrapping the real cause. You must catch CompletionException and inspect getCause(), or the original type won't match.
Parallel streams run on the JVM-wide common pool (~`cores - 1workers). A blocking JDBC call parks a worker; enough of them starve the pool, stalling *every* parallel stream and defaultCompletableFuturein the process. Fix: dedicated executor,ManagedBlocker`, or virtual threads.
Each worker has a deque; it pushes/pops its own tasks LIFO from the head (cache-hot). Idle workers steal FIFO from the tail of another worker's deque (oldest/largest task). This balances load with minimal contention.
Forking both then joining both wastes a queue slot and forces a steal; computing one branch inline keeps the current thread doing useful work while the other branch is available to be stolen. It's the canonical, most efficient form.
availableProcessors() - 1. Change via -Djava.util.concurrent.ForkJoinPool.common.parallelism=N. Note it's a global singleton shared across the JVM.
var cf = CompletableFuture.supplyAsync(() -> 1)
.thenApply(x -> x / 0)
.thenApply(x -> x + 1)
.exceptionally(ex -> -1);
System.out.println(cf.join());
Answer: -1. The / 0 throws ArithmeticException, completing that stage exceptionally; the next thenApply is skipped; exceptionally recovers to -1.
List<CompletableFuture<Integer>> fs = ...;
CompletableFuture.allOf(fs.toArray(new CompletableFuture[0]));
List<Integer> results = fs.stream().map(CompletableFuture::join).toList();
Bug: allOf(...) is created but never joined/awaited, so its result is dropped and there's no guarantee stages scheduled after aren't racing. More importantly, if any future fails, join() throws mid-stream and you lose the others. Fix: chain .thenApply(...) off the allOf future, or collect with per-future exception handling.
Eager — supplyAsync starts executing immediately. Mono/Flux are lazy: nothing runs until subscribe(). This affects retries, caching, and whether creating-but-not-consuming does work.
No. Virtual threads run on a separate ForkJoinPool in FIFO mode, one carrier thread per core by default (jdk.virtualThreadScheduler.parallelism). The common pool (LIFO) is used by parallel streams and default CompletableFuture.
All 100k callbacks land on the common pool; if any of them block, or if the volume exceeds cores - 1 productive slots, you get queueing/starvation. Pass a purpose-sized executor, or bound the fan-out.
When you need its composition semantics: combining heterogeneous results (thenCombine), waiting on many (allOf/anyOf), hedging (anyOf), declarative timeouts (orTimeout), and building reusable async graphs. For plain "run N blocking calls and join," virtual threads / StructuredTaskScope are simpler.
Because non-async combinators execute on the completing thread. If each stage completes on the same executor's thread (e.g. supplyAsync(..., io) then a non-async thenCompose whose inner future also completes on io), the whole chain can serialize onto io workers — a subtle source of pool starvation. Insert *Async(..., otherPool) to hop deliberately.
Senior notes & advanced edge cases
So far this chapter taught you the mechanics of CompletableFuture, ForkJoin, and virtual threads. But what separates a senior from a strong mid-level engineer isn't knowing these APIs — it's knowing where they quietly lie to you: the things you think they do but don't, and which are exactly what light up your pager at 3 a.m. This section is that layer.
(1) why cancelling a CompletableFuture is a lie and orTimeout leaks; (2) how anyOf/allOf never cancel their siblings; (3) how to protect a future from outside tampering; (4) the common pool's behavior inside containers and Kubernetes; (5) virtual-thread pinning and its Java 24 fix; (6) ThreadLocal, MDC and ScopedValue; (7) Structured Concurrency as the real successor; (8) retry with backoff and exceptionallyCompose; (9) memory guarantees (happens-before). Then several hard interview questions.
1. Cancelling a future is a lie — and orTimeout leaks
This one the Javadoc barely admits, and it draws blood in production. The common belief is: "if I'm done waiting on a future, I'll cancel(true) it and the work behind it stops." That's wrong.
CompletableFuture<String> f = CompletableFuture.supplyAsync(() -> slowDbCall(), io);
boolean cancelled = f.cancel(true); // true = mayInterruptIfRunning
The mayInterruptIfRunning parameter is completely ignored by CompletableFuture. Unlike FutureTask, a CompletableFuture holds no reference to the thread running the supplier, so it cannot interrupt it. All cancel does is transition the future itself to an exceptional (CancellationException) state so that downstream stages of the graph are abandoned. But slowDbCall() runs to completion, and the DB connection stays checked out until it returns.
The same logic applies to orTimeout:
CompletableFuture<Response> r = callServiceAsync().orTimeout(2, SECONDS);
orTimeout(2, SECONDS) completes your future with a TimeoutException after 2s, but the HTTP/JDBC request behind it is still alive, still holding its resource (connection, thread). Under load this becomes a connection-pool leak: timeouts fire quickly, but the real work piles up until the pool is exhausted. The real timeout must live in the client (e.g. HttpClient with .timeout(...)), not in the CompletableFuture layer.
Senior takeaway: cancellation in the CompletableFuture world is not cooperative unless you wire it yourself. Either use a client that actually propagates cancellation (like HttpClient.sendAsync, whose returned future is wired to real cancellation), or use StructuredTaskScope, which genuinely interrupts subtasks when the scope fails.
2. anyOf and allOf never kill the siblings
The chapter said allOf throws away values and anyOf returns the first answer. Two crucial subtleties were left out:
anyOfcompletes with "the first thing to complete" — even if that completion is a failure. A fast failure beats a slow success. If you're doing hedged requests withanyOf(send the same request to two servers, take whichever answers first), an immediate error from the first server fails the whole result while the second was about to answer correctly.- Neither one cancels the losing siblings. After
anyOfreturns the winner, all the other futures still run to completion and burn resources. LikewiseallOfcompletes exceptionally on the first failure, but the remaining futures keep running unmonitored.
If you truly want hedging, after anyOf yields a winner you must manually apply cancel/resource-release logic to the losers, or use futures whose cancellation actually frees the resource. anyOf alone is just "first signal," not "lifecycle management of the competitors."
3. Don't hand out a raw future — protect it
When your library method returns a CompletableFuture<T>, any caller can call complete(...), obtrudeValue(...), or cancel() on it and hijack your entire downstream graph — like exposing a private field as public.
// Bad: anyone can forge your result
public CompletableFuture<Price> quote() { return internalFuture; }
// Good: a view whose completion methods throw
public CompletionStage<Price> quote() { return internalFuture.minimalCompletionStage(); }
minimalCompletionStage() (Java 9) returns a CompletionStage that only exposes combinators; calling complete/cancel/get on it throws UnsupportedOperationException. copy() gives a dependent copy whose completion doesn't affect the original. And use obtrudeValue/obtrudeException only in tests — they forcibly rewrite an already-completed future and shatter every invariant.
4. The common pool inside containers — the Kubernetes trap
The chapter said the common pool defaults to availableProcessors() - 1 threads. In the container world that number bites:
- Since JDK 10 (and 8u191),
availableProcessors()is cgroup-aware and respects the CPU limit. - If your pod has a 1-vCPU limit,
cores - 1becomes zero, and the common pool effectively collapses to no parallelism:parallelStream()runs serially on the calling thread. That speedup you saw on your 8-core laptop vanishes in prod. - Conversely, if you set no limit, the JVM may see all the host node's cores and build a pool far larger than your real share.
Don't trust auto-sizing in a container. Either set -XX:ActiveProcessorCount=N or directly -Djava.util.concurrent.ForkJoinPool.common.parallelism=N. Also remember common-pool threads are daemon threads; on JVM exit their unfinished work is abandoned with no guarantee — so never hand critical "must-finish" work to the common pool.
5. Virtual-thread pinning — the modern trap and its Java 24 fix
The chapter said virtual threads make blocking cheap, but there's a big asterisk that bled through Java 23: pinning.
When a virtual thread entered a synchronized block, in versions up to Java 23 that monitor lock was registered on the platform carrier thread, not on the virtual thread itself. Result: while inside that block the virtual thread was pinned to its carrier, and if you blocked there (e.g. I/O inside synchronized), the carrier could not be released to serve another virtual thread. If all carriers got pinned this way, the whole scheduler choked — and in the worst case, deadlock.
JEP 491 in Java 24 made the monitor tie to the virtual thread itself; now synchronized no longer pins the carrier and you can use it freely again. But two things still pin: native (JNI) frames, and blocking inside a native method. So pinning isn't dead, just very rare now.
On I/O hot paths, replace synchronized with ReentrantLock — java.util.concurrent locks never pinned. To debug, enable -Djdk.tracePinnedThreads (deprecated in 24) or the JFR jdk.VirtualThreadPinned event. And two evergreen rules: don't pool virtual threads (they're cheap by design), and don't use virtual threads for CPU-bound work — they're for blocking, not heavy computation.
6. ThreadLocal, MDC and ScopedValue
With virtual threads, two context problems get worse:
- Memory: with millions of virtual threads each holding a heavy
ThreadLocal, memory usage explodes.InheritableThreadLocalis worse still, because it's copied into every child. - Context loss across async hops: values like MDC (your log
traceId) andSecurityContextride onThreadLocal. When athenApplyAsynchops to a different thread, thatThreadLocaldoesn't come along — suddenly your logs lose thetraceIdand your security context is empty. The same problem occurs across anexecutor.submit(...)boundary.
ScopedValue (finalized in Java 25 via JEP 506) gives an immutable, scope-bounded context: with ScopedValue.where(KEY, val).run(...) you bind the value only for the duration of that run and for all child subtasks in a StructuredTaskScope. Unlike ThreadLocal, it doesn't leak, needs no manual cleanup, and is built for the million-thread world. For propagating MDC in frameworks, Spring uses a TaskDecorator and Micrometer's context-propagation library.
7. Structured Concurrency — the real successor to fragile graphs
The problem CompletableFuture handles badly: if one of N subtasks fails, the others keep running as orphans, burning resources, because nobody owns their collective lifecycle. Structured Concurrency solves exactly this: subtasks become children of a scope, and like a braced try-block, either they all succeed or, when one fails, the rest are cancelled and interrupted.
The scope tree (each subtask is a child of the scope, joined in the same block):
flowchart TD
Scope[StructuredTaskScope] --> A[fork: userSvc.get]
Scope --> B[fork: orderSvc.recent]
Scope --> C[fork: riskSvc.score]
Scope --> J{join}
J -->|any fails| Cancel[cancel & interrupt siblings]
J -->|all ok| Combine[combine results]
// Java 25 shape (still preview): static open() factory + a Joiner for the policy
try (var scope = StructuredTaskScope.open(Joiner.<Object>allSuccessfulOrThrow())) {
var user = scope.fork(() -> userSvc.get(id)); // child of the scope
var order = scope.fork(() -> orderSvc.recent(id));
scope.join(); // wait for all; if one blows up, the rest are interrupted
return new Dashboard(user.get(), order.get());
} // closing the scope guarantees no orphaned subtask is left behind
Structured Concurrency is still preview, and its signature changed between releases: in Java 25 the public constructors were replaced by the static factory StructuredTaskScope.open(...) plus a Joiner interface (with policies like allSuccessfulOrThrow, anySuccessfulResultOrThrow). On an older JDK you'll still see the old ShutdownOnFailure/ShutdownOnSuccess pattern — same concept, different signature.
8. Retry with backoff — where exceptionally falls short
exceptionally can only recover to a plain value; it cannot kick off an async retry. For that, Java 12 added exceptionallyCompose, which maps the failure to another future. It pairs with delayedExecutor (Java 9), which defers a run without a separate Timer:
CompletableFuture<Resp> withOneRetry() {
return callAsync()
.exceptionallyCompose(ex -> // failure? build a fresh future (not a value)
CompletableFuture.supplyAsync(this::callAsyncOnce,
CompletableFuture.delayedExecutor(200, MILLISECONDS))); // retry after 200ms
}
If the whenComplete action itself throws, and the upstream stage had already failed, the original exception is lost and your new exception flows downstream. So never put fallible work in whenComplete unless you try/catch it yourself — otherwise the root cause of your bug disappears.
9. Memory guarantee (happens-before) and the pointless hop
Two fine-grained senior points to close:
CompletableFutureestablishes a happens-before relationship: everything the completing thread did beforecomplete(v)is visible to dependent stages. So to read the produced value you need no extravolatileor lock. (But that guarantee covers only the value itself, not any side mutable state you share between stages.)thenApplyAsync(fn, sameExecutor)still hops: even if you're already running on that executor, the*Asyncvariant submits a new task — extra queueing and a context switch. If you don't need the hop, the non-async combinator is cheaper. A deliberate hop is only worth it when you want to escape an expensive thread (e.g. an event-loop thread or an IO pool).
Senior interview questions (hard)
No. The mayInterruptIfRunning flag is ignored by CompletableFuture, because a CF holds no reference to that thread. cancel only makes the future itself exceptional so downstream stages are abandoned, but the underlying work runs to completion. Real cancellation must be cooperative, or come from a client that propagates it (or from StructuredTaskScope).
Nothing. The future is completed with a TimeoutException, but the query/connection is still alive and stays checked out. Under load this becomes a connection-pool leak: timeouts fire fast while the real work piles up until the pool is exhausted. The real timeout must be set in the client (JDBC/HTTP) itself, not just in the CompletableFuture layer.
The fast failure. anyOf completes with the first thing to complete — and an exceptional completion is still a completion. So if the first server errors immediately, the whole result fails even though the second was about to answer correctly. Also, anyOf doesn't cancel the losers, so they keep burning resources. Correct hedging needs manual cancellation of the competitors.
Because any caller can call complete(...), cancel(), or obtrudeValue(...) and hijack your downstream graph — like making a private field public. Fix: minimalCompletionStage() (Java 9) returns a CompletionStage without completion methods, so calling them throws UnsupportedOperationException; or copy() for a protected copy.
Up to Java 23, entering synchronized caused pinning: the monitor was registered on the platform carrier, keeping the virtual thread pinned to it. If you blocked inside that block (I/O), the carrier couldn't be freed; if all carriers pinned this way, parallelism dropped to zero and it could deadlock. JEP 491 in Java 24 tied the monitor to the virtual thread itself; now synchronized doesn't pin. Only native frames still pin. Pre-24 workaround: ReentrantLock instead of synchronized.
With allOf, on the first failure the combined future goes exceptional but the other futures keep running as orphans, burning resources — nobody owns their lifecycle. With StructuredTaskScope, subtasks are children of the scope; when one fails, the rest are cancelled and interrupted, and closing the scope guarantees no orphan thread is left. That's what "structured" means.
Since JDK 10, availableProcessors() is cgroup-aware and returns 1, so the common pool's parallelism becomes 1 - 1 = 0 and the pool effectively collapses: the parallel stream runs serially on the calling thread. The speedup you saw on a multi-core laptop vanishes in prod. Fix: set parallelism explicitly with -Djava.util.concurrent.ForkJoinPool.common.parallelism=N or -XX:ActiveProcessorCount=N, or run that stream on a dedicated ForkJoinPool.
Because exceptionally can only recover to a plain value; it can't kick off an async retry. exceptionallyCompose (Java 12) maps the failure to another future, so you can return a fresh supplyAsync (the retry) inside it. Pair it with delayedExecutor(200, MILLISECONDS) so the next attempt runs after a delay, without a separate Timer.
- Cancel and
orTimeoutdo not stop the underlying work — resources leak unless you truly propagate cancellation. anyOf/allOfnever kill the siblings;anyOfalso completes on a fast failure.- Hand out futures protected via
minimalCompletionStage(). - In containers, size the common pool explicitly; on 1 vCPU, parallel streams go serial.
- Pinning: pre-Java-24
synchronizedpinned the carrier;JEP 491fixed it. Don't pool virtual threads or use them for CPU work. - Context (MDC/security) is lost across async hops;
ScopedValue(Java 25) is the structuredThreadLocalreplacement. - Structured Concurrency is the real successor to fragile graphs: all succeed, or the rest are interrupted.
- Two tools, two jobs:
ForkJoinPoolis a CPU-bound work engine with work-stealing;CompletableFutureis an async composition graph of callbacks that creates no threads itself. - The common pool isn't yours: a JVM-wide singleton with
cores - 1workers, shared by parallel streams and default CF. Blocking work on it = starvation = the root of most bugs. - The thread rule: non-async combinators run on the completing thread;
*Asyncwith no executor runs on the common pool. For blocking work, always pass a dedicated executor. - map vs flatMap:
thenApplyfor a plain value,thenComposefor a nested future. - Errors slide down the chain wrapped in a
CompletionException— callgetCause().exceptionally/handlerecover,whenCompleteonly records. - fork/join: fork one, compute the other inline, then join; tune
THRESHOLD. - Virtual threads (Java 21, JEP 444) make blocking cheap and have their own separate FIFO pool; prefer them for I/O fan-out, and keep
CompletableFuturefor genuine composition graphs.