【发布时间】:2020-04-02 13:03:57
【问题描述】:
我目前正在尝试通过 Mosquitto MQQT 服务器发布和接收 Protobuf 消息。我已成功将正确的内容发布到服务器。但是,当客户端收到它时,方法 parseFrom() 挂起并且永远不会返回。这是一个与this one 非常相似的问题,当通过从未关闭的 Socket 发送 Protobuf 消息时发生。
出版商:
MqttClient adapterClient = new MqttClient(broker, clientID);
SpecsMessage.Specs protoNotifyMessage = SpecsMessage.Specs.newBuilder()
.setNodeType("basic")
.setAddress(serverSocket.getInetAddress().getHostName())
.setPort(serverSocket.getLocalPort())
.build();
MqttMessage notifyMessage = new MqttMessage(protoNotifyMessage.toString().getBytes());
adapterClient.publish("availableNodes", notifyMessage);
订阅者:
public class TestController implements MqttCallback {
public void messageArrived(String topic, MqttMessage message){
System.out.println("New node connected");
System.out.println("Payload: \n" + new String(message.getPayload()));
SpecsMessage.Specs protoMessage = SpecsMessage.Specs.parseFrom(message.getPayload());
}
}
我找不到向 MQQT 服务器指定发送消息的正确方式的方法。
我也尝试使用 writeDelimitedFrom() 方法。
MqttClient adapterClient = new MqttClient(broker, clientID);
SpecsMessage.Specs protoNotifyMessage = SpecsMessage.Specs.newBuilder()
.setNodeType("basic")
.setAddress(serverSocket.getInetAddress().getHostName())
.setPort(serverSocket.getLocalPort())
.build();
ByteArrayOutputStream output = new ByteArrayOutputStream();
protoNotifyMessage.writeDelimitedTo(output);
MqttMessage notifyMessage = new MqttMessage(output.toByteArray());
adapterClient.publish("availableNodes", notifyMessage);
但是,消息未正确转换为 byte[],如下所示:
nodeType: "basic"
address: "0.0.0.0"
port: 43101
这就是我得到的:
basic0.0.0.0��
有没有办法让这个工作,要么通过更正发送方法,要么通过解决字节[]转换问题?
【问题讨论】:
-
我怀疑这是问题所在,但如果
parseFrom本身挂起,而不是提供数据的方法挂起,我会感到惊讶。不过,new String(payload)有效。如果将有效负载存储在变量中,然后在 println 和 parseFrom 调用中使用,它是否有效? -
感谢您的评论。我试过了,但 parseFrom() 仍然挂在我身上。但是,正如您所建议的,这确实表明这不是沟通问题。
-
您还应该小心使用
getBytes()和new String的默认字符集:如果发布者和订阅者JVM 使用不同的默认字符集,您会得到意想不到的结果。 -
我会记住这一点,再次感谢您的宝贵时间!
标签: java protocol-buffers mqtt mosquitto protobuf-java