1. 为什么选择MQTTs?从HTTP到MQTT的升级之路

几年前我做物联网项目,第一反应也是用HTTP API去连接云平台,就像原始文章里那样,用OkHttp发个请求,等个响应,看起来简单直接。但实际用起来,尤其是设备数量一多、数据需要频繁上报或者实时控制的时候,HTTP那套“一问一答”的模式就开始捉襟见肘了。想象一下,你的手机App想随时知道家里温湿度传感器的读数,如果用HTTP,你就得隔几秒让App去“问”一次云平台:“数据有更新吗?” 大部分时候问到的都是旧数据,白白浪费流量和电量,体验还很差。

这时候,MQTT协议的优势就体现出来了。它就像一个高效的“订阅-发布”消息系统。你的Android应用(作为客户端)只需要一次连接到OneNET平台(作为MQTT Broker),然后“订阅”你关心的主题(Topic),比如 /your_product_id/your_device_id/sensor/temperature。之后,只要这个温度传感器设备往这个主题发布了新数据,平台会立刻把这条消息“推送”给你的App。整个过程是实时的、双向的,App不用反复去“问”,数据来了自动收。这特别适合物联网场景里设备状态监控、远程指令下发这些需要即时性的功能。

那MQTTs后面的这个“s”又是什么意思呢?它就是SSL/TLS,给MQTT这条通信通道加了一把牢固的锁。所有数据在传输前都会加密,防止在网络上被窃听或篡改。OneNET平台强制使用MQTTs进行连接,这其实是件好事,意味着从你的App到云端,整个链路都是安全的。相比之下,原始的HTTP连接(非HTTPS)在物联网这种涉及实际控制和敏感数据的场景下,风险就太高了。所以,这次我们彻底告别旧式的请求-响应模式,拥抱更现代、更高效、也更安全的MQTTs长连接方案。

我理解很多Android开发者可能对HTTP更熟悉,觉得MQTT又要搞连接、又要管订阅,会不会很复杂?其实不然,用对了库,整套逻辑比管理一堆异步HTTP回调还要清晰。接下来,我就带你一步步,用一个最主流的库,把Android App稳稳当当地接到OneNET上。

2. 前期准备:兵马未动,粮草先行

在开始写代码之前,我们需要把“战场”布置好。这包括在OneNET云平台上创建设备,以及在Android项目中引入强大的“武器库”。

2.1 在OneNET平台上的配置

首先,你得有个OneNET的账号。登录后,进入控制台,找到“多协议接入”部分。这里我们选择MQTT(旧版) 或者 MQTT物联网套件,根据你的产品情况来定。我以“MQTT(旧版)”为例,因为它比较通用。

  1. 创建产品:点击创建产品,产品名称随便取,比如“智能家居测试”。在联网方式和协议这里,一定要选择“设备接入协议”为 MQTT。其他选项像数据格式、认证方式保持默认即可。创建成功后,你会得到一个重要的 Product ID(产品ID),记下来。
  2. 创建设备:进入刚创建的产品,添加一个设备。设备名称也可以自定义,比如“客厅温湿度计”。创建成功后,平台会生成三个关键信息:
    • Device ID(设备ID):设备的唯一标识。
    • Auth Info(鉴权信息):在旧版MQTT中,这通常是一个自定义的字符串,你可以把它理解为设备的密码。我们连接时会用到它。
    • API Key:注意,这个Key是用于通过HTTP API管理设备的(比如原始文章里查询数据),它不能用于MQTT连接认证。很多新手会在这里搞混,切记。

把Product ID、Device ID和Auth Info这三样妥善保存,它们就是我们App连接平台的“钥匙”。

2.2 Android项目依赖与网络配置

打开你的Android Studio项目,我们开始集成MQTT客户端库。业界最成熟、最通用的选择是 Eclipse Paho 的Android版本。它非常稳定,社区支持也好。

在你的 app 模块下的 build.gradle 文件里,添加以下依赖:

dependencies {
    implementation 'org.eclipse.paho:org.eclipse.paho.client.mqttv3:1.2.5'
    implementation 'org.eclipse.paho:org.eclipse.paho.android.service:1.1.1'
}

