I have this scenario chalenging me: Kafka topic with messages that I must consume and expose via SSE to a Webcomponent. I have been done several tentatives in all layers looking for a stable and more reliable approach or at least some I fill more confortable to support.
Now I create a very simple Kafka topic and create two different consumers and both seems to work pretty the same way. One is using org.springframework.kafka.annotation.KafkaListener and other reactor.kafka.receiver.KafkaReceiver.
The final goal is expose via via Spring WebFlux an event posting the message when consumer from topic.
If I am not wrong I have read somewhere that Spring..KafkaListener is blocking code but as far as I can see it isn't. It is just a listener been triggered exactlly as Reactor..KafkaReceiver. Since I am coding a non-blocking code I should avoid blocking code but I can't Spring..KafkaListener blocking anywhere.
Here are the basic comparision resulting in exact same result:
Reactor KafkaReceiver:
ReceiverOptions<Object, Object> consumerOptions = ReceiverOptions.create(consumerProps)
.subscription(Collections.singleton("test"))
.addAssignListener(partitions -> logger.debug("onPartitionsAssigned {}", partitions))
.addRevokeListener(partitions -> logger.debug("onPartitionsRevoked {}", partitions));
kafkaReceiver = KafkaReceiver.create(consumerOptions);
((Flux<ReceiverRecord>) kafkaReceiver.receive()).doOnNext(r -> {
logger.info(String.format("Consumed Message using KafkaListener -> %s", r.value()));
r.receiverOffset().acknowledge();
}).subscribe();
Spring Kafka Listener:
@KafkaListener(topics = "test")
public void consume(String message) {
logger.info(String.format("Consumed Message using KafkaListener -> %s", message));
}
If you want to reproduce the comparision:
1 - clone https://github.com/jimisdrpc/simplest-comparison-kafkaconsumer
2 - create default Kafka topic name test
3 - produce any message to such topic
4 - enable either Reactor Kafka and run again after enable Kafka Listener (I didn't find how make both work in same application at same moment but let's say it isn't that important now)
PS.: I am not asking which is better. I am trying to compare for a very specific scenario (reactive programming aimed to produce events for a webcomponent via Spring WebFlux) and trying to garante that I am not choosing one that might be blocking code somehow.