How to pause consumption until a flux of messages has been processed

Viewed 183

I am looking for a way to pause to the consumption from a Kafka partition until a flux of messages has been processed.

For example

kafkaReceiver.receive() // I would like to pause this consumption until the current flux is processed
.flatmapIterable(records -> process(records))
.concatMap(record -> commit(record))

Currently I am thinking of adding a delay or adding onBackpressureBuffer to slow down the consumption.

Any advice or pointers would be very helpful.

0 Answers
Related