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;
}
}
}