Clear in-flight elements in a stream when an upstream publisher is restarted in Spring Project Reactor?

Viewed 79

I have a publisher that executes a long-running and large query on MongoDB and returns the data in a Flux. Entities that are marked in the database as "processed" will be filtered out and the entities are then buffered and passed to a concatMap operator (so that all buffered ≤elements are processed before elements in the next buffer are processed). It looks something like this:

Flux<Entity> entitiesFromMongoDb = myMongoRepository.executeLargeQuery();
entitiesFromMongoDb.filter(entity -> !entity.isProcessed())
                   .buffer(10)
                   .concatMap(bufferedEntityList ->  
                                    Flux.fromIterable(bufferedEntityList)
                                        .flatMap(makeExternalCall)
                                        .then()));

Where makeExternalCall calls a third-party remote server and sets the entity to processed after the call has been made. In most cases this works fine, but when the remote server is really slow or has an error then makeExternalCall will retry (with exponential backoff) the operation to the remote server. In some cases it can take quite a while before all 10 external calls have been processed. In fact it can take so long that the myMongoRepository.executeLargeQuery() publisher is restarted and the query is executed again. Now we run into a problem that I'll try to describe here:

  1. Entity A is read from the database (i.e. it's returned in the flux produced by myMongoRepository.executeLargeQuery()). It's not yet marked as "processed" which means that entity.isProcessed() will return false and it'll be retained in the stream.
  2. The external server is really slow or down so that makeExternalCall is forced to retry the operation before entity A has been marked as "processed" in the DB.
  3. myMongoRepository.executeLargeQuery() is restarted and the query is executed again.
  4. Entity A is read from the database once more. But the problem is that there's already another instance of entity A in-flight since it has not yet been marked as "processed" by the previous call to myMongoRepository.executeLargeQuery().
  5. This means that makeExternalCall will be called twice for entity A, which is not optimal!

I could make an additional request to the DB and check the status of processed for each entity in the makeExternalCall method, but this will cause additional load (since an extra request is necessary for each entity) to the DB which is not optimal.

So my question is:

Is there a way to somehow "restart" the entire stream, and thus clear intermediary buffers (i.e. remove entity A that is in-flight from the ongoing stream) when the MongoDB query triggered by myMongoRepository.executeLargeQuery() is restarted/re-executed? Or is there a better way to handle this?

I'm using Spring Boot 2.2.4.RELEASE, project reactor 3.3.2.RELEASE and spring-boot-starter-data-mongodb-reactive 2.2.4.RELEASE.

1 Answers

Not sure If I understood the problem completely. But trying to answer as it sounds interesting.

As you need to be aware of the requests which are already being processed by the makeExternalCall, can you maintain a set / local cache which contains the entities which are being processed?

Set<Entity> inProgress = new HashSet<>(); 

Flux<Entity> entitiesFromMongoDb = myMongoRepository.executeLargeQuery();

entitiesFromMongoDb.filter(entity -> !entity.isProcessed())
                   .buffer(10)
                   .map(bufferedEntityList -> {  // remove the inprogress requests to avoid redundant processing
                        bufferedEntityList.removeIf(inProgress::contains);
                        return bufferedEntityList;
                   })
                   .concatMap(bufferedEntityList ->  
                                    inProgress.addAll(bufferedEntityList);
                                    Flux.fromIterable(bufferedEntityList)
                                        .flatMap(makeExternalCall) //assuming once processed, it emits the entity object
                                        .map(entity -> {   //clear growing set
                                            inProgress.remove(entity);
                                            return entity;
                                        })
                                        .then()));

This approach is not a good solution when you need to scale your application horizontally. In that case instead of maintaining a local cache, you could go for an external cache server like redis.

Related