Caching parallel request in Spring Webflux Mono

Viewed 428

We are using spring webflux (project reactor), as part of the requirement we need to call one API from our server.

For the API call, we need to cache the response. So we are using Mono.cache operator.

It caches the response Mono<ResponseDto> and the next time the same API call happens, it will get it from the cache. Following is example implementation

public Mono<ResponseDto> getResponse() {
    if (res == null) {
      res =
          fetchResponse()
              .onErrorMap(Exception.class, (error) -> new CustomException())
              .cache(
                  r -> Duration.ofSeconds(r.expiresIn()),
                  error -> Duration.ZERO,
                  () -> Duration.ZERO);
    }
    return res;
  }

The problem is if the server calls the same API call twice ( for example Mono.zip) at the same time, then the response is not cached and we actually call it twice.

Is there any out of box solution available to this problem? Instead of caching the Response, can we cache the Mono itself so that both requests subscribe to the same Mono hence both are executed after a Single API call response?

It should also work with sequential execution too - I am afraid that if we cache the Mono then once the request is completed, the subscription is over and no other process can subscribe to it.

Cases

3 Answers

You can initialize the Mono in the constructor (assuming it doesn't depend on any request time parameter). Using cache operator will prevent multiple subscriptions to the source.

class MyService {
    private final Mono<ResponseBodyDto> response;

    public MyService() {
        response = fetchResponse()
            .onErrorMap(Exception.class, (error) -> new CustomException())
            .cache(
                r -> Duration.ofSeconds(r.expiresIn()),
                error -> Duration.ZERO,
                () -> Duration.ZERO);
    }

    public Mono<ResponseDto> getResponse() {
        return response;
    }
}

If there is a dependency on request time parameters, you should consider some custom caching solution.

Project Reactor provides a cache utility CacheMono that is non-blocking but can stampede.

AsyncCache will be better integration, for the first lookup with key "K" will result in a cache miss, it will return a CompletableFuture of the API call and for the second lookup with the same key "K" will get the same CompletableFuture object.

The returned future object can be converted to/from Mono with Mono.fromFuture()

 public Mono<ResponseData> lookupAndWrite(AsyncCache<String, ResponseData> cache, String key) {
return Mono.defer(
    () ->
        Mono.fromFuture(
            cache.get(
                key,
                (searchKey, executor) -> {
                  CompletableFuture<ResponseData> future = callAPI(searchKey).toFuture();
                  return future.whenComplete(
                      (r, t) -> {
                        if (t != null) {
                          cache.synchronous().invalidate(key);
                        }
                      });
                })));}

You could use CacheMono from io.projectreactor.addons:reactor-extra to wrap non-reactive cache implementation like Guava Cache or simple ConcurrentHashMap. It doesn't provide an "exactly-once" guarantee and parallel requests could result in cache misses, but in many scenarios, it should not be an issue.

Here is an example with Guava Cache

public class GlobalSettingsCache {
    private final GlobalSettingsClient globalSettingsClient;
    private final Cache<String, GlobalSettings> cache;

    public GlobalSettingsCache(GlobalSettingsClient globalSettingsClient, Duration cacheTtl) {
        this.globalSettingsClient = globalSettingsClient;
        this.cache = CacheBuilder.newBuilder()
                .expireAfterWrite(cacheTtl)
                .build();
    }

    public Mono<GlobalSettings> get(String tenant) {
        return CacheMono.lookup(key -> Mono.justOrEmpty(cache.getIfPresent(key)).map(Signal::next), tenant)
                .onCacheMissResume(() -> fetchGlobalSettings(tenant))
                .andWriteWith((key, signal) -> Mono.fromRunnable(() ->
                        Optional.ofNullable(signal.get())
                                .ifPresent(value -> cache.put(key, value))));
    }

    private Mono<GlobalSettings> fetchGlobalSettings(String tenant) {
        return globalSettingsClient.getGlobalSettings(tenant);
    }
}
Related