Spring WebFlux 与 Kafka 和 Websockets

Yas*_*lan 5 spring apache-kafka spring-websocket spring-webflux

现在,我在我的 SpringBoot 应用程序中实现了一个简单的 Kafka Consumer 和 Producer,它工作得很好,接下来我想做的是,我的消费者获取消费的消息并将其直接广播给所有订阅的客户端。我发现我不能将 STOMP 消息传递与 WebFlux 一起使用,那么我该如何完成这个任务,我看到了反应式 WebSocket 实现,但我不知道如何将我使用的数据发送到我的 websocket。

这是我的简单 KafkaProducer:

fun addMessage(message: Message){
        val headers : MutableMap<String, Any> = HashMap()
        headers[KafkaHeaders.TOPIC] = topicName
        kafkaTemplate.send(GenericMessage<Message>(message, headers))
    }
Run Code Online (Sandbox Code Playgroud)

我的简单消费者看起来像这样:

@KafkaListener(topics = ["mytopic"], groupId = "test-consumer-group")
    fun receiveData(message:Message) :Message{
        //Take consumed data and send to websocket
    }
Run Code Online (Sandbox Code Playgroud)

Art*_*lan 4

我会考虑有一个Sinks.many().multicast().onBackpressureBuffer()作为全球中间容器。然后,receiveData()您只需数据放入 Reactor 抽象中即可。

对于您的 WebSocket 连接会话,我建议在API 中实现org.springframework.web.reactive.socket.WebSocketHandler并使用。这样,只要连接到此 WebSocket 服务器,所有会话都将使用相同的 Kafka 数据。Sinks.Many.asFlux()WebSocketSession.send(Publisher<WebSocketMessage> messages)

在文档中查看更多信息:https://docs.spring.io/spring-framework/docs/current/reference/html/web-reactive.html#webflux-websockethandler

更新

您可以在这里找到一些示例:https ://github.com/artembilan/sandbox/tree/master/so-65667450