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

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 و completeOnTimeout — از جاوا ۹

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 از قبل تمام شده باشد (مثلاً completedFuturethenApply تو به‌صورت 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() یا get()؟

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/join

یک شاخه را 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
common pool مالِ تو نیست

چون این استخر بین 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 تقسیم می‌کند.

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 استفاده کن — که البته هزینه دارد چون باید هماهنگی برقرار کند.
شرکت‌پذیر (associative) یعنی چه؟

یعنی (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ها استخرِ خودشان را دارند

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.


سؤالات مصاحبه

حالا بخشی که برایش آمده‌ای. هر سؤال را جوری جواب می‌دهیم که هم درست باشد هم نشان‌دهنده‌ی عمق.

۱. یک callback از نوع `thenApply` (غیر-async) روی چه نخی اجرا می‌شود؟

روی نخی که future بالادست را کامل کرده. اگر بالادست هنگام اتصالِ callback از قبل کامل شده باشد، به‌صورت همگام روی نخِ فراخوان اجرا می‌شود. فقط نسخه‌های *Async پرش به یک executor را تضمین می‌کنند (اگر هیچ executorی داده نشود، common pool).

۲. تفاوت `thenApply` و `thenCompose`؟

thenApply همان map است (T -> UthenCompose همان flatMap (T -> CompletableFuture<U>). وقتی تابعت خودش ناهمگام است و یک future برمی‌گرداند، از thenCompose استفاده کن تا از CompletableFuture<CompletableFuture<U>> تودرتو پرهیز شود.

۳. `exceptionally` در برابر `handle` در برابر `whenComplete`؟

exceptionally فقط روی شکست اجرا می‌شود و می‌تواند بازیابی کند (همان نوع). handle روی هر دو مسیر اجرا می‌شود، می‌تواند مقدار را تبدیل و بازیابی کند و نوعِ نتیجه را عوض کند. whenComplete روی هر دو اجرا می‌شود اما فقط side-effect است — نمی‌تواند بازیابی کند و استثنای اصلی را به پایین‌دست دوباره پرتاب می‌کند.

۴. (تله) چرا `catch (SomeException e)` دورِ `future.join()` گاهی مطابقت نمی‌کند؟

چون join() یک CompletionException پرتاب می‌کند که علت واقعی را می‌پیچد. باید CompletionException را بگیری و getCause() را بررسی کنی، وگرنه نوعِ اصلی مطابقت نخواهد کرد.

۵. (سخت) چرا انجام فراخوان‌های JDBC درون `parallelStream()` خطرناک است؟

parallel streamها روی common pool در سطح کل JVM اجرا می‌شوند (تقریباً هسته - 1 worker). یک فراخوان مسدودکننده‌ی JDBC یک worker را پارک می‌کند؛ به‌اندازه‌ی کافی از این‌ها استخر را قحطی‌زده می‌کنند و در همین فرآیند هر parallel stream و CompletableFuture پیش‌فرض را متوقف می‌کنند. راه‌حل: executor اختصاصی، ManagedBlocker، یا virtual threadها.

۶. سرقت کار (work-stealing) چطور وظایف را زمان‌بندی می‌کند؟

هر worker یک deque دارد؛ وظایف خودش را LIFO از سر push/pop می‌کند (کش‌داغ). workerهای بیکار از دُمِ deque یک worker دیگر FIFO سرقت می‌کنند (قدیمی‌ترین/بزرگ‌ترین وظیفه). این کار بار را با حداقل رقابت متعادل می‌کند.

۷. در یک `RecursiveTask`، چرا `left.fork(); right.compute(); left.join();` به‌جای fork کردن هر دو؟

fork کردنِ هر دو و سپس join هر دو، یک جای صف را هدر می‌دهد و یک سرقت را اجباری می‌کند؛ محاسبه‌ی یک شاخه به‌صورت inline نخِ فعلی را مشغولِ کارِ مفید نگه می‌دارد در حالی که شاخه‌ی دیگر برای سرقت در دسترس است. این فرمِ متعارف و کارآمدترین است.

۸. درجه‌ی موازیِ پیش‌فرضِ common pool چیست و چطور تغییرش می‌دهی؟

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 جداگانه جمع‌آوری کن.

۱۱. (سخت) آیا `CompletableFuture` تنبل است یا مشتاق؟ با `Mono` در Reactor مقایسه کن.

مشتاق — supplyAsync بلافاصله شروع به اجرا می‌کند. Mono/Flux تنبل‌اند: تا subscribe() هیچ چیز اجرا نمی‌شود. این روی retry، caching و اینکه آیا ساختن-اما-مصرف‌نکردن کاری انجام می‌دهد اثر می‌گذارد.

۱۲. virtual threadها از چه executorی استفاده می‌کنند، و آیا همان common pool است؟

خیر. virtual threadها روی یک ForkJoinPool جداگانه در حالت FIFO اجرا می‌شوند، به‌صورت پیش‌فرض یک نخِ carrier به‌ازای هر هسته (‏jdk.virtualThreadScheduler.parallelism). common pool (‏LIFO) توسط parallel streamها و CompletableFuture پیش‌فرض استفاده می‌شود.

۱۳. (تله) `thenApplyAsync(fn)` (بدون executor) را ۱۰۰ هزار بار در یک سرویسِ تحتِ بار صدا می‌زنی و می‌بینی توان عملیاتی فرو می‌ریزد. چرا؟

همه‌ی ۱۰۰ هزار callback روی common pool فرود می‌آیند؛ اگر هر کدام مسدود شود یا حجم از هسته - 1 جای بهره‌ور فراتر رود، صف‌بندی/قحطی می‌گیری. یک executor با اندازه‌ی مناسب پاس بده، یا fan-out را محدود کن.

۱۴. چه زمانی روی جاوا ۲۱ همچنان `CompletableFuture` را به virtual threadها ترجیح می‌دهی؟

وقتی به معناشناسیِ ترکیب آن نیاز داری: ترکیبِ نتایجِ ناهمگون (thenCombine)، انتظار برای چندتا (allOf/anyOf)، hedging (anyOf)، timeoutهای اعلانی (orTimeout)، و ساختنِ گراف‌های ناهمگامِ قابل‌استفاده‌ی مجدد. برای «N فراخوان مسدودکننده را اجرا کن و join کن»، virtual threadها / StructuredTaskScope ساده‌ترند.

۱۵. (سخت) چرا یک زنجیره‌ی `thenCompose` می‌تواند کاملاً روی یک نخ از استخرِ IO اجرا شود؟

چون 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` کارِ زیرین را متوقف نمی‌کند

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 درست

اگر واقعاً 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‌ای بسیار بزرگ‌تر از سهمِ واقعی‌ات بسازد.
عددِ 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.

جاوا ۲۴ تقریباً همه‌ی pinningها را برداشت

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` جانشینِ ساختاریافته است

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ِ یتیمی جا نمی‌ماند
API هنوز در حالِ تثبیت است

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 یک استثنا پرتاب شود و مرحله‌ی بالادست هم قبلاً شکست خورده بوده، استثنای اصلی گم می‌شود و استثنای جدیدِ تو پایین می‌رود. پس در 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) فرار کنی.

سؤالات مصاحبه‌ی سنیور (سطحِ سخت)

۱. با `cancel(true)` روی یک `CompletableFuture`، آیا نخی که `supplyAsync` را اجرا می‌کند interrupt می‌شود؟

نه. پارامترِ mayInterruptIfRunning در CompletableFuture نادیده گرفته می‌شود؛ چون CF هیچ ارجاعی به آن نخ ندارد. cancel فقط خودِ future را استثنایی می‌کند تا مرحله‌های پایین‌دست رها شوند، ولی کارِ زیرین تا آخر اجرا می‌شود. کنسلِ واقعی باید تعاونی باشد یا از کلاینتی بیاید که کنسل را propagate می‌کند (یا از StructuredTaskScope).

۲. `orTimeout(2, SECONDS)` شلیک می‌کند؛ چه بلایی سرِ کانکشنِ دیتابیسِ زیرین می‌آید؟

هیچ. future با TimeoutException بسته می‌شود اما کوئری/کانکشن همچنان زنده است و checked-out می‌ماند. زیرِ بار، این می‌شود نشتیِ کانکشن‌پول: تایم‌اوت‌ها سریع رخ می‌دهند ولی کارِ واقعی تلنبار می‌شود تا pool خفه شود. تایم‌اوتِ واقعی باید در خودِ کلاینت (JDBC/HTTP) تنظیم شود، نه فقط در لایه‌ی CompletableFuture.

۳. در `anyOf`، بینِ یک شکستِ سریع و یک موفقیتِ کند، کدام برنده می‌شود؟ چرا این برای hedging خطرناک است؟

شکستِ سریع. anyOf با اولین چیزی که کامل شود کامل می‌شود — و کامل‌شدنِ استثنایی هم کامل‌شدن است. پس اگر سرورِ اول فوری خطا بدهد، کلِ نتیجه شکست می‌خورد حتی اگر سرورِ دوم داشت درست جواب می‌داد. ضمناً anyOf بازنده‌ها را کنسل نمی‌کند، پس منبع می‌سوزانند. hedgingِ درست نیاز به منطقِ کنسلِ دستیِ رقبا دارد.

۴. چرا برگرداندنِ یک `CompletableFuture` خام از یک متدِ کتابخانه یک بوی بدِ طراحی است، و راه‌حلش چیست؟

چون هر فراخوانی می‌تواند complete(...), cancel() یا obtrudeValue(...) صدا بزند و گرافِ پایین‌دستِ تو را بدزدد — مثلِ public کردنِ یک فیلدِ private. راه‌حل: minimalCompletionStage() (جاوا ۹) که یک CompletionStage بدونِ متدهای completion می‌دهد و صدا زدنشان UnsupportedOperationException می‌دهد؛ یا copy() برای یک کپیِ محافظت‌شده.

۵. پیش از جاوا ۲۴، چطور یک بلاکِ `synchronized` می‌توانست کلِ scheduler‌ِ virtual threadها را خفه کند، و در ۲۴ چه عوض شد؟

تا جاوا ۲۳، ورود به synchronized باعثِ pinning می‌شد: monitor روی carrierِ پلتفرمی ثبت می‌شد و virtual thread به carrier میخکوب می‌ماند. اگر داخلِ آن بلاک بلاک می‌شدی (I/O)، carrier آزاد نمی‌شد؛ اگر همه‌ی carrierها این‌طور می‌شدند، موازی‌سازی به صفر می‌رسید و می‌شد deadlock. JEP 491 در جاوا ۲۴ monitor را به خودِ virtual thread گره زد؛ حالا synchronized پین نمی‌کند. فقط فریم‌های native هنوز پین می‌کنند. راه‌حلِ پیش از ۲۴: ReentrantLock به‌جای synchronized.

۶. تفاوتِ رفتاری بینِ شکستِ یک `allOf` و شکستِ یک `StructuredTaskScope` نسبت به subtaskهای دیگر چیست؟

در allOf، به‌محضِ اولین شکست future‌ِ ترکیبی استثنایی می‌شود ولی بقیه‌ی futureها یتیم و بی‌مراقب ادامه می‌دهند و منبع می‌سوزانند — چون هیچ‌کس مالکِ چرخه‌ی عمرِ آن‌ها نیست. در StructuredTaskScope، subtaskها فرزندِ scope‌اند؛ با شکستِ یکی، بقیه کنسل و interrupt می‌شوند و بستنِ scope تضمین می‌کند هیچ نخِ یتیمی جا نمی‌ماند. این همان «ساختاریافته» بودن است.

۷. روی یک pod با `cpu limit: 1`، رفتارِ `parallelStream()` چه می‌شود و چرا؟

از 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 اختصاصی بساز.

۸. چرا برای retry با backoff به `exceptionallyCompose` نیاز داری نه `exceptionally`؟

چون exceptionally فقط می‌تواند خطا را به یک مقدارِ ساده برگرداند؛ نمی‌تواند یک عملیاتِ asyncِ دوباره را راه بیندازد. exceptionallyCompose (جاوا ۱۲) خطا را به یک future دیگر می‌بندد، پس می‌توانی داخلش یک supplyAsync جدید (تلاشِ دوباره) برگردانی. همراهش delayedExecutor(200, MILLISECONDS) را بگذار تا بدونِ Timer جداگانه، تلاشِ بعدی با تأخیر اجرا شود.

جمع‌بندیِ لایه‌ی سنیور
  • کنسل و orTimeout کارِ زیرین را متوقف نمی‌کنند — منبع نشت می‌کند مگر کنسل را واقعاً propagate کنی.
  • anyOf/allOf برادرها را نمی‌کشند؛ anyOf با شکستِ سریع هم کامل می‌شود.
  • future را با minimalCompletionStage() محافظت‌شده بیرون بده.
  • در کانتینر، سایزِ common pool را صریح بگذار؛ روی ۱ vCPU، parallel stream سریال می‌شود.
  • pinning: پیش از جاوا ۲۴ synchronized carrier را میخکوب می‌کرد؛ JEP 491 رفعش کرد. virtual thread را pool نکن و برای CPU به کار نبر.
  • context (MDC/security) در پرش‌های async گم می‌شود؛ ScopedValue (جاوا ۲۵) جانشینِ ساختاریافته‌ی ThreadLocal است.
  • Structured Concurrency جانشینِ واقعیِ گراف‌های شکننده است: یا همه موفق، یا بقیه interrupt.
جمع‌بندی
  • دو ابزار، دو کار: ForkJoinPool موتورِ کارِ CPU-bound با سرقت کار است؛ CompletableFuture یک گراف ترکیبِ ناهمگام از callbackهاست که خودش نخ نمی‌سازد.
  • common pool مالِ تو نیست: یک singletonِ سراسری با هسته - 1 worker که بین 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.

Roadmap for this lesson

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.
One sentence that holds it all together

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.

The kitchen vs the recipe card

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:

  • ForkJoinPool is 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 of parallelStream() and the default executor for CompletableFuture.
  • CompletableFuture is 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.
The single most important sentence in this chapter

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:

  • supplyAsync means "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 in f2.
  • completedFuture means "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 CompletableFuture and later fill it from any thread with complete(...). 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.

The receipt is issued right now

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.

A gift box vs a nested box

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 boxCompletableFuture<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());
The big allOf trap

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:

The exception isn't thrown here — it "slides down"

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>
Three people at the assembly line

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 CompletionException wrapping the real cause. Always call ex.getCause() to inspect the real one.
  • whenComplete does not swallow the exception — the returned future still fails. It's for logging/cleanup, not recovery.
  • handle runs 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 *Async variant (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 and completeOnTimeout — since Java 9

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:

The two "which thread?" rules

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.

The clerk who does the next person's job at their own desk

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, then step1 also occupies an io thread — a long thenApply chain can starve your IO pool even though only one stage was "async."
  • If the future is already done (e.g. completedFuture), your thenApply runs 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() - 1 of 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() or get()?

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.

Your own desk and the idle coworker

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));
The golden fork/join rule

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
The common pool isn't yours

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;
    }
});
What's a compensation thread?

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.

