返回工作笔记
NOTE ARCHIVE现场记录无人机2026-07-08
GIS

OSD接入开发详解

OSD接入开发详解

Clark 更新于 2026-08-31 阅读约 6 分钟 阅读 0
阅读导航

本文档只说明如何从 0 搭建一个工程,通过 MQTT 订阅无人机 OSD 数据。本文不包含 OSD 模型转换、InfluxDB 入库、WebSocket 推送、历史查询等内容。

参考模块:

ruoyi-modules/other-drone-manage

核心目标:

无人机 / 第三方平台
  -> MQTT Broker
  -> Java 服务订阅 thing/product/+/osd
  -> 代码拿到 topic、设备 SN、原始 OSD JSON

1. 创建模块

在 RuoYi Cloud 项目下新建模块,例如:

ruoyi-modules/other-drone-manage

基础目录:

ruoyi-modules/other-drone-manage
├── pom.xml
└── src/main
    ├── java/com/ruoyi/otherdrone
    │   ├── OtherDroneManageApplication.java
    │   ├── config
    │   │   └── MqttPropertyConfiguration.java
    │   ├── mqtt
    │   │   ├── OsdMqttSubscriber.java
    │   │   └── model
    │   │       ├── MqttClientOptions.java
    │   │       ├── MqttProtocolEnum.java
    │   │       └── MqttUseEnum.java
    │   └── service
    │       ├── DroneOsdMessageHandler.java
    │       └── impl/DroneOsdMessageHandlerImpl.java
    └── resources
        ├── bootstrap.yml
        └── bootstrap-dev.yml

2. pom.xml 依赖

如果是在现有 RuoYi Cloud 工程中建模块,pom.xml 里至少需要 Spring Boot、RuoYi 公共依赖和 Paho MQTT 客户端。

关键依赖是:

<dependency>
    <groupId>org.eclipse.paho</groupId>
    <artifactId>org.eclipse.paho.client.mqttv3</artifactId>
    <version>1.2.5</version>
</dependency>

一个最小示例:

<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <parent>
        <groupId>com.ruoyi</groupId>
        <artifactId>ruoyi-modules</artifactId>
        <version>3.6.6</version>
    </parent>

    <artifactId>other-drone-manage</artifactId>
    <packaging>jar</packaging>

    <dependencies>
        <dependency>
            <groupId>com.ruoyi</groupId>
            <artifactId>ruoyi-common-core</artifactId>
        </dependency>

        <dependency>
            <groupId>com.ruoyi</groupId>
            <artifactId>ruoyi-common-security</artifactId>
        </dependency>

        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-web</artifactId>
        </dependency>

        <dependency>
            <groupId>org.eclipse.paho</groupId>
            <artifactId>org.eclipse.paho.client.mqttv3</artifactId>
            <version>1.2.5</version>
        </dependency>

        <dependency>
            <groupId>org.projectlombok</groupId>
            <artifactId>lombok</artifactId>
        </dependency>
    </dependencies>
</project>

如果只做独立 Spring Boot 工程,不依赖 RuoYi,也可以只保留 spring-boot-starter-webpaho mqttlombok

3. 启动类

package com.ruoyi.otherdrone;

import com.ruoyi.common.security.annotation.EnableCustomConfig;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;

@EnableCustomConfig
@SpringBootApplication
public class OtherDroneManageApplication {

    public static void main(String[] args) {
        SpringApplication.run(OtherDroneManageApplication.class, args);
        System.out.println("other-drone-manage started");
    }
}

如果这个模块需要注册到 Nacos,再按 RuoYi 现有模块加上服务注册相关依赖和配置即可;仅测试 MQTT 接收时不是必须。

4. MQTT 配置

bootstrap.yml

server:
  port: 9300

spring:
  application:
    name: other-drone-manage
  profiles:
    active: dev

bootstrap-dev.yml