第一个是核心的MQTT库,第二个是Android Service封装,它帮我们在后台管理MQTT连接的生命周期,即使App退到后台,连接也能保持(取决于你的需求和服务策略)。

接下来,由于我们使用的是MQTTs(加密连接),Android系统从某个版本开始对非加密流量(Cleartext Traffic)限制越来越严。虽然我们用的是SSL,但为了兼容性和避免一些未知问题,最好配置一下网络安全。在 app/src/main/res/xml 目录下(如果没有xml目录就新建一个),创建一个文件 network_security_config.xml:

<?xml version="1.0" encoding="utf-8"?>
<network-security-config>
    <base-config cleartextTrafficPermitted="false">
        <trust-anchors>
            <certificates src="system" />
            <!-- 如果你有自定义CA证书,可以在这里添加 -->
            <!-- <certificates src="@raw/my_custom_ca"/> -->
        </trust-anchors>
    </base-config>
    <!-- 如果需要针对特定域名放开明文(不推荐),可以在这里配置 -->
    <!-- <domain-config cleartextTrafficPermitted="true">
        <domain includeSubdomains="true">insecure.example.com</domain>
    </domain-config> -->
</network-security-config>

这个配置告诉Android系统:我们默认不允许明文传输,信任系统内置的证书颁发机构(CA)。OneNET的SSL证书是由正规CA签发的,所以系统会自动信任。

然后,在 AndroidManifest.xml 中应用这个配置,并声明必要的权限和服务:

<manifest ...>
    <!-- 网络权限必不可少 -->
    <uses-permission android:name="android.permission.INTERNET" />
    <!-- 如果需要后台运行,可能需要这个 -->
    <uses-permission android:name="android.permission.WAKE_LOCK" />

    <application
        ...
        android:networkSecurityConfig="@xml/network_security_config"
        ...>
        <!-- 注册Paho的MqttService -->
        <service android:name="org.eclipse.paho.android.service.MqttService" />
        ...
    </application>
</manifest>

做完这些,我们的开发环境就准备好了。比起原始文章里只加一个网络权限,我们这一步做得更细致,为安全的MQTTs连接打好了基础。

3. 连接建立:与OneNET的第一次安全握手

万事俱备,只欠连接。建立MQTTs连接是整个流程的核心,这里面的参数配置是关键,配错一个就连不上。我会把每个参数都掰开讲清楚。

首先,我们需要构造连接用的服务器地址(URI)。OneNET的MQTTs接入点通常是这样的:ssl://mqtts.heclouds.com:1883。注意,协议头是 ssl://,不是 tcp://,也不是 http://。端口 1883 是MQTT的标准端口(SSL)。

然后,我们需要准备连接选项 MqttConnectOptions。这是重头戏:

import org.eclipse.paho.client.mqttv3.MqttClient;
import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
import org.eclipse.paho.client.mqttv3.MqttException;
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;

public class MQTTManager {
    private MqttClient mqttClient;
    private static final String SERVER_URI = "ssl://mqtts.heclouds.com:1883";
    private static final String CLIENT_ID = "Your_Android_Client_ID"; // 自定义,需唯一
    private static final String PRODUCT_ID = "你的产品ID";
    private static final String DEVICE_ID = "你的设备ID";
    private static final String AUTH_INFO = "你的设备鉴权信息";

