首页

java 连接阿里云的mqtt服务(客户端源码)

java

2020-6-29

/**
 * aliyun.com Inc.
 * Copyright (c) 2004-2017 All Rights Reserved.
 */
package com.aliyun.iot.demo.iothub;

import java.net.InetAddress;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.ThreadPoolExecutor.CallerRunsPolicy;
import java.util.concurrent.TimeUnit;

import javax.net.ssl.SSLContext;
import javax.net.ssl.SSLSocketFactory;
import javax.net.ssl.TrustManager;

import org.eclipse.paho.client.mqttv3.IMqttDeliveryToken;
import org.eclipse.paho.client.mqttv3.IMqttMessageListener;
import org.eclipse.paho.client.mqttv3.MqttCallback;
import org.eclipse.paho.client.mqttv3.MqttClient;
import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;

import com.aliyun.iot.util.LogUtil;
import com.aliyun.iot.util.SignUtil;

/**
 * IoT套件JAVA版设备接入demo
 */
public class SimpleClient4IOT {
	
	/******这里是客户端需要的参数*******/
    public static String deviceName = "huakuangxxxxx";
    public static String productKey = "axxxxxxx";
    public static String secret = "6sllz6mmMeuZpri9Jnnhwxxxx";

    //用于测试的topic
    private static String subTopic = "/"   productKey   "/"   deviceName   "/get";
    private static String pubTopic = "/"   productKey   "/"   deviceName   "/update";

    public static void main(String... strings) throws Exception {
        //客户端设备自己的一个标记,建议是MAC或SN,不能为空,32字符内
        String clientId = InetAddress.getLocalHost().getHostAddress();

        //设备认证
        Map<String, String> params = new HashMap<String, String>();
        params.put("productKey", productKey); //这个是对应用户在控制台注册的 设备productkey
        params.put("deviceName", deviceName); //这个是对应用户在控制台注册的 设备name
        params.put("clientId", clientId);
        String t = System.currentTimeMillis()   "";
        params.put("timestamp", t);

        //MQTT服务器地址,TLS连接使用ssl开头
        String targetServer = "ssl://"   productKey   ".iot-as-mqtt.cn-shanghai.aliyuncs.com:1883";

        //客户端ID格式,两个||之间的内容为设备端自定义的标记,字符范围[0-9][a-z][A-Z]
        String mqttclientId = clientId   "|securemode=2,signmethod=hmacsha1,timestamp="   t   "|";
        String mqttUsername = deviceName   "&"   productKey; //mqtt用户名格式
        String mqttPassword = SignUtil.sign(params, secret, "hmacsha1"); //签名

        System.err.println("mqttclientId="   mqttclientId "&mqttPassword=" mqttPassword);

        connectMqtt(targetServer, mqttclientId, mqttUsername, mqttPassword, deviceName);
    }

    public static void connectMqtt(String url, String clientId, String mqttUsername,
                                   String mqttPassword, final String deviceName) throws Exception {
        MemoryPersistence persistence = new MemoryPersistence();
        SSLSocketFactory socketFactory = createSSLSocket();
        final MqttClient sampleClient = new MqttClient(url, clientId, persistence);
        MqttConnectOptions connOpts = new MqttConnectOptions();
        connOpts.setMqttVersion(4); // MQTT 3.1.1
        connOpts.setSocketFactory(socketFactory);

        //设置是否自动重连
        connOpts.setAutomaticReconnect(true);

        //如果是true,那么清理所有离线消息,即QoS1或者2的所有未接收内容
        connOpts.setCleanSession(false);

        connOpts.setUserName(mqttUsername);
        connOpts.setPassword(mqttPassword.toCharArray());
        connOpts.setKeepAliveInterval(65);

        LogUtil.print(clientId   "进行连接, 目的地: "   url);
        sampleClient.connect(connOpts);

        sampleClient.setCallback(new MqttCallback() {
            @Override
            public void connectionLost(Throwable cause) {
                LogUtil.print("连接失败,原因:"   cause);
                cause.printStackTrace();
            }

            @Override
            public void messageArrived(String topic, MqttMessage message) throws Exception {
                LogUtil.print("接收到消息,来至Topic ["   topic   "] , 内容是:["
                      new String(message.getPayload(), "UTF-8")   "],  ");
            }

            @Override
            public void deliveryComplete(IMqttDeliveryToken token) {
                //如果是QoS0的消息,token.resp是没有回复的
                LogUtil.print("消息发送成功! "   ((token == null || token.getResponse() == null) ? "null"
                    : token.getResponse().getKey()));
            }
        });
        LogUtil.print("连接成功:---");

        //这里测试发送一条消息
        String content = "{'content':'msg from :"   clientId   ","   System.currentTimeMillis()   "'}";

        MqttMessage message = new MqttMessage(content.getBytes("utf-8"));
        message.setQos(0);
        //System.out.println(System.currentTimeMillis()   "消息发布:---");
        sampleClient.publish(pubTopic, message);

        //一次订阅永久生效 
        //这个是第一种订阅topic方式,回调到统一的callback
        sampleClient.subscribe(subTopic);

        //这个是第二种订阅方式, 订阅某个topic,有独立的callback
        //sampleClient.subscribe(subTopic, new IMqttMessageListener() {
        //    @Override
        //    public void messageArrived(String topic, MqttMessage message) throws Exception {
        //
        //        LogUtil.print("收到消息:"   message   ",topic="   topic);
        //    }
        //});

        //回复RRPC响应
        final ExecutorService executorService = new ThreadPoolExecutor(2,
            4, 600, TimeUnit.SECONDS,
            new ArrayBlockingQueue<Runnable>(100), new CallerRunsPolicy());

        String reqTopic = "/sys/"   productKey   "/"   deviceName   "/rrpc/request/ ";
        sampleClient.subscribe(reqTopic, new IMqttMessageListener() {
            @Override
            public void messageArrived(String topic, MqttMessage message) throws Exception {
                LogUtil.print("收到请求:"   message   ", topic="   topic);
                String messageId = topic.substring(topic.lastIndexOf('/')   1);
                final String respTopic = "/sys/"   productKey   "/"   deviceName   "/rrpc/response/"   messageId;
                String content = "hello world";
                final MqttMessage response = new MqttMessage(content.getBytes());
                response.setQos(0); //RRPC只支持QoS0
                //不能在回调线程中调用publish,会阻塞线程,所以使用线程池
                executorService.submit(new Runnable() {
                    @Override
                    public void run() {
                        try {
                            sampleClient.publish(respTopic, response);
                            LogUtil.print("回复响应成功,topic="   respTopic);
                        } catch (Exception e) {
                            e.printStackTrace();
                        }
                    }
                });
            }
        });
    }

