I'm new to reactive programming and trying to fetch records from Flux if Consumer records and forward that record for further processing.
@Slf4j
public class OMSKafkaConsumer<K, V> {
private final KafkaReceiver<K, V> receiver;
public OMSKafkaConsumer(KafkaConsumerConfig<K, V> consumerConfig) {
final ReceiverOptions<K, V> receiverOptions = ReceiverOptions.<K, V>create(consumerConfig.consumerConfigs())
.subscription(Collections.singleton(consumerConfig.getTopic()));
receiver = KafkaReceiver.create(receiverOptions);
}
public Flux<ConsumerRecord<K, V>> receiveMessage() {
final Flux<ConsumerRecord<K, V>> inboundFlux = receiver.receiveAtmostOnce();
inboundFlux.subscribe(r -> log.info("Received message: %s \n", r));
return inboundFlux;
}
}
This is how I'm trying to access the record received from Kafka in my consumer class but I'm getting an error at fromIterable as it only expects list. Not sure how I can loop for consumer records and process them in a reactive way.
return kafkaConsumer.receiveMessage()
.map(consumerRecords -> Flux.fromIterable(consumerRecords)
.doOnNext(() -> {
for (ConsumerRecord<K, V> r : consumerRecords)
processMessage(r.getValue());
}));