    public void connect() {
        try {
            // 持久化方式,这里用内存,简单场景够用
            MemoryPersistence persistence = new MemoryPersistence();
            // 创建客户端实例。CLIENT_ID需要保证唯一,避免多个设备冲突。
            mqttClient = new MqttClient(SERVER_URI, CLIENT_ID, persistence);

            MqttConnectOptions options = new MqttConnectOptions();
            // 设置不清除会话,这样重连后还能收到离线期间的消息(如果QoS支持)
            options.setCleanSession(false);
            // 设置连接超时时间
            options.setConnectionTimeout(10);
            // 设置心跳间隔,单位秒。OneNET可能有要求,比如60-120秒。
            options.setKeepAliveInterval(60);
            // !!!最重要的:OneNET MQTT旧版鉴权方式 !!!
            // 用户名格式为:产品ID
            options.setUserName(PRODUCT_ID);
            // 密码格式为:鉴权信息
            options.setPassword(AUTH_INFO.toCharArray());

            // 设置SSL/TLS相关(Paho库通常会根据ssl://前缀自动处理,但显式设置更稳妥)
            // 可以配置SSL Socket Factory,这里使用系统默认的,一般无需额外配置
            // options.setSocketFactory(SSLSocketFactory.getDefault());

            // 设置遗嘱消息(Last Will)。这是个好习惯,设备异常断开时,平台会发布此消息到指定主题。
            // 例如,设备离线时,通知其他订阅者。
            String willTopic = PRODUCT_ID + "/" + DEVICE_ID + "/status";
            String willPayload = "offline";
            options.setWill(willTopic, willPayload.getBytes(), 2, true); // QoS=2, retained=true

            // 设置回调,用于接收连接状态、消息到达等事件
            mqttClient.setCallback(new MyMqttCallback());

            // 开始连接
            mqttClient.connect(options);
            Log.d("MQTT", "连接命令已发送");

        } catch (MqttException e) {
            e.printStackTrace();
            Log.e("MQTT", "连接失败: " + e.getMessage());
            // 这里可以根据不同的错误码进行重试或提示用户
        }
    }
}

我来解释几个容易踩坑的点:

  1. Client ID:这个需要你自己生成一个唯一的字符串。可以用“产品ID+设备ID+时间戳”的组合,或者UUID。确保每次连接(尤其是重连)时,如果CleanSession为false,Client ID要保持一致,否则服务器会认为是新会话。
  2. 鉴权:这是和HTTP API完全不同的地方。MQTT连接不使用API Key,而是用UserName和Password字段。对于OneNET旧版MQTT,UserName填产品ID,Password填设备鉴权信息。这个信息在设备详情页能找到,千万别填错了。
  3. 遗嘱消息:这个功能非常实用。你设定了遗嘱主题和消息后,一旦你的App(设备)网络异常断开,且没有发送正常的断开包,OneNET平台就会自动帮你把这条遗嘱消息发布出去。其他订阅了这个主题的应用就能立刻知道这个设备离线了,实现了设备在线的状态同步。
  4. 回调设置:我们在连接前就设置了回调MyMqttCallback,这个类需要实现MqttCallback接口。这样,连接成功、丢失、收到消息等事件,我们都能第一时间知道。

连接操作本身是异步的,结果会通过回调通知。所以我们需要耐心等待,并在UI上给用户适当的连接状态提示。

4. 消息收发实战:订阅、发布与回调处理

连接成功只是第一步,我们的目标是交换数据。在MQTT的世界里,这围绕着主题(Topic) 进行。OneNET对主题有固定的格式要求,我们必须遵守。

4.1 主题格式与订阅

OneNET的主题格式通常为:$sys/{pid}/{device-name}/thing/property/post 或更通用的 {product_id}/{device_id}/{自定义路径}。具体需要查阅你所用套件的文档。我们假设一个简单的自定义格式:{产品ID}/{设备ID}/data 用于上报数据,{产品ID}/{设备ID}/cmd 用于接收命令。

连接成功后,我们就可以订阅我们关心的主题了。比如,我们的App想接收平台下发的控制指令:

public void subscribeToCommandTopic() {
    if (mqttClient != null && mqttClient.isConnected()) {
        String commandTopic = PRODUCT_ID + "/" + DEVICE_ID + "/cmd";
        try {
            // 订阅主题,QoS级别设为1(至少送达一次)
            mqttClient.subscribe(commandTopic, 1);
            Log.d("MQTT", "已订阅主题: " + commandTopic);
        } catch (MqttException e) {
            e.printStackTrace();
            Log.e("MQTT", "订阅失败: " + e.getMessage());
        }
    }
}

QoS(服务质量) 是个重要概念。它有三个级别:

  • 0(最多一次):发出去就不管了,可能丢失。
  • 1(至少一次):确保对方至少收到一次,但可能重复。
  • 2(恰好一次):确保对方恰好收到一次,最可靠但开销大。

