Build a reactive API server messages from Kafka in server sent event manner with Spring webflux

Viewed 345

I'm working on building a reactive web service that's connected to a Kafka topic, and it will support a reactive API that all requestor calling this API will receive new messages from Kafka as server sent events. (Originally thinking about the API will support keywords and will only receive kafka messages containing this keywords but figured out that Kafka doesn't support that). I'm pretty new to reactive and spring webflux, my current plan is to subscribe to Kafka, send the message from Kafka to a sink, then the controller returns sink.asFlux. But obviously, my understanding of the sink was not correct. Any suggestion would be very helpful :)

Controller

@RestController

@RequestMapping(path = "/sse")

public class SSEController {

    
    @Autowired
    private Flux<String> flux;

 
   
    @CrossOrigin(allowedHeaders = "*")
    
    @GetMapping(value = "/listen/{keyword}", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    
    public Flux<String> listenToKeyword(@PathVariable String keyword) {

        // TODO: how to filter out by keyword as data stream

        return flux;

    }



    @CrossOrigin(allowedHeaders = "*")

    @GetMapping(value = "/listenall", produces = MediaType.TEXT_EVENT_STREAM_VALUE)

    public Flux<String> listenToAll() {

        return flux;

    }
}

Configuration

@Configuration

public class KafkaConfiguration {

    Logger logger = LoggerFactory.getLogger(getClass());



    @Value("${spring.kafka.bootstrap-servers}")

    private String bootStrapServers;



    @Value("${spring.kafka.consumer.topic}")

    private String kafkaTopic;
    


    @Bean

    public Sinks.Many<String> sink() {

        return Sinks.many().replay().latest();

    }
  
  
    
    @Bean

    public Flux<String> flux(Sinks.Many<String> sink) {

        return sink.asFlux();

    }



    @Bean(name = "kafkaReceiver")

    public KafkaReceiver kafkaReceiver() {


        Map<String, Object> props = new HashMap<>();
        
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootStrapServers);
        
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());

        props.put(

                ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());

        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");

        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true);


        try {

            ReceiverOptions<String, String> receiverOptions = ReceiverOptions.<String, String>create(props).commitInterval(Duration.ZERO).commitBatchSize(0).subscription(Collections.singleton(kafkaTopic));

            KafkaReceiver<String, String> kafkaReceiver = KafkaReceiver.create(receiverOptions).receive().publishOn(Schedulers.newSingle("sample", true)).concatMap(m -> this.sink()::emitNext); // TODO: how to emit to the sink

            return kafkaReceiver;

        } catch (Exception e) {

            logger.debug("", e);

            throw e;

        }

    }

}

0 Answers
Related