দুটি Future-এর বাইরে

thenCombine কেবল দুটি স্বাধীন ফিউচার জোড়ার জন্য ভালো, কিন্তু তালিকার জন্য এটি স্কেল করে না। বাস্তবে ফ্যান-আউট/ফ্যান-ইন খুব সাধারণ। যেমন:

  • একটি কালেকশনের প্রতিটি আইটেমের জন্য async কল চালু করা
  • দশজন সাপ্লায়ারের কাছ থেকে দাম আনা
  • পাঁচটি রেপ্লিকাকে ping করা

এরপর সব ফিউচার শেষ হলে বা যেকোনো একটি শেষ হলেই পরবর্তী কাজ করা হয়। CompletableFuture.allOf এবং CompletableFuture.anyOf ঠিক এই কাজের জন্য তৈরি।

Fan-Out: একসাথে অনেকগুলো Future চালু করা

ব্যাচ শুরু করা সোজা: প্রতিটি আইটেমকে CompletableFuture-এ map করুন supplyAsync দিয়ে, আদর্শভাবে workload অনুযায়ী sized pool ব্যবহার করুন।

import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class QuoteFanOut {
    public static void main(String[] args) {
        List<String> supplierIds = List.of("s1", "s2", "s3", "s4", "s5");
        ExecutorService pool = Executors.newFixedThreadPool(5);

        List<CompletableFuture<String>> quoteFutures = supplierIds.stream()
                .map(id -> CompletableFuture.supplyAsync(() -> fetchQuote(id), pool))
                .toList();

        pool.shutdown();
    }

    private static String fetchQuote(String supplierId) {
        // বাস্তবে এখানে HTTP কল বা ব্লকিং I/O হবে
        return supplierId + "-quote";
    }
}

এই কোডে পাঁচটি fetchQuote কল একসাথে চালু হয়। এখানে মূল থ্রেড এখনো ব্লক হয়নি; প্রতিটি ফিউচার নিজের ফলাফল নিয়ে অপেক্ষা করছে।

Fan-In: allOf পুরো ব্যাচের জন্য অপেক্ষা করে

CompletableFuture.allOf varargs হিসেবে ফিউচার নেয় এবং একটি CompletableFuture<Void> ফেরত দেয়। এটি তখন complete হয় যখন প্রতিটি ইনপুট ফিউচার সফলভাবে বা exceptionally শেষ হয়ে যায়।

CompletableFuture<?>[] futuresArray = quoteFutures.toArray(new CompletableFuture<?>[0]);
CompletableFuture<Void> allDone = CompletableFuture.allOf(futuresArray);

প্রথম awkward অংশ হলো Void রিটার্ন টাইপ: allOf শুধু জানায় সব শেষ হয়েছে, কিন্তু আলাদা ফলাফল দেয় না। ফলাফল পেতে হলে allDone-এর ওপর thenApply চালিয়ে প্রতিটি ফিউচারে join() করতে হয়।

CompletableFuture<List<String>> allQuotes = allDone.thenApply(v ->
        quoteFutures.stream()
                .map(CompletableFuture::join)
                .toList()
);

যেহেতু allDone তখনই complete হয় যখন তালিকার সব ফিউচার শেষ হয়ে গেছে, তাই ভেতরে join() কল করলে আর ব্লক হয় না। join() অনেকটা get()-এর মতো, তবে এটি checked ExecutionException-এর বদলে unchecked CompletionException ছোড়ে, তাই ল্যাম্বডার ভেতরে ব্যবহার করা সহজ।

একটি ব্যর্থতা বাকিদের থামায় না, কিন্তু allOf-কে বিষিয়ে তোলে

পাঁচটি quote কলের একটি ব্যর্থ হলে বাকি চারটি চলতেই থাকে। allOf অন্য ফিউচারগুলোকে cancel করে না। তবে allDone নিজে তখন exceptionally complete হয়। সফল ফিউচারগুলোর ওপর আলাদাভাবে join() করলে ফলাফল পাওয়া যেতে পারে, কিন্তু allDone.join() নিজে, অথবা exceptionally/handle ছাড়া সরাসরি chained stage, সেই ব্যর্থতাটাই ছড়িয়ে দেবে।

একটি সাধারণ প্যাটার্ন হলো প্রতিটি ফিউচারকে allOf-এ ঢোকানোর আগে exceptionally দিয়ে নিজস্ব fallback দেওয়া। এতে একজন সাপ্লায়ার টাইমআউট করলে পুরো ব্যাচের ফলাফল নষ্ট হয় না।

List<CompletableFuture<String>> safeQuoteFutures = supplierIds.stream()
        .map(id -> CompletableFuture.supplyAsync(() -> fetchQuote(id), pool)
                .exceptionally(ex -> "fallback-for-" + id))
        .toList();

CompletableFuture<?>[] safeFuturesArray = safeQuoteFutures.toArray(new CompletableFuture<?>[0]);
CompletableFuture<Void> safeAllDone = CompletableFuture.allOf(safeFuturesArray);

CompletableFuture<List<String>> safeAllQuotes = safeAllDone.thenApply(v ->
        safeQuoteFutures.stream()
                .map(CompletableFuture::join)
                .toList()
);

এখানে প্রতিটি ফিউচার ব্যর্থ হলে একটি fallback মান পায়। ফলে safeAllDone স্বাভাবিকভাবে complete হয় এবং join() কোনো exception ছাড়াই সব ফলাফল দেয়।

Fan-In: anyOf প্রথম শেষ হওয়া Future-কে নেয়

CompletableFuture.anyOf একই রকমভাবে varargs নেয়, তবে এটি complete হয় যখন যেকোনো একটি ইনপুট ফিউচার শেষ হয়। এর রিটার্ন টাইপ CompletableFuture<Object> হওয়াই এখানকার awkward অংশ।

CompletableFuture<?>[] futuresArray = quoteFutures.toArray(new CompletableFuture<?>[0]);
CompletableFuture<Object> firstQuote = CompletableFuture.anyOf(futuresArray);

firstQuote.thenAccept(q -> System.out.println("প্রথম দাম পৌঁছেছে: " + q));

যে ফিউচারটি সবার আগে complete হয়, তার ফলাফলই anyOf ফেরত দেয়। টাইপ Object হওয়ায় প্রয়োজন অনুযায়ী cast করতে হয়। যদি প্রথম complete হওয়া ফিউচারটি exceptionally শেষ হয়, তাহলে anyOf-এর ফিউচারটিও exceptionally complete হয়।

Share