mqtt:
  config:
    BASIC:
      protocol: ${MQTT_PROTOCOL:MQTT}
      host: ${MQTT_HOST:[IP已隐藏]}
      port: ${MQTT_PORT:1883}
      username: ${MQTT_USERNAME:}
      password: ${MQTT_PASSWORD:}
      client-id: ${MQTT_CLIENT_ID:other-drone-manage}
      connection-timeout: ${MQTT_CONNECTION_TIMEOUT:30}
      path: ${MQTT_PATH:}

说明:

配置说明
protocolMQTT 协议,可用 MQTTMQTTSWSWSS
hostMQTT Broker 地址
portMQTT Broker 端口
username/passwordMQTT 用户名和密码
client-id客户端 ID,同一个 Broker 下不要重复
pathWebSocket MQTT 才需要,例如 /mqtt

5. MQTT 配置模型

5.1 协议枚举

package com.ruoyi.otherdrone.mqtt.model;

public enum MqttProtocolEnum {
    MQTT("tcp://"),
    MQTTS("ssl://"),
    WS("ws://"),
    WSS("wss://");

    private final String protocolAddr;

    MqttProtocolEnum(String protocolAddr) {
        this.protocolAddr = protocolAddr;
    }

    public String getProtocolAddr() {
        return protocolAddr;
    }
}

5.2 使用类型枚举

package com.ruoyi.otherdrone.mqtt.model;

public enum MqttUseEnum {
    BASIC
}

5.3 客户端配置对象

package com.ruoyi.otherdrone.mqtt.model;

import lombok.Data;

@Data
public class MqttClientOptions {

    private MqttProtocolEnum protocol = MqttProtocolEnum.MQTT;

    private String host;

    private Integer port;

    private String username;

    private String password;

    private String clientId;

    private Integer connectionTimeout = 30;

    private String path;
}

5.4 配置读取类

package com.ruoyi.otherdrone.config;

import com.ruoyi.otherdrone.mqtt.model.MqttClientOptions;
import com.ruoyi.otherdrone.mqtt.model.MqttProtocolEnum;
import com.ruoyi.otherdrone.mqtt.model.MqttUseEnum;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.Configuration;

import jakarta.annotation.PostConstruct;
import java.util.Map;

@Slf4j
@Data
@Configuration
@ConfigurationProperties(prefix = "mqtt")
public class MqttPropertyConfiguration {

    private Map<MqttUseEnum, MqttClientOptions> config;

    @PostConstruct
    public void logConfig() {
        if (config == null || config.isEmpty()) {
            log.warn("MQTT 配置为空,请检查 mqtt.config.BASIC");
            return;
        }
        config.forEach((key, opts) ->
                log.info("MQTT 配置加载: type={}, protocol={}, host={}:{}, clientId={}",
                        key, opts.getProtocol(), opts.getHost(), opts.getPort(), opts.getClientId()));
    }

    public MqttClientOptions getBasicClientOptions() {
        if (config == null || !config.containsKey(MqttUseEnum.BASIC)) {
            throw new IllegalStateException("请配置 mqtt.config.BASIC");
        }
        return config.get(MqttUseEnum.BASIC);
    }

    public String getBasicMqttAddress() {
        return getMqttAddress(getBasicClientOptions());
    }

    public static String getMqttAddress(MqttClientOptions options) {
        StringBuilder address = new StringBuilder()
                .append(options.getProtocol().getProtocolAddr())
                .append(options.getHost().trim())
                .append(":")
                .append(options.getPort());

        if ((options.getProtocol() == MqttProtocolEnum.WS || options.getProtocol() == MqttProtocolEnum.WSS)
                && options.getPath() != null && !options.getPath().isBlank()) {
            address.append(options.getPath());
        }
        return address.toString();
    }
}

6. OSD 消息处理接口

这里不做模型转换,只把 MQTT 收到的原始 OSD JSON 交给业务层。

package com.ruoyi.otherdrone.service;

public interface DroneOsdMessageHandler {

    /**
     * 处理无人机 OSD 原始消息。
     *
     * @param deviceSn MQTT topic 中的设备 SN
     * @param topic MQTT topic
     * @param payload 原始 JSON 字符串
     */
    void handleOsd(String deviceSn, String topic, String payload);
}

