使用 Java 开发 MQTT 客户端

本章节以简单的例子讲解如何写一个 Java MQTT 客户端,社区中有几个 Java MQTT 库,包括有,

  • Paho: Paho是 Eclipse 的一个开源 MQTT 项目,包含多种语言实现,Java 是其中之一
  • Fusesource:Fusesource是另外一个开源的 Java MQTT 库,但是该项目已经不活跃,大约有 2 年时间没有更新

此处以社区中比较活跃的 java-paho 为例来说明使用 Java 开发 MQTT 客户端的过程。

安装 java-paho

本例中使用 Maven 来管理依赖的库文件,打开 pom.xml,加入以下的 JAR 依赖,等待完成相关 JAR 包的下载。

  1. <dependency>
  2. <groupId>org.eclipse.paho</groupId>
  3. <artifactId>org.eclipse.paho.client.mqttv3</artifactId>
  4. <version>1.2.0</version>
  5. </dependency>

实现一个简单的客户端

这部分实现的例子比较简单,利用 paho 提供的 MqttClient,实现了以下的功能,

  • EMQ X 服务器的连接
  • 连接建立成功后,订阅主题 demo/topics,并且设置了回调的 Listener,当有消息转发到该主题的时候就调用此方法
  • 调用 publish 方法来实现对主题 demo/topics 的消息发送

初始化和建立连接

构造 MqttClient 的构造函数如下,

  1. MqttClient(String serverURI, String clientId, MqttClientPersistence persistence)
  • serverURI:EMQ X 的服务器地址,在此处为 tcp://localhost:1883
  • clientId:标识该客户端的唯一 ID,此 ID 在同一个 EMQ X 服务器中必须保证唯一,否则在服务器端在处理 session 的时候会有问题
  • MqttClientPersistence:本地消息的持久化实例,在本地消息处理过程在涉及到服务器端忙碌或者不可用等状态的时候,需要对消息进行持久化的处理,在这里可以传入持久化处理的类实例。

建立连接的代码如下所示,

  1. package paho_demo;
  2. import org.eclipse.paho.client.mqttv3.MqttClient;
  3. import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
  4. import org.eclipse.paho.client.mqttv3.MqttException;
  5. import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
  6. public class Demo {
  7. public static void main(String[] args) {
  8. String broker = "tcp://localhost:1883";
  9. String clientId = "JavaSample";
  10. //Use the memory persistence
  11. MemoryPersistence persistence = new MemoryPersistence();
  12. try {
  13. MqttClient sampleClient = new MqttClient(broker, clientId, persistence);
  14. MqttConnectOptions connOpts = new MqttConnectOptions();
  15. connOpts.setCleanSession(true);
  16. System.out.println("Connecting to broker:" + broker);
  17. sampleClient.connect(connOpts);
  18. System.out.println("Connected");
  19. } catch (MqttException me) {
  20. System.out.println("reason" + me.getReasonCode());
  21. System.out.println("msg" + me.getMessage());
  22. System.out.println("loc" + me.getLocalizedMessage());
  23. System.out.println("cause" + me.getCause());
  24. System.out.println("excep" + me);
  25. me.printStackTrace();
  26. }
  27. }
  28. }

执行该段代码后,如果能够成功连上服务器,在控制台中将打印出如下内容。如果发生异常,请根据异常信息来定位和修复问题。

  1. Connecting to broker: tcp://localhost:1883
  2. Connected

订阅消息

连接建立成功之后,可以进行主题订阅。MqttClient 提供了多个 subscribe 方法,可以实现不同方式的主题订阅。主题可以是明确的单个主题,也可以用通配符 # 或者 +

  1. subscribe(java.lang.String topicFilter)

