当前位置: 首页 > article >正文

Android MQTT开发之 Hivemq MQTT Client

使用一个开源库:hivemq-mqtt-client,这是Java生态的一个MQTT客户端框架,需要Java 8,Android上使用的话问题不大,需要一些额外的配置,下面列出了相关的配置,尤其是 packagingOptions,不然编译不过,因为框架使用了Java8新增的语言特性,所以 minSdk 设置为24,即Android7.0,如果要兼容Android7.0以下系统,可以参考这份详细文档配置一下语法脱糖的SDK: Installation on Android

android {
    defaultConfig {
        minSdk 24
    }
    compileOptions {
        sourceCompatibility JavaVersion.VERSION_8
        targetCompatibility JavaVersion.VERSION_8
    }
    kotlinOptions {
        jvmTarget = '8'
    }
    packagingOptions {
        resources {
            excludes += ['META-INF/INDEX.LIST', 'META-INF/io.netty.versions.properties']
        }
    }
}

dependencies {
    implementation 'com.hivemq:hivemq-mqtt-client:1.3.3'
}

刚开始在自动连接这块花了好多时间,最后才发现是设置用户名和密码的地方不对,一定要在设置自动重连(初始化Client)的地方设置,而不是连接的时候!下面是一个简单的使用示例代码

MqttManager.kt

import android.util.Log
import com.hivemq.client.mqtt.datatypes.MqttQos
import com.hivemq.client.mqtt.lifecycle.MqttClientConnectedContext
import com.hivemq.client.mqtt.lifecycle.MqttClientConnectedListener
import com.hivemq.client.mqtt.lifecycle.MqttClientDisconnectedContext
import com.hivemq.client.mqtt.lifecycle.MqttClientDisconnectedListener
import com.hivemq.client.mqtt.mqtt5.Mqtt5AsyncClient
import com.hivemq.client.mqtt.mqtt5.Mqtt5Client
import com.hivemq.client.mqtt.mqtt5.message.connect.connack.Mqtt5ConnAckReasonCode
import com.hivemq.client.mqtt.mqtt5.message.publish.Mqtt5Publish
import com.hivemq.client.mqtt.mqtt5.message.subscribe.suback.Mqtt5SubAck
import java.util.UUID
import java.util.concurrent.CompletableFuture
import java.util.concurrent.Executors
import java.util.function.Consumer

open class MqttListener {
    open fun onConnected() {}
    open fun onDisconnected() {}
    open fun onSubscribed(vararg topics: String) {}
    open fun onReceiveMessage(topic: String, data: ByteArray) {}
    open fun onSendMessage(topic: String, data: ByteArray) {}
}

/*
文档
https://github.com/hivemq/hivemq-mqtt-client
https://hivemq.github.io/hivemq-mqtt-client/docs/installation/android/
*/
class MqttManager private constructor() : MqttClientConnectedListener, MqttClientDisconnectedListener {
    private val executor = Executors.newFixedThreadPool(Runtime.getRuntime().availableProcessors()) {
        Thread(it).apply { isDaemon = true }
    }

    private val mqttAsynClient: Mqtt5AsyncClient = Mqtt5Client.builder()
        .identifier(UUID.randomUUID().toString())
        .serverHost(SERVER_HOST)
        .serverPort(SERVER_PORT)
        .addConnectedListener(this)
        .addDisconnectedListener(this)
        .simpleAuth()//在初始化的时候设置账号密码,重连才能成功
        .username(USERNAME)
        .password(PASSWORD.toByteArray())
        .applySimpleAuth()
        .automaticReconnectWithDefaultConfig()//自动重连
        .buildAsync()

    private val listeners = mutableListOf<MqttListener>()

    private val subTopics
        get() = arrayOf("top1", "top2", "top3")

    fun addMqttListener(listener: MqttListener) {
        if (!listeners.contains(listener)) {
            listeners.add(listener)
        }
    }

    fun removeMqttListener(listener: MqttListener) {
        listeners.remove(listener)
    }

    override fun onConnected(context: MqttClientConnectedContext) {
        Log.i(TAG, "onConnected()")
        for (l in listeners) {
            l.onConnected()
        }
        subscribeAll()
    }

    private fun subscribeAll() {
        CompletableFuture.supplyAsync({
            val futures = subTopics.map(::subscribe)
                .map {
                    it.thenCompose {
                        CompletableFuture.supplyAsync({
                            val success = !it.reasonString.isPresent
                            if (success) {
                                Log.i(TAG, "subscribe success")
                            } else {
                                Log.e(
                                    TAG, "subscribe() - reasonCodes=[${it.reasonCodes.joinToString(", ")}]" +
                                            ", reasonString=${it.reasonString}"
                                )
                            }
                            success
                        }, executor)
                    }
                }
                .toTypedArray()
            CompletableFuture.allOf(*futures).join()//等待所有订阅结果
            if(futures.all { it.get() }) {
                Log.i(TAG, "subscribeAll() - 全部订阅成功")
            }
            for (l in listeners) {
                l.onSubscribed(*subTopics)
            }
        }, executor)
    }