示例实现:

package com.ruoyi.otherdrone.service.impl;

import com.ruoyi.otherdrone.service.DroneOsdMessageHandler;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;

@Slf4j
@Service
public class DroneOsdMessageHandlerImpl implements DroneOsdMessageHandler {

    @Override
    public void handleOsd(String deviceSn, String topic, String payload) {
        log.info("收到无人机 OSD: deviceSn={}, topic={}, payload={}", deviceSn, topic, payload);

        // 这里按业务需要处理原始 OSD:
        // 1. 直接转发给其他服务
        // 2. 放入消息队列
        // 3. 缓存在 Redis
        // 4. 后续再接入存储或 WebSocket
    }
}

7. MQTT 订阅器

核心类:启动后连接 Broker,订阅无人机 OSD 主题,收到消息后提取设备 SN 和原始 payload。

package com.ruoyi.otherdrone.mqtt;

import com.ruoyi.otherdrone.config.MqttPropertyConfiguration;
import com.ruoyi.otherdrone.service.DroneOsdMessageHandler;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.IMqttDeliveryToken;
import org.eclipse.paho.client.mqttv3.MqttCallbackExtended;
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.MqttMessage;
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
import org.springframework.boot.context.event.ApplicationReadyEvent;
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Component;

import java.nio.charset.StandardCharsets;

@Slf4j
@Component
public class OsdMqttSubscriber implements MqttCallbackExtended {

    private static final String OSD_TOPIC = "thing/product/+/osd";

    private final MqttPropertyConfiguration mqttConfig;
    private final DroneOsdMessageHandler osdMessageHandler;

    private MqttClient mqttClient;

    public OsdMqttSubscriber(MqttPropertyConfiguration mqttConfig,
                             DroneOsdMessageHandler osdMessageHandler) {
        this.mqttConfig = mqttConfig;
        this.osdMessageHandler = osdMessageHandler;
    }

    @EventListener(ApplicationReadyEvent.class)
    public void connectAndSubscribe() {
        try {
            var basicOptions = mqttConfig.getBasicClientOptions();
            String brokerUrl = mqttConfig.getBasicMqttAddress();
            String clientId = basicOptions.getClientId();

            if (clientId == null || clientId.isBlank()) {
                clientId = "other-drone-manage-" + System.currentTimeMillis();
            }

            MqttConnectOptions options = new MqttConnectOptions();
            options.setServerURIs(new String[]{brokerUrl});
            options.setUserName(basicOptions.getUsername());
            options.setPassword(toPassword(basicOptions.getPassword()));
            options.setAutomaticReconnect(true);
            options.setKeepAliveInterval(10);
            options.setConnectionTimeout(
                    basicOptions.getConnectionTimeout() == null ? 30 : basicOptions.getConnectionTimeout());
            options.setCleanSession(true);

            mqttClient = new MqttClient(brokerUrl, clientId, new MemoryPersistence());
            mqttClient.setCallback(this);
            mqttClient.connect(options);

            mqttClient.subscribe(OSD_TOPIC, 1);
            log.info("MQTT 订阅成功: broker={}, clientId={}, topic={}", brokerUrl, clientId, OSD_TOPIC);
        } catch (MqttException ex) {
            log.error("MQTT 连接或订阅失败: reason={}, message={}", ex.getReasonCode(), ex.getMessage(), ex);
        } catch (Exception ex) {
            log.error("MQTT 连接或订阅失败", ex);
        }
    }

    @Override
    public void connectComplete(boolean reconnect, String serverURI) {
        log.info("MQTT 连接完成: reconnect={}, server={}", reconnect, serverURI);
        if (!reconnect || mqttClient == null || !mqttClient.isConnected()) {
            return;
        }
        try {
            mqttClient.subscribe(OSD_TOPIC, 1);
            log.info("MQTT 重连后重新订阅成功: topic={}", OSD_TOPIC);
        } catch (MqttException ex) {
            log.error("MQTT 重连后重新订阅失败", ex);
        }
    }

