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.