CompletableFuture recursive restart on exception from exceptionaly() block

Viewed 376

Having a hard time to do a reliable retry of a background task which sends request to let's say mail service in order to get latest emails. Once emails successfully received the execution should continue in thenAccept() block - persist emails, however if exception occurs I have to rerun mail retrieval until successful attempt and on success should persist mails and stop. Please take a look and advice if I do it wrong.

   private void retrieveMailsAsync(User user) {
    CompletableFuture.supplyAsync(() -> {
        try {
            return mailService.getEmails(user.getName(), user.getPassword());
        } catch (InvalidAuthentication | TimeoutException | BadGatewayException e) {
            throw new CompletionException(e);
        }
    }).thenAccept(email -> {
        mailService.persist(email);
    }).exceptionally(ex -> {
        log.log(Level.SEVERE, "Exception retrieveMailsAsync emails, Retrying retrieveMailsAsync:: ", ex.getCause());
        retrieveMailsAsync(user);
        return null;
    });
}

P.S please also take a look at how I'm handling checked exception wrapping it into CompletionException and rethrowing - the main idea here to handle all exceptions (defined checked and runtime) in one exceptionally() block rather than logging them in catch block and return null.

Thanks guys in advance, hope I'm not doing pretty stupid stuff, or at least there is already reliable solutions exists for Java 8.

2 Answers

I think you can achieve that via :

public static void main(String[] args) {
    String result = call(new User().setName("name").setPassword("p")).join();
    System.out.println(result);
}

private static CompletableFuture<String> call(User user) {
    CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> retrieveMailsAsync(user));
    return future.handleAsync((String result, Throwable ex) -> {
        // or any other Predicate that is satisfied against ex
        if(ex != null) {
            return call(user);
        } else {
            return future;
        }
    }).thenCompose(Function.identity());
}

EDIT

So what stays in your way to change the code above, for example, to:

static ExecutorService service = Executors.newFixedThreadPool(1);

public static void main(String[] args) {
    call(new User().setName("name").setPassword("p"))
         // chain any other action here, like mailService.persist(email);
         .thenAcceptAsync(
            System.out::println,
            service
    );
    System.out.println("Continue main thread");
}

private static CompletableFuture<String> call(User user) {
    CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> retrieveMailsAsync(user), service);
    return future.handleAsync((String result, Throwable ex) -> {
        // or any other Predicate that is satisfied against ex
        if(ex != null) {
            return call(user);
        } else {
            return future;
        }
    }).thenCompose(Function.identity());
}

What I meant in my comment was this:

   private void retrieveMailsAsync(User user) {
    CompletableFuture.supplyAsync(() -> {
        while (continueQuery()) { // true for infinite retries, or some other logic
          try {
              return mailService.getEmails(user.getName(), user.getPassword());
          } catch (InvalidAuthentication | TimeoutException | BadGatewayException e) {
              log.log(Level.SEVERE, "Exception retrieveMailsAsync emails, Retrying retrieveMailsAsync: ", e);
          }
        }
        return null;
    }).thenAccept(email -> {
        mailService.persist(email);
    });
}

Ie you just retry in the submitted runnable until you don't get an exception anymore.

Related