订阅主题后,设置一个回调实例 MqttCallback,在消息转发过来的时候将调用该实例的方法。消息订阅部分的代码如下所示。

  1. String topic = "demo/topics";
  2. System.out.println("Subscribe to topic:" + topic);
  3. sampleClient.subscribe(topic);
  4. sampleClient.setCallback(new MqttCallback() {
  5. public void messageArrived(String topic, MqttMessage message) throws Exception {
  6. String theMsg = MessageFormat.format("{0} is arrived for topic {1}.", new String(message.getPayload()), topic);
  7. System.out.println(theMsg);
  8. }
  9. public void deliveryComplete(IMqttDeliveryToken token) {
  10. }
  11. public void connectionLost(Throwable throwable) {
  12. }
  13. });

在消息转发成功后,控制台会打印转发的消息以及针对的主题。

发布消息

MqttClientpublish 方法用于发布消息

  1. publish(java.lang.String topic, MqttMessage message)
  • topic:主题名称
  • MqttMessage:消息内容

MqttClient 还提供了以下的方法,用户可以在发布消息的时候指定 QoS,以及消息是否需要保持。

  1. publish(java.lang.String topic, byte[] payload, int qos, boolean retained)

发布消息的代码如下所示,

  1. String topic = "demo/topics";
  2. String content = "Message from MqttPublishSample";
  3. int qos = 2;
  4. System.out.println("Publishing message:" + content);
  5. MqttMessage message = new MqttMessage(content.getBytes());
  6. message.setQos(qos);
  7. sampleClient.publish(topic, message);
  8. System.out.println("Message published");

完整例程

  1. package paho_demo;
  2. import java.text.MessageFormat;
  3. import org.eclipse.paho.client.mqttv3.IMqttDeliveryToken;
  4. import org.eclipse.paho.client.mqttv3.MqttCallback;
  5. import org.eclipse.paho.client.mqttv3.MqttClient;
  6. import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
  7. import org.eclipse.paho.client.mqttv3.MqttException;
  8. import org.eclipse.paho.client.mqttv3.MqttMessage;
  9. import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
  10. public class Demo {
  11. public static void main(String[] args) {
  12. String broker = "tcp://localhost:1883";
  13. String clientId = "JavaSample";
  14. //Use the memory persistence
  15. MemoryPersistence persistence = new MemoryPersistence();
  16. try {
  17. MqttClient sampleClient = new MqttClient(broker, clientId, persistence);
  18. MqttConnectOptions connOpts = new MqttConnectOptions();
  19. connOpts.setCleanSession(true);
  20. System.out.println("Connecting to broker:" + broker);
  21. sampleClient.connect(connOpts);
  22. System.out.println("Connected");
  23. String topic = "demo/topics";
  24. System.out.println("Subscribe to topic:" + topic);
  25. sampleClient.subscribe(topic);
  26. sampleClient.setCallback(new MqttCallback() {
  27. public void messageArrived(String topic, MqttMessage message) throws Exception {
  28. String theMsg = MessageFormat.format("{0} is arrived for topic {1}.", new String(message.getPayload()), topic);
  29. System.out.println(theMsg);
  30. }
  31. public void deliveryComplete(IMqttDeliveryToken token) {
  32. }
  33. public void connectionLost(Throwable throwable) {
  34. }
  35. });
  36. String content = "Message from MqttPublishSample";
  37. int qos = 2;
  38. System.out.println("Publishing message:" + content);
  39. MqttMessage message = new MqttMessage(content.getBytes());
  40. message.setQos(qos);
  41. sampleClient.publish(topic, message);
  42. System.out.println("Message published");
  43. } catch (MqttException me) {
  44. System.out.println("reason" + me.getReasonCode());
  45. System.out.println("msg" + me.getMessage());
  46. System.out.println("loc" + me.getLocalizedMessage());
  47. System.out.println("cause" + me.getCause());
  48. System.out.println("excep" + me);
  49. me.printStackTrace();
  50. }
  51. }
  52. }

运行结果

  1. Connecting to broker: tcp://localhost:1883
  2. Connected
  3. Subscribe to topic: demo/topics
  4. Publishing message: Message from MqttPublishSample
  5. Message published
  6. Message from MqttPublishSample is arrived for topic demo/topics.