【问题标题】:Spring Websocket integration with KafkaSpring Websocket 与 Kafka 的集成
【发布时间】:2018-03-07 23:10:57
【问题描述】:

我正在尝试通过 Spring MVC 项目中的 Spring-Websockets 将使用的 Kafka 数据发送到前端 (JavaScript)。

为了建立服务器和客户端之间的通信,我有以下内容。

客户端 (app.js)

function connect() {
    var socket = new SockJS('/kafka-data-websocket');
    stompClient = Stomp.over(socket);
    stompClient.connect({}, function (frame) {
        console.log('Connected: ' + frame);
        stompClient.send("/app/fetchData");
        stompClient.subscribe('/data/records', function (message) {
            console.log(JSON.parse(message.body).content);
        });
    });
}

服务器 (KafkaController.java)

@Controller
public class KafkaController {

    @MessageMapping("/fetchData")
    @SendTo("/data/records")
    public String fetchMetrics() {
        //...
    }
}

要使用来自特定 Kafka 主题的数据,我正在使用 @KafkaListener 注释,如下所示:

public class KafkaReceiver {
    @KafkaListener(topics = "mytopic")
    public void receive(ConsumerRecord<?, ?> record) {
        MyRecord m = new MyRecord(new Long(record.offset()), record.key().toString(), record.value().toString());
           //...
    }
}

我有一个适当的 KafkaConfig 类,其中包含所有必要的 bean (like explained here)。

如何将数据从 receive 方法发送到 KafkaController 的 fetchMetrics(进而发送到 websocket)上的每条传入/消费消息?

【问题讨论】:

  • 你有解决办法吗?

标签: javascript java spring spring-websocket spring-kafka


【解决方案1】:

您应该将SimpMessagingTemplate 注入KafkaReceiver 并从receive() 方法中使用它:

 this.template.convertAndSend("/data/records", m);

在 Spring Framework Reference Manual 中查看更多信息。

【讨论】:

  • 嘿!非常感谢您的回答。我现在收到此错误:SLF4J: Failed toString() invocation on an object of type [org.apache.kafka.clients.NodeApiVersions] java.lang.NullPointerException at org.apache.kafka.clients.NodeApiVersions.apiVersionToText(NodeApiVersions.java:167) 有什么想法吗?
猜你喜欢
  • 2023-03-31
  • 2017-06-07
  • 2019-03-03
  • 2018-06-21
  • 2016-07-09
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-05-22
相关资源
最近更新 更多