    @Override
    public void connectionLost(Throwable cause) {
        log.warn("MQTT 连接断开,将自动重连: {}", cause == null ? null : cause.getMessage());
    }

    @Override
    public void messageArrived(String topic, MqttMessage message) {
        String payload = new String(message.getPayload(), StandardCharsets.UTF_8);
        String deviceSn = parseDeviceSn(topic);

        if (deviceSn == null || deviceSn.isBlank()) {
            log.warn("OSD topic 无法解析设备 SN: topic={}", topic);
            return;
        }

        log.debug("MQTT 收到 OSD: topic={}, deviceSn={}, size={}", topic, deviceSn, payload.length());
        osdMessageHandler.handleOsd(deviceSn, topic, payload);
    }

    @Override
    public void deliveryComplete(IMqttDeliveryToken token) {
    }

    private String parseDeviceSn(String topic) {
        if (topic == null) {
            return null;
        }
        String[] parts = topic.split("/");
        // thing/product/{deviceSn}/osd
        return parts.length >= 3 ? parts[2] : null;
    }

    private char[] toPassword(String password) {
        return password == null ? new char[0] : password.toCharArray();
    }
}

8. OSD Topic 约定

当前订阅主题:

thing/product/+/osd

厂商推送时 topic 应该是:

thing/product/{无人机SN}/osd

示例:

thing/product/DRONE_SN_001/osd

代码会从 topic 第三段取出设备 SN:

thing / product / DRONE_SN_001 / osd
                 ^
                 deviceSn

9. OSD Payload 约定

本文档不要求转换模型,所以 payload 只要求是合法 JSON 字符串。业务处理时可以先完整保存或转发。

推荐格式:

{
  "timestamp": 1719990000000,
  "data": {
    "[坐标已隐藏]
    "[坐标已隐藏]
    "height": 120.5,
    "mode_code": 5,
    "horizontal_speed": 8.2,
    "vertical_speed": 0.1,
    "battery": {
      "capacity_percent": 78
    }
  }
}

如果只做接收,代码不会关心 data 里面有哪些字段。

10. 从 OSD 中获取关键字段

DJI Dock 到云端的 M3D/M3TD OSD 结构中,飞机状态字段通常在 MQTT payload 的 data 节点下。也就是:

{
  "timestamp": 1719990000000,
  "data": {
    "[坐标已隐藏]
    "[坐标已隐藏]
    "height": 120.5,
    "elevation": 88.3,
    "attitude_head": 12,
    "attitude_pitch": -2.1,
    "attitude_roll": 0.5,
    "position_state": {
      "gps_number": 24,
      "quality": 10,
      "is_fixed": 2,
      "rtk_number": 18
    },
    "81-0-0": {
      "payload_index": "81-0-0",
      "gimbal_pitch": -45,
      "gimbal_yaw": 18.36,
      "gimbal_roll": 0,
      "zoom_factor": 0.56,
      "thermal_global_temperature_max": 43.6,
      "thermal_global_temperature_min": 26.8
    },
    "payloads": [
      {
        "payload_index": "81-0-0",
        "gimbal_pitch": -90,
        "gimbal_yaw": 12,
        "gimbal_roll": 0,
        "zoom_factor": 1
      }
    ],
    "cameras": [
      {
        "payload_index": "81-0-0",
        "zoom_factor": 1
      }
    ]
  }
}

字段提取建议:

信息类别OSD 字段 / JSON 路径用途
无人机位置data.[坐标已隐藏]dedata.heightdata.elevation判断拍摄点经[坐标已隐藏]高度;height 通常是椭球高 / 绝对高度,elevation 通常是相对起飞点高度
飞行姿态data.attitude_headdata.attitude_pitchdata.attitude_roll判断机头方向、俯仰、横滚,辅助判断拍摄姿态是否稳定
云台姿态data["81-0-0"].gimbal_pitchdata["81-0-0"].gimbal_yawdata["81-0-0"].gimbal_roll;兼容 data.payloads[].gimbal_*判断相机视线方向,常用于确认是否垂直向下拍摄
相机参数data["81-0-0"].payload_indexdata["81-0-0"].zoom_factordata.cameras[].payload_indexdata.cameras[].zoom_factor通过负载索引和变焦倍率判断视场范围;镜头焦距、传感器尺寸等固定参数建议从设备能力或配置表补充
时间信息timestamp、识别结果里的 detect_time将识别结果和飞行状态按时间对齐
定位质量data.position_state.gps_numberdata.position_state.qualitydata.position_state.is_fixeddata.position_state.rtk_number判断定位可信度;RTK 固定、卫星数充足时可信度更高