Spliterator means "splitter"

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; but LinkedList and Iterator-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. reduce must be associative; collect uses a combiner to 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.)
  • forEach guarantees no ordering; if you need encounter order, use forEachOrdered — which costs, because it has to coordinate.
What does "associative" mean?

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
Why does this trick work?

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).

One parcel vs a conveyor belt

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."

A crucial semantic difference: eager vs lazy

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.

The cheap worker vs the expensive worker

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 have their own pool

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 CompletableFuture for 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 dedicated Executor.
  • CompletionException wrapping. Your catch/exceptionally sees the wrapper; call getCause().
  • whenComplete is not recovery. It re-throws; only handle/exceptionally recover.
  • allOf loses values. Re-join each future after allOf completes.
  • Non-async combinators steal the completing thread. A long thenApply chain runs on whatever thread finished the previous stage — possibly your IO pool, possibly a common-pool worker.
  • Forgetting the future is eager. supplyAsync starts 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

  1. Own your executors for anything that blocks. Never trust the common pool for I/O.
  2. Name your threads (ThreadFactory) so stack dumps are debuggable.
  3. Prefer join() at the edge of the graph; keep the interior non-blocking.
  4. Add orTimeout/completeOnTimeout to every external call.
  5. Handle errors with handle/exceptionally; always getCause().
  6. Use RecursiveTask only for CPU-bound divide-and-conquer with a tuned threshold.
  7. On Java 21+, prefer virtual threads / structured concurrency for I/O fan-out; reserve CompletableFuture for 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.