    override fun onDisconnected(context: MqttClientDisconnectedContext) {
        Log.e(
            TAG, "onDisconnected() - isConnected=${mqttAsynClient.state.isConnected}" +
                    ", isConnectedOrReconnect=${mqttAsynClient.state.isConnectedOrReconnect}"
        )
        for (l in listeners) {
            l.onDisconnected()
        }
    }

    fun connect() {
        mqttAsynClient
            .connectWith()
            .cleanStart(true)
            .keepAlive(30)
            .send()
            .thenAccept {
                if (it.reasonCode == Mqtt5ConnAckReasonCode.SUCCESS) {
                    Log.i(TAG, "connect() - SUCCESS")
                } else {
                    Log.e(TAG, "connect() - ${it.reasonCode}")
                }
            }
    }

    fun disconnect() {
        mqttAsynClient.disconnect().thenAccept {
            Log.i(TAG, "disconnect()")
        }
    }


    private val callback = Consumer<Mqtt5Publish> {
        val topic = it.topic.toString()
        val data = it.payloadAsBytes
        processReceivedMessage(topic, data)
    }

    private fun processReceivedMessage(topic: String, data: ByteArray) {
        //处理接收的数据
        for (l in listeners) {
            l.onReceiveMessage(topic, data)
        }
    }

    fun subscribe(topic: String): CompletableFuture<Mqtt5SubAck> {
        return mqttAsynClient.subscribeWith()
            .topicFilter(topic)
            .noLocal(true)// we do not want to receive our own message
            .qos(MqttQos.AT_MOST_ONCE)
            .callback(callback)
            .executor(executor)
            .send()
    }

    fun unsubscribe(topic: String) {
        mqttAsynClient.unsubscribeWith()
            .topicFilter(topic)
            .send().thenAccept {
                Log.i(TAG, "unsubscribe() - $it")
            }
    }

    /**
     * 发送数据
     */
    fun publish(topic: String, payload: ByteArray) {
        mqttAsynClient.publishWith()
            .topic(topic)
            .qos(MqttQos.AT_MOST_ONCE)
            .payload(payload)
            .send()
            .thenAccept { mqtt5PublishResult ->
                mqtt5PublishResult.publish.let { mqtt5Publish ->
//                    val topic = mqtt5Publish.topic.toString()
                    val data = mqtt5Publish.payloadAsBytes
                    for (l in listeners) {
                        l.onSendMessage(topic, data)
                    }
                }
            }
    }

    companion object {
        private const val TAG = "MqttManager"

        private const val SERVER_HOST = "example.com"
        private const val SERVER_PORT = 1883 // 1883即TCP协议,host不要再加上"tcp://",否则连不成功
        private const val USERNAME = "admin"
        private const val PASSWORD = "123456"

        val instance = MqttManager()
    }
}


http://www.kler.cn/a/135704.html

相关文章:

  • 通用项目工程的过程视图概览
  • Oracle OCP认证考试考点详解082系列16
  • go T 泛型
  • 壹连科技IPO闯关成功!连接器行业上市企业+1
  • Redis安装(Windows环境)
  • SpringBoot(十三)SpringBoot配置webSocket
  • 全志R128内存泄漏调试案例
  • 鸿蒙4.0开发笔记之DevEco Studio之配置代码片段快速生成(三)
  • 【Python 千题 —— 基础篇】输出可以被5整除的数
  • 嵌入式QTGit面试题
  • 计算机毕业设计选题推荐-高校后勤报修微信小程序/安卓APP-项目实战
  • 可逆矩阵的性质
  • 获取阿里云Docker镜像加速器
  • Arduino驱动DS18B20数字温度传感器(温湿度传感器)
  • OpenCV快速入门:直方图、掩膜、模板匹配和霍夫检测
  • 第四篇 《随机点名答题系统》——基础设置详解(类抽奖系统、在线答题系统、线上答题系统、在线点名系统、线上点名系统、在线考试系统、线上考试系统)
  • AtCoder Beginner Contest 329 题解A~F
  • 【数据机构】最小生成树(prim算法)
  • Harmony Ble 蓝牙App (一)扫描
  • .babyk勒索病毒解析:恶意更新如何威胁您的数据安全
  • SpringCloud相关
  • Mac安装win程序另一个方案
  • TCP传输的三次握手、四次挥手策略是什么
  • 【苏州元德维康生物医药-注册】
  • 2.3IP详解及配置
  • 给大伙讲个笑话:阿里云服务器开了安全组防火墙还是无法访问到服务