注意:

  • 真实 DJI Dock OSD 中,负载状态可能不是 payloads[] 数组,而是 data 下以负载编号作为动态 key 的对象,例如 data["81-0-0"]
  • data["81-0-0"]payloads[]cameras[] 都可能包含 payload_indexzoom_factor,实际使用时建议按 payload_index 对齐同一个负载。
  • gimbal_pitch/yaw/roll 通常跟负载挂载相关,优先从 data[payload_index] 这种动态负载节点取,取不到再回退到 payloads[]
  • detect_time 不是 DJI OSD 固有字段,一般来自 AI 识别结果;需要和 MQTT 外层 timestamp 做时间对齐。
  • DJI 文档中 position_state.quality = 10 表示 RTK fixed,可作为高可信定位条件之一。 你贴的真实报文里 quality = 5is_fixed = 2rtk_number = 40,说明 RTK 固定状态存在,但质量枚举不是 fixed 档,业务上建议把它标记为“可用但非最高质量”,不要简单等同 quality = 10

10.1 Java 读取示例

继续沿用原始 JSON 处理,不需要 DTO 转换。可以在 DroneOsdMessageHandlerImpl.handleOsd 里使用 Jackson JsonNode 读取:

package com.ruoyi.otherdrone.service.impl;

import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.ruoyi.otherdrone.service.DroneOsdMessageHandler;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;

@Slf4j
@Service
public class DroneOsdMessageHandlerImpl implements DroneOsdMessageHandler {

    private final ObjectMapper objectMapper;

    public DroneOsdMessageHandlerImpl(ObjectMapper objectMapper) {
        this.objectMapper = objectMapper;
    }

    @Override
    public void handleOsd(String deviceSn, String topic, String payload) {
        try {
            JsonNode root = objectMapper.readTree(payload);
            JsonNode data = root.path("data");

            [坐标已隐藏]ot.path("timestamp").asLong(System.currentTimeMillis());

            Double [坐标已隐藏]ta, "[坐标已隐藏]
            Double [坐标已隐藏]ata, "[坐标已隐藏]
            Double height = getDouble(data, "height");
            Double elevation = getDouble(data, "elevation");

            Double attitudeHead = getDouble(data, "attitude_head");
            Double attitudePitch = getDouble(data, "attitude_pitch");
            Double attitudeRoll = getDouble(data, "attitude_roll");

            JsonNode positionState = data.path("position_state");
            Integer gpsNumber = getInt(positionState, "gps_number");
            Integer quality = getInt(positionState, "quality");
            Integer isFixed = getInt(positionState, "is_fixed");
            Integer rtkNumber = getInt(positionState, "rtk_number");

            JsonNode payloadNode = firstPayloadNode(data);
            String payloadIndex = getText(payloadNode, "payload_index");
            Double gimbalPitch = getDouble(payloadNode, "gimbal_pitch");
            Double gimbalYaw = getDouble(payloadNode, "gimbal_yaw");
            Double gimbalRoll = getDouble(payloadNode, "gimbal_roll");
            Double zoomFactor = getDouble(payloadNode, "zoom_factor");

            if (zoomFactor == null) {
                JsonNode cameraNode = findCameraByPayloadIndex(data.path("cameras"), payloadIndex);
                zoomFactor = getDouble(cameraNode, "zoom_factor");
            }

            boolean rtkFixed = Integer.valueOf(2).equals(isFixed);
            boolean highQuality = Integer.valueOf(10).equals(quality);
            boolean rtkUsable = rtkFixed && rtkNumber != null && rtkNumber >= 10;
            boolean rtkReliable = highQuality || rtkUsable;

            log.info("OSD关键字段: sn={}, time={}, [坐标已隐藏]{}, elevation={}, " +
                            "head={}, pitch={}, roll={}, gimbalPitch={}, gimbalYaw={}, gimbalRoll={}, " +
                            "payloadIndex={}, zoomFactor={}, gps={}, quality={}, fixed={}, rtk={}, reliable={}",
                    deviceSn, mqttTimestamp, [坐标已隐藏]ight, elevation,
                    attitudeHead, attitudePitch, attitudeRoll,
                    gimbalPitch, gimbalYaw, gimbalRoll,
                    payloadIndex, zoomFactor, gpsNumber, quality, isFixed, rtkNumber, rtkReliable);
        } catch (Exception e) {
            log.warn("解析 OSD 关键字段失败: deviceSn={}, topic={}, payload={}", deviceSn, topic, payload, e);
        }
    }

