【发布时间】:2015-01-26 22:22:02
【问题描述】:
我创建了一个MqttClient 类型的client,如下面的代码所示,我创建了一个客户端并设置了它的Asynchronous callback。问题是,
1-当我运行程序时,出现System.out.println("Client is Connected");,但我没有收到来自onSuccess 或来自oonFailure 的响应,为什么?我在代码中做错了什么。
2-我实现了static IMqttAsyncClient asynchClientCB = new IMqttAsyncClient() 接口,但是由于我有一个MqttClient 类型的客户端,我不能使用这个IMqttAsyncClient 接口。我尝试使用mqttAsynchClien,但因为我为java而不是Android编程,所以我不能使用它。 IMqttAsyncClient接口如何使用?
Update_1
在下面的代码“Updated_code_1”中,我稍微修改了代码,但我希望每次成功连接到broker 时,都会打印onSuccess 同步回调中的消息,并且onFailure 同步回调中的消息在连接终止的情况下打印callbck,例如当我故意断开网络连接时。但是当我连接到broker 时,onSuccess 和onFailur 都没有显示任何内容。那么,它们的设计目的是什么?
*Update_2_17_Dec_2014
我有一个询问可能会引导我们找到解决方案,也就是说,我通过有线/无线网络连接到代理是否重要?这会改变同步和异步监听器的行为吗?
Updated_1_code:
MqttConnectOptions opts = getClientOptions();
client = MQTTClientFactory.newClient(broker, port, clientID);
if (client != null) {
System.out.println("Client is not Null");
client.setCallback(AsynchCallBack);
if (opts != null) {
iMQTTToken = client.connectWithResult(opts);
publishMSG(client, TOPIC,"010101".getBytes(), QoS, pub_isRetained);
iMQTTToken.setActionCallback(synchCallBack);
if (client.isConnected()) {
System.out.println("Client CONNECTED.");
publishMSG(client, TOPIC,"010101".getBytes(), QoS, pub_isRetained);
}
}
}
....
....
....
....
IMqttToken iMQTTToken = new IMqttToken() {
@Override
public void waitForCompletion(long arg0) throws MqttException {
// TODO Auto-generated method stub
System.out.println("@waitForCompletion(): waiting " + (arg0 * 1000) + " seconds for connection to be established.");
}
@Override
public void waitForCompletion() throws MqttException {
// TODO Auto-generated method stub
System.out.println("@waitForCompletion(): waiting for connection to be established.");
}
@Override
public void setUserContext(Object arg0) {
// TODO Auto-generated method stub
}
@Override
public void setActionCallback(IMqttActionListener arg0) {
// TODO Auto-generated method stub
arg0.onSuccess(iMQTTToken);
//System.out.println(" " + arg0.onSuccess());
//System.out.println(" " + arg0.onSuccess(iMQTTToken));
iMQTTToken.setActionCallback(synchCallBack);
}
@Override
public boolean isComplete() {
// TODO Auto-generated method stub
return false;
}
@Override
public Object getUserContext() {
// TODO Auto-generated method stub
return null;
}
@Override
public String[] getTopics() {
// TODO Auto-generated method stub
return null;
}
@Override
public boolean getSessionPresent() {
// TODO Auto-generated method stub
return false;
}
@Override
public MqttWireMessage getResponse() {
// TODO Auto-generated method stub
return null;
}
@Override
public int getMessageId() {
// TODO Auto-generated method stub
return 0;
}
@Override
public int[] getGrantedQos() {
// TODO Auto-generated method stub
return null;
}
@Override
public MqttException getException() {
// TODO Auto-generated method stub
return null;
}
@Override
public IMqttAsyncClient getClient() {
// TODO Auto-generated method stub
return null;
}
@Override
public IMqttActionListener getActionCallback() {
// TODO Auto-generated method stub
return null;
}
};
IMqttActionListener synchCallBack = new IMqttActionListener() {
@Override
public void onSuccess(IMqttToken arg0) {
// TODO Auto-generated method stub
System.out.println("@onSuccess: Connection Successful.");
}
@Override
public void onFailure(IMqttToken arg0, Throwable arg1) {
// TODO Auto-generated method stub
System.out.println("@onFailure: Connection Failed.");
setViewEnableState(Bconnect, true);
}
};
MqttCallback AsynchCallBack = new MqttCallback() {
@Override
public void messageArrived(String topic, MqttMessage msg) throws Exception {
// TODO Auto-generated method stub
System.out.println("@messageArrived: Message Delivered.");
}
@Override
public void deliveryComplete(IMqttDeliveryToken token) {
// TODO Auto-generated method stub
System.out.println("@deliveryComplete: Delivery Completed.");
}
@Override
public void connectionLost(Throwable thrw) {
// TODO Auto-generated method stub
System.out.println("@Connection Lost: Connection Lost.");
setViewEnableState(Bconnect, true);
}
};
新客户:
MqttConnectOptions opts = new MqttConnectOptions();
opts.setCleanSession(CS);
opts.setKeepAliveInterval(KATimer);
HashMap<Integer, WILL> LWTData = WILLFactory.newWILL("LWT", "LWT MS".getBytes(), 1, false);
opts.setWill(LWTData.get(0).getWILLTopic(),
LWTData.get(0).getWILLPayLoad(),
LWTData.get(0).getWILLQoS(),
LWTData.get(0).isWILLRetained());
client = MQTTClientFactory.newClient(IP, PORT, clientID);
if (client != null) {
System.out.println("client is not null");
client.setCallback(AsynchCB);
IMqttToken token = client.connectWithResult(opts);
if (client.isConnected()) {
System.out.println("Client is Connected");
token.setActionCallback(new IMqttActionListener() {
public void onSuccess(IMqttToken arg0) {
// TODO Auto-generated method stub
System.out.println("synchCB->@onSuccess(): Connection Successful");
try {
client.subscribe(TOPIC, QoS);
} catch (MqttException e) {
// TODO Auto-generated catch block
e.printStackTrace();
}
try {
client.disconnect();
} catch (MqttException e) {
// TODO Auto-generated catch block
e.printStackTrace();
}
}
public void onFailure(IMqttToken arg0, Throwable arg1) {
// TODO Auto-generated method stub
System.out.println("synchCB->@onFailure(): Connection Failed");
}
});
}else {
System.out.println("client is not connected");
}
}else {
System.out.println("client = null");
}
异步回调:
/**
* Asynchronous Callback to inform the user about events that might happens Asynchronously. If it is not used, any pending
* messages destined to the client would not be received.
*/
private static MqttCallback AsynchCB = new MqttCallback() {
public void messageArrived(String topic, MqttMessage msg) throws Exception {
// TODO Auto-generated method stub
System.out.println("AsynchCB->@messageArrived(): ");
System.out.println("Topic: " + topic);
System.out.println("MSG: " + msg.toString());
}
public void deliveryComplete(IMqttDeliveryToken arg0) {
// TODO Auto-generated method stub
System.out.println("AsynchCB->@deliveryComplete(): ");
}
public void connectionLost(Throwable arg0) {
// TODO Auto-generated method stub
System.out.println("AsynchCB->@connectionLost(): ");
}
};
【问题讨论】:
-
如果您从 onSuccess() 中取出订阅并在检查连接成功后立即添加它会发生什么?我的意思是把这条线 client.subscribe(TOPIC, QoS);在检查 isConnected() 返回 true 之后。我有点困惑,为什么除了设置回调侦听器之外,您还希望在实际订阅或对连接进行任何操作之前调用 onSuccess()。
-
@kha 谢谢你的评论。实际上,在阅读了您的 cmets 之后,我似乎误解了 onSuccess 和 onFailure 的作用。因为我认为,onSuccess() 和 onFailure 是同步回调,当连接成功时调用“onSuccess”或失败“onFailure”,这就是为什么我在 onSucess() 中订阅,认为当连接建立/成功然后订阅。我是对还是错?请指导
-
您的连接已经成功并建立。您已经在这一行中正确地检查了它: if (client.isConnected()) ... 由于您有一个工作连接,您应该可以订阅您的主题。试一试,看看主题订阅是否有效。如果是这样,您应该能够开始接收发布在这些主题上的消息。
-
@kha [您的连接已经成功并建立。您已经在这一行中正确检查了它: if (client.isConnected())]===> yes, 好的,我检查正确,但是,如果网络连接发生意外断开并重新连接,我没有收到任何来自“onSuccess”或“onFailure”的响应。 AFAIU,上面提到的同步回调旨在报告此类动作“连接/断开连接”。请参阅上面的更新
-
与网络的连接类型无关紧要,我在 GridGain 网络中工作过,所有节点都有无线网络连接,所以即使进程是异步的,这也无所谓。跨度>
标签: java mqtt messagebroker mosquitto paho