How would one implement a doAfterNext operator in Project Reactor?

Viewed 50

RxJava2 has a doAfterNext operator that emits items downstream, and then invokes the consumer. It doesn't seem like Project Reactor has such an operator so I'd like to get some pointers on the best way to create my own to achieve the same thing.

The use case is freeing memory after the subscriber has received the item

1 Answers

Not sure if leavering doOnEach is a valid solution:

public class ByteBufferSafeReleaseConsumer implements Consumer<Signal<ByteBuffer<?>>> {

    private final List<ByteBuffer<?>> elements = new ArrayList<>();

    @Override
    public void accept(Signal<ByteBuffer<?>> signal) {
        if (signal.isOnNext()) {
            ByteBuffer<?> next = signal.get();
            if (next != null) {
                elements.add(next);
            }
        }
        if (signal.isOnComplete() || signal.isOnError()) {
            for (ByteBuffer<?> buffer : elements) {
                ByteBufferUtils.safeRelease(buffer);
            }
        }
    }
}

ByteBufferSafeReleaseConsumer consumer = new ByteBufferSafeReleaseConsumer()
Flux.from(byteBufferPublisher).doOnEach(consumer)
Related