    private static SSLSocketFactory createSSLSocket() throws Exception {
        SSLContext context = SSLContext.getInstance("TLSV1.2");
        context.init(null, new TrustManager[] {new ALiyunIotX509TrustManager()}, null);
        SSLSocketFactory socketFactory = context.getSocketFactory();
        return socketFactory;
    }
}
资源下载此资源下载价格为3D币(VIP免费),请先
资源文件列表
.idea/.name , 10
.idea/compiler.xml , 632
.idea/encodings.xml , 172
.idea/libraries/Maven__com_alibaba_fastjson_1_2_28.xml , 514
.idea/libraries/Maven__org_eclipse_paho_org_eclipse_paho_client_mqttv3_1_1_0.xml , 681
.idea/misc.xml , 439
.idea/modules.xml , 260
.idea/workspace.xml , 35593
README.md , 255
mqttclient.iml , 957
pom.xml , 1241
src/main/java/com/aliyun/iot/demo/iothub/ALiyunIotX509TrustManager.java , 3301
src/main/java/com/aliyun/iot/demo/iothub/SimpleClient4IOT.java , 7880
src/main/java/com/aliyun/iot/demo/shadow/SimpleClient4Shadow.java , 14481
src/main/java/com/aliyun/iot/util/AliyunWebUtils.java , 13905
src/main/java/com/aliyun/iot/util/CipherUtils.java , 1050
src/main/java/com/aliyun/iot/util/DynamicByteBuffer.java , 9063
src/main/java/com/aliyun/iot/util/LogUtil.java , 1046
src/main/java/com/aliyun/iot/util/Md5.java , 13000
src/main/java/com/aliyun/iot/util/ProtocolUtils.java , 3983
src/main/java/com/aliyun/iot/util/SecureUtils.java , 3708
src/main/java/com/aliyun/iot/util/SignUtil.java , 3151
src/main/resources/root.crt , 1280
target/classes/com/aliyun/iot/demo/iothub/ALiyunIotX509TrustManager.class , 3373
target/classes/com/aliyun/iot/demo/iothub/SimpleClient4IOT$1.class , 2303
target/classes/com/aliyun/iot/demo/iothub/SimpleClient4IOT$2$1.class , 1581
target/classes/com/aliyun/iot/demo/iothub/SimpleClient4IOT$2.class , 2387
target/classes/com/aliyun/iot/demo/iothub/SimpleClient4IOT.class , 6002
target/classes/com/aliyun/iot/demo/shadow/SimpleClient4Shadow$1.class , 2297
target/classes/com/aliyun/iot/demo/shadow/SimpleClient4Shadow$2.class , 2702
target/classes/com/aliyun/iot/demo/shadow/SimpleClient4Shadow.class , 10324
target/classes/com/aliyun/iot/util/AliyunWebUtils$1.class , 775
target/classes/com/aliyun/iot/util/AliyunWebUtils$DefaultTrustManager.class , 1230
target/classes/com/aliyun/iot/util/AliyunWebUtils.class , 11495
target/classes/com/aliyun/iot/util/CipherUtils.class , 1490
target/classes/com/aliyun/iot/util/DynamicByteBuffer.class , 9893
target/classes/com/aliyun/iot/util/LogUtil.class , 1622
target/classes/com/aliyun/iot/util/Md5.class , 8362
target/classes/com/aliyun/iot/util/ProtocolUtils.class , 2632
target/classes/com/aliyun/iot/util/SecureUtils.class , 3603
target/classes/com/aliyun/iot/util/SignUtil.class , 3355
target/classes/root.crt , 1280
没有账号? 忘记密码?

社交账号快速登录