ComputableFuture Parallel Calls

Viewed 82

I want to run 10 APIs in parallel. Each API returns some values. I want to stop when total count of values are equal to 100 i.e I don't want to wait for all the 10 APIs if I get my 100 results before getting results of all APIs. So I was thinking to play with CompletableFuture.anyOf() in loop and return but I am unable to figure out the right syntax for the same. Also if there is some other efficient way to do so?

Please reply.

Thanks in advance!

1 Answers

You can use CountDownLatch

    ExecutorService executor = Executors.newFixedThreadPool(10);

    List<Integer> results = new ArrayList<>();
    int maxResults = 100;
    CountDownLatch latch = new CountDownLatch(maxResults);

    for (int i = 0; i < 10; ++i) {
        int apiNumber = i;
        executor.execute(() -> {
            while (results.size() < maxResults) {
                try {
                    Thread.sleep(new Random().nextInt(1000));
                    System.out.println("API-" + apiNumber + " call");
                    int[] newValues = new Random().ints(0, 10).limit(new Random().nextInt(10)).toArray(); // API call
                    for (int value : newValues) {
                        synchronized (results) {
                            if (results.size() < maxResults) {
                                results.add(value);
                                latch.countDown();
                            }
                        }
                    }
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            }
        });
    }

    latch.await();

    System.out.println("The results have been calculated (" + results.size() + ")");


    executor.shutdown();

or Phaser

    ExecutorService executor = Executors.newFixedThreadPool(10);

    List<Integer> results = new ArrayList<>();
    int maxResults = 100;
    Phaser phaser = new Phaser(1);
    phaser.register();

    for (int i = 0; i < 10; ++i) {
        int apiNumber = i;
        executor.execute(() -> {
            while (results.size() < maxResults) {
                try {
                    Thread.sleep(new Random().nextInt(1000));
                    System.out.println("API-" + apiNumber + " call");
                    int[] newValues = new Random().ints(0, 10).limit(new Random().nextInt(10)).toArray(); // API call
                    for (int value : newValues) {
                        synchronized (results) {
                            results.add(value);
                            if (results.size() >= maxResults) {
                                phaser.arrive();
                                break;
                            }
                        }
                    }
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            }
        });
    }

    phaser.arriveAndAwaitAdvance();

    System.out.println("The results have been calculated (" + results.size() + ")");

    executor.shutdown();
Related