对于控制指令,我们通常选择QoS 1,平衡了可靠性和效率。订阅之后,所有发往这个主题的消息,都会推送到我们的回调类中。

4.2 实现回调处理消息

现在来实现前面提到的 MyMqttCallback 类:

import org.eclipse.paho.client.mqttv3.IMqttDeliveryToken;
import org.eclipse.paho.client.mqttv3.MqttCallback;
import org.eclipse.paho.client.mqttv3.MqttMessage;

public class MyMqttCallback implements MqttCallback {

    @Override
    public void connectionLost(Throwable cause) {
        // 连接丢失,通常在这里触发重连逻辑
        Log.w("MQTT", "连接丢失", cause);
        // 可以启动一个定时任务,尝试重新连接
    }

    @Override
    public void messageArrived(String topic, MqttMessage message) throws Exception {
        // 收到消息!在这里处理
        String payload = new String(message.getPayload());
        Log.d("MQTT", "收到消息. 主题: " + topic + ", 内容: " + payload);
        int qos = message.getQos();

        // 根据不同的主题,解析并处理payload(通常是JSON格式)
        if (topic.endsWith("/cmd")) {
            // 处理控制命令
            handleCommand(payload);
        } else if (topic.endsWith("/status")) {
            // 处理其他设备的状态遗嘱消息
            updateDeviceStatus(payload);
        }

        // 注意:此方法运行在非UI线程,如果需要更新UI,请使用Handler或runOnUiThread
        runOnUiThread(() -> {
            // 更新UI,例如将收到的命令显示在TextView上
            textViewReceivedMsg.setText("收到命令: " + payload);
        });
    }

    @Override
    public void deliveryComplete(IMqttDeliveryToken token) {
        // 消息发布成功后的回调(当QoS>0时)
        Log.d("MQTT", "消息发布完成");
        try {
            // 可以获取已发布的消息
            // MqttMessage deliveredMessage = token.getMessage();
        } catch (Exception e) {
            e.printStackTrace();
        }
    }

    private void handleCommand(String jsonCommand) {
        // 解析JSON,执行对应的操作,比如控制一个开关
        try {
            JSONObject obj = new JSONObject(jsonCommand);
            String cmd = obj.getString("command");
            boolean switchState = obj.getBoolean("value");
            // 这里可以调用硬件接口或改变App状态
            Log.d("MQTT", "执行命令: " + cmd + ", 状态: " + switchState);
        } catch (JSONException e) {
            e.printStackTrace();
        }
    }
}

messageArrived 方法是数据流的终点,也是业务逻辑的起点。所有订阅主题的消息都会涌向这里。你需要根据topic来区分消息类型,并解析payload(通常是JSON字符串)来执行具体操作。记住,这个方法不在主线程,更新UI一定要做线程切换。

4.3 发布消息上报数据

有来有往。我们的App也需要向平台上报数据,比如传感器读数。这就是发布消息。

public void publishSensorData(float temperature, float humidity) {
    if (mqttClient != null && mqttClient.isConnected()) {
        String dataTopic = PRODUCT_ID + "/" + DEVICE_ID + "/data";
        // 构造JSON格式的数据
        JSONObject dataJson = new JSONObject();
        try {
            dataJson.put("temp", temperature);
            dataJson.put("humi", humidity);
            dataJson.put("timestamp", System.currentTimeMillis());
        } catch (JSONException e) {
            e.printStackTrace();
        }

        String payload = dataJson.toString();
        MqttMessage message = new MqttMessage(payload.getBytes());
        // 设置QoS级别,数据上报用1比较合适
        message.setQos(1);
        // 是否保留消息。如果为true,Broker会保存这条消息,后续新订阅者能立刻收到。
        message.setRetained(false);

        try {
            mqttClient.publish(dataTopic, message);
            Log.d("MQTT", "数据发布成功: " + payload);
        } catch (MqttException e) {
            e.printStackTrace();
            Log.e("MQTT", "数据发布失败", e);
        }
    } else {
        Log.w("MQTT", "客户端未连接,无法发布消息");
        // 可以在这里缓存数据,等连接恢复后再发送
    }
}

