【发布时间】:2018-05-05 05:03:15
【问题描述】:
我正在编写一个从 MQTT 代理接收消息的程序 服务器从客户端获取ID并使其成为主题名称:例如topic1,topic2。 然后在订阅时,服务器传递主题的名称,然后从该主题读取消息。 这是我的服务器:
public class AnalyticServer {
// The server socket.
private static ServerSocket serverSocket = null;
// The client socket.
private static Socket clientSocket = null;
// This server can accept up to maxClientsCount clients' connections.
private static final int maxClientsCount = 5;
private static final clientThread[] threads = new clientThread[maxClientsCount];
public static void main(String args[]) throws MqttException, InterruptedException {
// The default port number.
int portNumber = 4544;
//Open Server
try {
serverSocket = new ServerSocket(portNumber);
} catch (IOException e) {
System.out.println(e);
}
//
//When server is listening
System.out.println("Server is now listening at port 4544");
while (true) {
try {
//Make connection
clientSocket = serverSocket.accept();
System.out.println("Connected");
int i = 0;
//Find thread null to run the connection
for (i = 0; i < maxClientsCount; i++) {
if (threads[i] == null) {
(threads[i] = new clientThread(clientSocket, threads)).start();
break;
}
}
if (i == maxClientsCount) {
PrintStream os = new PrintStream(clientSocket.getOutputStream());
os.println("Server is now, please try again later");
os.close();
clientSocket.close();
}
} catch (IOException e) {
System.out.println(e);
}
}
}
}
//Thread control each Request
class clientThread extends Thread {
private Socket clientSocket = null;
private final clientThread[] threads;
private int maxClientsCount;
public clientThread(Socket clientSocket, clientThread[] threads) {
this.clientSocket = clientSocket;
this.threads = threads;
maxClientsCount = threads.length;
}
public void run() {
int maxClientsCount = this.maxClientsCount;
clientThread[] threads = this.threads;
try {
int identifier=0;
//get id
InputStream input = null;
input = clientSocket.getInputStream();
identifier=input.read();
String topic="topic".concat(String.valueOf(identifier));
//Subscribe
try {
System.out.println("subscribing");
Subscribe receive=new Subscribe(topic);
} catch (MqttException e) {
// TODO Auto-generated catch block
e.printStackTrace();
} catch (URISyntaxException e) {
// TODO Auto-generated catch block
e.printStackTrace();
}
// open image
FileInputStream imgPath = new FileInputStream("image.jpg");
BufferedImage bufferedImage = ImageIO.read(imgPath);
Thread.sleep(1200);
ByteArrayOutputStream baos = new ByteArrayOutputStream();
ImageIO.write( bufferedImage, "jpg", baos );
baos.flush();
byte[] imageInByte = baos.toByteArray();
baos.close();
//SendImage
DataOutputStream outToClient = new DataOutputStream(clientSocket.getOutputStream());
outToClient.write(imageInByte);
System.out.println(outToClient.size());
clientSocket.close();
} catch (IOException e) {
} catch (InterruptedException e) {
// TODO Auto-generated catch block
e.printStackTrace();
}
}
}
'这是我的订阅课程
public class Subscribe implements MqttCallback {
private final int qos = 1;
static String topic=null;
private MqttClient client;
String subText = "abc";
public Subscribe(String topic) throws MqttException, URISyntaxException {
this.topic=topic;
String host = "tcp://m14.cloudmqtt.com:19484";
String username = "***";
String password = "********";
String clientId = MqttClient.generateClientId();
MqttConnectOptions conOpt = new MqttConnectOptions();
conOpt.setCleanSession(true);
conOpt.setUserName(username);
conOpt.setPassword(password.toCharArray());
this.client = new MqttClient(host, clientId, new MemoryPersistence());
;
this.client.setCallback(this);
this.client.connect(conOpt);
this.client.subscribe(topic,1);
System.out.println("subscribe topic: " +this.topic);
}
/**
* @see MqttCallback#connectionLost(Throwable)
*/
public void connectionLost(Throwable cause) {
System.out.println("Connection lost because: " + cause);
System.exit(1);
}
/**
* @see MqttCallback#deliveryComplete(IMqttDeliveryToken)
*/
public void deliveryComplete(IMqttDeliveryToken token) {
}
/**
* @throws IOException
* @see MqttCallback#messageArrived(String, MqttMessage)
*/
public void messageArrived(String topic, MqttMessage message) throws MqttException, IOException {
System.out.println("1");
subText = message.getPayload().toString();
System.out.println("Received"+subText);
}
}
订阅构造函数中的主题是正确的,但是 this.client.setCallback(this) 似乎没有调用方法 messageArrived。所以我什么都收不到。
有人知道吗? 非常感谢
【问题讨论】:
-
您已在来源中包含您的用户名和密码,最好的办法是删除问题并隐藏该信息再次询问
-
哦,我忘了,非常感谢。
-
@ngoc-anh 该信息仍在编辑历史记录中。完全删除问题并再次提问。同时更改该用户的密码
标签: java tcp cloud mqtt broker