    private Double getDouble(JsonNode node, String field) {
        JsonNode value = node == null ? null : node.get(field);
        return value == null || value.isNull() || !value.isNumber() ? null : value.asDouble();
    }

    private Integer getInt(JsonNode node, String field) {
        JsonNode value = node == null ? null : node.get(field);
        return value == null || value.isNull() || !value.isNumber() ? null : value.asInt();
    }

    private String getText(JsonNode node, String field) {
        JsonNode value = node == null ? null : node.get(field);
        return value == null || value.isNull() ? null : value.asText();
    }

    private JsonNode firstArrayItem(JsonNode arrayNode) {
        return arrayNode != null && arrayNode.isArray() && !arrayNode.isEmpty() ? arrayNode.get(0) : null;
    }

    private JsonNode firstPayloadNode(JsonNode data) {
        JsonNode dynamicPayload = firstDynamicPayloadNode(data);
        if (dynamicPayload != null) {
            return dynamicPayload;
        }
        return firstArrayItem(data.path("payloads"));
    }

    private JsonNode firstDynamicPayloadNode(JsonNode data) {
        if (data == null || !data.isObject()) {
            return null;
        }
        var fields = data.fields();
        while (fields.hasNext()) {
            var entry = fields.next();
            JsonNode value = entry.getValue();
            if (value != null && value.isObject() && value.has("payload_index")) {
                return value;
            }
        }
        return null;
    }

    private JsonNode findCameraByPayloadIndex(JsonNode cameras, String payloadIndex) {
        if (cameras == null || !cameras.isArray()) {
            return null;
        }
        for (JsonNode camera : cameras) {
            if (payloadIndex != null && payloadIndex.equals(getText(camera, "payload_index"))) {
                return camera;
            }
        }
        return firstArrayItem(cameras);
    }
}

10.2 按 payload_index 获取指定负载

如果无人机挂载多个负载,不建议直接取第一个负载。应该按识别结果或业务配置里的 payload_index 找到对应负载。

真实 OSD 可能使用动态 key:

{
  "data": {
    "81-0-0": {
      "payload_index": "81-0-0",
      "gimbal_pitch": -45,
      "zoom_factor": 0.5678
    }
  }
}

读取时先尝试 data.get(payloadIndex),再回退扫描数组:

private JsonNode findPayloadByIndex(JsonNode data, String payloadIndex) {
    if (data == null || payloadIndex == null || payloadIndex.isBlank()) {
        return null;
    }

    JsonNode dynamicPayload = data.get(payloadIndex);
    if (dynamicPayload != null && dynamicPayload.isObject()) {
        return dynamicPayload;
    }

    return findPayloadByIndexFromArray(data.path("payloads"), payloadIndex);
}
private JsonNode findPayloadByIndexFromArray(JsonNode payloads, String payloadIndex) {
    if (payloads == null || !payloads.isArray()) {
        return null;
    }
    for (JsonNode payload : payloads) {
        JsonNode value = payload.get("payload_index");
        if (value != null && payloadIndex.equals(value.asText())) {
            return payload;
        }
    }
    return null;
}