发布消息相对简单。确定主题、构造消息体、设置QoS和保留标志,然后调用publish方法。这里我建议将数据格式封装成JSON,结构清晰,也方便OneNET平台的数据流解析。如果遇到网络断开,一个好的实践是先将数据缓存到本地(比如用Room数据库),等网络恢复后,再将积压的数据按顺序发布出去,确保数据不丢失。

5. 避坑指南与进阶优化

走通了基本流程,我们再来聊聊实战中会遇到的那些“坑”,以及如何让我们的连接更健壮、体验更好。

5.1 连接保活与断线重连

移动网络环境不稳定,App切后台也可能被系统限制。MQTT的KeepAliveInterval(心跳间隔)就是用来保活的。客户端会定期发送一个PING请求告诉服务器:“我还活着”。如果服务器在1.5 * KeepAliveInterval时间内没收到任何消息(包括PING),就会认为客户端死了,断开连接。

在Paho Android Service中,重连机制可以这样实现:

// 在MqttConnectOptions中设置自动重连
options.setAutomaticReconnect(true);
options.setMaxReconnectDelay(30000); // 最大重连延迟30秒

设置setAutomaticReconnect(true)后,Paho库会在连接断开后自动尝试重连,并且重连间隔会指数级增长(直到最大值),避免频繁重连浪费资源。同时,我们也要在connectionLost回调里加入自己的重连逻辑,作为双重保险。

5.2 后台服务与通知

如果你的App需要持续接收消息(比如即时聊天、实时监控),那么必须考虑后台运行。单纯在Activity里持有MQTT客户端,一旦Activity销毁,连接就断了。这时候就需要用到MqttService以及一个前台服务(Foreground Service)。

你可以创建一个继承自Service的类,在里面初始化和管理MqttAndroidClient(这是Paho Android库提供的,封装了Service绑定逻辑)。然后通过startForeground()方法启动一个前台服务,并显示一个持续的通知,告诉用户App正在后台保持连接。这样即使App退到后台,系统也不容易杀死你的服务,连接得以维持。这部分涉及Android后台机制,需要仔细处理,避免耗电过高。

5.3 消息去重与有序性

当QoS设置为1时,可能会收到重复的消息。这是因为发送方没收到确认,可能会重发。所以,在messageArrived里处理业务时,最好能做到幂等性。简单说,就是同一条命令执行多次,结果应该和执行一次一样。比如,“开灯”命令,无论收到多少次,最终灯都是开的状态,而不会因为收到两次就出问题。可以为重要的消息添加唯一ID,在本地记录已处理过的ID,来过滤重复消息。

对于消息顺序,MQTT协议本身只保证在单个主题(Topic)上,对同一个发送者,消息是按顺序送达的。如果你的业务对全局顺序有严格要求,可能需要自己在消息体里加序列号,或者在服务端做更复杂的逻辑。

5.4 资源释放与生命周期管理

这是一个很容易被忽略但至关重要的问题。一定要在适当的时候(比如Activity的onDestroy,或者Service的onDestroy)断开MQTT连接并释放资源。

@Override
protected void onDestroy() {
    super.onDestroy();
    if (mqttClient != null && mqttClient.isConnected()) {
        try {
            // 发送遗嘱消息(如果设置了),然后断开连接
            mqttClient.disconnect();
            // 关闭客户端,释放所有资源
            mqttClient.close();
        } catch (MqttException e) {
            e.printStackTrace();
        }
    }
}

不妥善关闭连接,可能会导致服务器端资源泄露,或者下次连接时出现“客户端ID已存在”的错误。特别是CleanSession设为false的情况下,服务器会为客户端保存会话状态(如未确认的QoS1/2消息),正确断开连接能让服务器清理这些状态。

最后,记得多打Log,把连接状态、收发的消息主题和内容都打印出来。OneNET控制台也提供了设备日志查询功能,两边对照着看,绝大部分问题都能定位。MQTTs连接OneNET,一旦跑通,你会发现它比HTTP轮询要优雅和高效得多,整个数据流是“活”的。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