Deal efficiently with a lot of discarded objects

Viewed 64

I have the following reactive stream where logs is Flux<String>:

logs.bufferTimeout(50, Duration.ofSeconds(20))
    .doOnNext(logs -> sendLogs(logs))
    .doOnDiscard(String.class, log -> sendLogs(List.of(log)));

The purpose of the stream is to send logs to another microservice for long-term storage. As you can see, the logs are buffered and sent in batches of 50. The intent is to limit the number of requests made to the second microservice.

For reasons that are not relevant to this question, the logs stream can be canceled at any time. Whenever this happens, I want all log messages that are still in the buffer to also be sent to the second microservice. The doOnDiscard operator I've added does just that, but perhaps not in the most performant/safe way. Let's say that there are 49 logs in the buffer at the time of cancellation. The lambda in doOnDiscard will then be called 49 times (once for each discarded log message) and this will result in 49 requests to the second microservice is a short amount of time.

You can reproduce this behavior by executing the following code:

    public static void main(String[] args) throws InterruptedException {
        Disposable disposable = Flux.interval(Duration.ofSeconds(1))
                .bufferTimeout(10, Duration.ofSeconds(5))
                .doOnNext(System.out::println)
                .doOnDiscard(Long.class, i -> System.out.println("Discarded: " + i))
                .subscribe();
        Thread.sleep(8000);
        disposable.dispose();
        Thread.sleep(2000);
    }

Is there any way to buffer the discarded elements and process them together? I've tried the following:

    public static void main(String[] args) throws InterruptedException {
        List<Long> discardedElements = new ArrayList<>();
        Disposable disposable = Flux.interval(Duration.ofSeconds(1))
                .bufferTimeout(10, Duration.ofSeconds(5))
                .doOnNext(System.out::println)
                .doOnDiscard(Long.class, discardedElements::add)
                .doFinally(signal -> System.out.println("Discarded elements: " + discardedElements))
                .subscribe();
        Thread.sleep(8000);
        disposable.dispose();
        Thread.sleep(2000);
    }

It does what I want, but I feel like using a variable outside of the stream for temporary storage is not the "cleanest" solution. Is there a better one? Is the pattern of using an external variable a "recommended" one?

0 Answers
Related