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)
我会考虑有一个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 ://github.com/artembilan/sandbox/tree/master/so-65667450
| 归档时间: |
|
| 查看次数: |
2956 次 |
| 最近记录: |