1. What thread runs a `thenApply` (non-async) callback?

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).

2. Difference between `thenApply` and `thenCompose`?

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>>.

3. `exceptionally` vs `handle` vs `whenComplete`?

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.

4. (Gotcha) Why does `catch (SomeException e)` around `future.join()` sometimes not match?

Because join() throws a CompletionException wrapping the real cause. You must catch CompletionException and inspect getCause(), or the original type won't match.

5. (Hard) Why is doing JDBC calls inside `parallelStream()` dangerous?

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.

6. How does work-stealing schedule tasks?

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.

7. In a `RecursiveTask`, why `left.fork(); right.compute(); left.join();` instead of forking both?

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.

8. What is the default parallelism of the common pool, and how do you change it?

availableProcessors() - 1. Change via -Djava.util.concurrent.ForkJoinPool.common.parallelism=N. Note it's a global singleton shared across the JVM.

9. (Gotcha) What does this print?
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.

10. (Find the bug)
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.

11. (Hard) Is `CompletableFuture` lazy or eager? Contrast with Reactor `Mono`.

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.

12. Which executor do virtual threads use, and is it the common pool?

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.

13. (Gotcha) You call `thenApplyAsync(fn)` (no executor) 100k times across a service under load and see throughput collapse. Why?

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.

14. When would you still choose `CompletableFuture` over virtual threads on Java 21?

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.