10.3 识别结果与 OSD 时间对齐

AI 识别结果一般会有自己的识别时间,例如:

{
  "detect_time": 1719990000123,
  "payload_index": "81-0-0",
  "target": "person"
}

对齐时建议:

  1. 读取 OSD 外层 timestamp
  2. 读取识别结果 detect_time
  3. 在 OSD 缓存中查找 abs(osd.timestamp - detect_time) 最小的一帧。
  4. 设置最大允许误差,例如 1000ms 或 2000ms;超过阈值则认为无法可靠对齐。

示例:

[坐标已隐藏]dTimestamp - detectTime);
if (diff > 2000L) {
    log.warn("识别结果与 OSD 时间差过大: diff={}ms", diff);
}

11. 本地调试

11.1 启动 MQTT Broker

可以使用本地 Mosquitto:

docker run --rm -it -p 1883:1883 eclipse-mosquitto:2

如果启用了账号密码,需要按实际 Broker 配置。

11.2 启动服务

mvn compile -pl ruoyi-modules/other-drone-manage -am -DskipTests

然后启动:

java -jar ruoyi-modules/other-drone-manage/target/other-drone-manage.jar

或直接在 IDE 里启动:

com.ruoyi.otherdrone.OtherDroneManageApplication

正常日志:

MQTT 配置加载
MQTT 订阅成功

11.3 模拟推送 OSD

mosquitto_pub \
  -h [IP已隐藏] \
  -p 1883 \
  -t 'thing/product/DRONE_SN_001/osd' \
  -m '{
    "timestamp": 1719990000000,
    "data": {
      "[坐标已隐藏]
      "[坐标已隐藏]
      "height": 120.5,
      "elevation": 88.3,
      "attitude_head": 12,
      "attitude_pitch": -2.1,
      "attitude_roll": 0.5,
      "mode_code": 5,
      "position_state": {
        "gps_number": 24,
        "quality": 10,
        "is_fixed": 2,
        "rtk_number": 18
      },
      "payloads": [
        {
          "payload_index": "81-0-0",
          "gimbal_pitch": -90,
          "gimbal_yaw": 12,
          "gimbal_roll": 0,
          "zoom_factor": 1
        }
      ]
    }
  }'

服务日志应能看到:

收到无人机 OSD: deviceSn=DRONE_SN_001, topic=thing/product/DRONE_SN_001/osd, payload=...

12. 常见问题

12.1 MQTT 连接失败

检查:

  • MQTT_HOSTMQTT_PORT 是否正确。
  • Broker 是否允许当前网络访问。
  • 用户名密码是否正确。
  • 协议是否匹配:普通 MQTT 用 MQTT,TLS 用 MQTTS,WebSocket 用 WS/WSS

12.2 订阅成功但收不到消息

检查:

  • 设备推送 topic 是否是 thing/product/{sn}/osd
  • Broker ACL 是否允许订阅 thing/product/+/osd
  • 推送端和订阅端是否连接到了同一个 Broker。
  • MQTT QoS 和 Clean Session 是否符合 Broker 策略。

12.3 能收到消息但解析不到 SN

检查 topic 层级是否正确:

正确:thing/product/DRONE_SN_001/osd
错误:thing/DRONE_SN_001/osd
错误:product/DRONE_SN_001/osd

12.4 重启后 clientId 冲突

同一个 Broker 下 MQTT clientId 必须唯一。如果多个实例同时启动,建议配置:

mqtt:
  config:
    BASIC:
      client-id: other-drone-manage-${server.port}

或在代码里拼接机器 IP、端口、随机后缀。

13. 后续扩展点

当前文档只完成“拿到 OSD 原始消息”。后续如果需要继续扩展,可以在 DroneOsdMessageHandlerImpl.handleOsd 中加入:

  • JSON 字段校验
  • OSD 数据存储
  • 推送到 WebSocket
  • 转发到 Kafka / Redis Stream
  • 按 SN 更新设备在线状态
  • 按厂商协议做字段转换