15. (Hard) Why can a `thenCompose` chain still run entirely on one IO-pool thread?

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.

Roadmap for this section

(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` does not stop the underlying work

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:

  • anyOf completes 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 with anyOf (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 anyOf returns the winner, all the other futures still run to completion and burn resources. Likewise allOf completes exceptionally on the first failure, but the remaining futures keep running unmonitored.
The correct hedging pattern

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 - 1 becomes 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.
Size the pool explicitly in containers

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.

Java 24 removed almost all pinning

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.

If you're stuck on pre-Java-24

On I/O hot paths, replace synchronized with ReentrantLockjava.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. InheritableThreadLocal is worse still, because it's copied into every child.
  • Context loss across async hops: values like MDC (your log traceId) and SecurityContext ride on ThreadLocal. When a thenApplyAsync hops to a different thread, that ThreadLocal doesn't come along — suddenly your logs lose the traceId and your security context is empty. The same problem occurs across an executor.submit(...) boundary.
`ScopedValue` is the structured replacement

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
The API is still stabilizing

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
}
An exception in `whenComplete` masks the original

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:

  • CompletableFuture establishes a happens-before relationship: everything the completing thread did before complete(v) is visible to dependent stages. So to read the produced value you need no extra volatile or 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 *Async variant 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)

1. Does `cancel(true)` on a `CompletableFuture` interrupt the thread running the supplier?

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).

2. `orTimeout(2, SECONDS)` fires — what happens to the underlying DB connection?

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.

3. In `anyOf`, between a fast failure and a slow success, which wins? Why is this dangerous for hedging?

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.

4. Why is returning a raw `CompletableFuture` from a library method a design smell, and what's the fix?

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.

5. Before Java 24, how could a `synchronized` block choke the entire virtual-thread scheduler, and what changed in 24?

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.

6. What's the behavioral difference between an `allOf` failing and a `StructuredTaskScope` failing, regarding the other subtasks?

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.

7. On a pod with `cpu limit: 1`, what happens to `parallelStream()` and why?

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.

8. Why do you need `exceptionallyCompose` rather than `exceptionally` for retry with backoff?

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.

The senior layer in a nutshell
  • Cancel and orTimeout do not stop the underlying work — resources leak unless you truly propagate cancellation.
  • anyOf/allOf never kill the siblings; anyOf also 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 synchronized pinned the carrier; JEP 491 fixed 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 structured ThreadLocal replacement.
  • Structured Concurrency is the real successor to fragile graphs: all succeed, or the rest are interrupted.
In a nutshell
  • Two tools, two jobs: ForkJoinPool is a CPU-bound work engine with work-stealing; CompletableFuture is an async composition graph of callbacks that creates no threads itself.
  • The common pool isn't yours: a JVM-wide singleton with cores - 1 workers, 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; *Async with no executor runs on the common pool. For blocking work, always pass a dedicated executor.
  • map vs flatMap: thenApply for a plain value, thenCompose for a nested future.
  • Errors slide down the chain wrapped in a CompletionException — call getCause(). exceptionally/handle recover, whenComplete only 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 CompletableFuture for genuine composition graphs.