上一篇用 MQTTX 手动发了第一条消息,跑通了「设备→EMQX→Dashboard」这条链路。这次用 Java 代码把 MQTTX 替换掉,收到的消息直接写进 TDengine。


一、整体目标

上一篇的链路:

MQTTX 手动发消息 → EMQX → Dashboard 看流量
                    ↓
              (消息进了 Broker,但没存起来)

这次要做的:

Spring Boot 订阅 → EMQX 转发消息 → 代码解析 JSON → JDBC 写入 TDengine

跑完之后,用 MQTTX 发一条消息,刷新 TDengine 控制台就能查到数据,不再需要手动 INSERT


二、新建 Spring Boot 项目

用 IDEA 新建一个项目,依赖选这几个就够了:

  • Spring Web(后面用 REST 接口测试)
  • Lombok(省 getter/setter)
  • 后面手动加 MQTT 和 TDengine 的依赖

JDK 版本选 8 或者 17 都行,我用的 17。

然后在 pom.xml 里加三个依赖:

<!-- MQTT 客户端 -->
<dependency>
    <groupId>org.eclipse.paho</groupId>
    <artifactId>org.eclipse.paho.client.mqttv3</artifactId>
    <version>1.2.5</version>
</dependency>

<!-- TDengine JDBC 驱动 -->
<dependency>
    <groupId>com.taosdata.jdbc</groupId>
    <artifactId>taos-jdbcdriver</artifactId>
    <version>3.3.0</version>
</dependency>

<!-- fastjson,解析 MQTT 消息 -->
<dependency>
    <groupId>com.alibaba</groupId>
    <artifactId>fastjson</artifactId>
    <version>1.2.83</version>
</dependency>

三个依赖:Paho 做 MQTT 客户端,taos-jdbcdriver 连 TDengine,fastjson 解析消息体。


三、配 MQTT 连接

application.yml

mqtt:
  broker-url: tcp://localhost:1883
  client-id: vehicle-server-001
  topic: vehicle/+/data
  qos: 1

tdengine:
  url: jdbc:TAOS-RS://localhost:6041/vehicle
  username: root
  password: taosdata
  • broker-url:EMQX 的 MQTT 端口,tcp 协议
  • client-id:每个客户端唯一,重复登录会把前一个踢掉
  • topic: vehicle/+/data+ 是 MQTT 的单层通配符,匹配 vehicle/001/datavehicle/002/data……任意一辆车的数据都能收到
  • qos: 1:至少送达一次,车况数据这个级别够用

四、写 MQTT 订阅代码

项目结构

src/main/java/com/wang/
├── VehicleApplication.java              ← 启动类
├── config/
│   └── MqttProperties.java             ← MQTT 配置
├── mqtt/
│   └── MqttSubscriber.java             ← MQTT 订阅
└── service/
    └── TDengineService.java            ← TDengine 写入

4.1 读取配置

package com.wang.config;

import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.stereotype.Component;

/**
 * @author xiaoman
 * @date 2026/7/5 12:21
 */
@Data
@Component
@ConfigurationProperties(prefix = "mqtt")
public class MqttProperties {
    private String brokerUrl;
    private String clientId;
    private String topic;
    private int qos;
}

4.2 核心:MqttSubscriber

package com.wang.mqtt;

import com.wang.config.MqttProperties;
import com.wang.service.TDengineService;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.*;
import org.springframework.stereotype.Component;

import javax.annotation.PostConstruct;
import java.nio.charset.StandardCharsets;

/**
 * @author xiaoman
 * @date 2026/7/5 12:21
 */
@Slf4j
@Component
public class MqttSubscriber {

    private final MqttProperties mqttProperties;
    private final TDengineService tdengineService;

    public MqttSubscriber(MqttProperties mqttProperties,
                          TDengineService tdengineService) {
        this.mqttProperties = mqttProperties;
        this.tdengineService = tdengineService;
    }

    @PostConstruct
    public void connect() {
        try {
            MqttClient client = new MqttClient(
                mqttProperties.getBrokerUrl(),
                mqttProperties.getClientId()
            );

            MqttConnectOptions options = new MqttConnectOptions();
            options.setAutomaticReconnect(true);       // 断线自动重连
            options.setCleanSession(true);             // 每次新建会话
            options.setConnectionTimeout(10);

            client.connect(options);
            log.info("MQTT 连接成功");

            client.subscribe(mqttProperties.getTopic(), mqttProperties.getQos(),
                (topic, message) -> handleMessage(topic, message));
        } catch (Exception e) {
            log.error("MQTT 连接失败", e);
        }
    }

    private void handleMessage(String topic, MqttMessage message) {
        String payload = new String(message.getPayload(), StandardCharsets.UTF_8);
        log.info("收到消息 topic={}, payload={}", topic, payload);

        // 解析 JSON → 写入 TDengine
        tdengineService.saveVehicleData(payload);
    }
}

注意:JDK 8 / Spring Boot 2.x 用 javax.annotation.PostConstruct,JDK 17+ / Spring Boot 3.x 改成 jakarta.annotation.PostConstruct

@PostConstruct:Spring 启动后自动执行 connect(),建立 MQTT 长连接。

MqttConnectOptions

  • setAutomaticReconnect(true):网络断了自动重连,车联网场景必须开
  • setCleanSession(true):每次连接清空之前的未消费消息。生产环境可能要改成 false,这样设备离线后 Broker 会缓存消息,重连后补发

client.subscribe:订阅 vehicle/+/data,收到消息回调 handleMessage。这个回调在 MQTT 的接收线程里执行,不是主线程。

handleMessage:拿到的 message.getPayload()byte[],转成 UTF-8 字符串后传给 tdengineService 存库。


五、写 TDengine 存储代码

5.1 JDBC 工具类

package com.wang.service;

import com.alibaba.fastjson.JSON;
import lombok.extern.slf4j.Slf4j;
import com.alibaba.fastjson.JSONObject;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;

import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;

/**
 * @author xiaoman
 * @date 2026/7/5 12:21
 */
@Slf4j
@Component
public class TDengineService {

    private final String url;
    private final String username;
    private final String password;

    public TDengineService(
            @Value("${tdengine.url}") String url,
            @Value("${tdengine.username}") String username,
            @Value("${tdengine.password}") String password) {
        this.url = url;
        this.username = username;
        this.password = password;
    }

    public void saveVehicleData(String jsonPayload) {
        try (Connection conn = DriverManager.getConnection(url, username, password)) {

            JSONObject json = JSON.parseObject(jsonPayload);
            String vin = json.getString("vin");
            int speed = json.getIntValue("speed");
            int battery = json.getInteger("battery");
            double lat = json.getDouble("lat");
            double lng = json.getDouble("lng");

            //String sql = "INSERT INTO vehicle.v_" + vin
            //        + " VALUES (NOW, ?, ?, ?, ?)";
            // 用超级表 USING 语法,同VIN的第一个设备自动创建子表
            String sql = "INSERT INTO vehicle.v_" + vin
                    + " USING vehicle.vehicles TAGS('" + vin + "')"
                    + " VALUES (NOW, ?, ?, ?, ?)";
           try (PreparedStatement ps = conn.prepareStatement(sql)) {
                ps.setInt(1, speed);
                ps.setInt(2, battery);
                ps.setDouble(3, lat);
                ps.setDouble(4, lng);
                ps.executeUpdate();
            }

            log.info("数据入库成功 vin={}", vin);
        } catch (Exception e) {
            log.error("数据入库失败 payload={}", jsonPayload, e);
        }
    }
}

子表名拼接vehicle.v_ + vin,上一篇建子表就是这个命名规则,直接拼。

PreparedStatement:预编译比字符串拼接快,防 SQL 注入。

JSON 解析:fastjson,pom.xml 里已经加了。

5.2 驱动注册

在启动类上加一行注册 TDengine 驱动:

package com.wang;

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;

/**
 * @author xiaoman
 * @date 2026/7/5 12:21
 */
@SpringBootApplication
public class VehicleApplication {
    public static void main(String[] args) {
        // 注册 TDengine JDBC 驱动
        try {
            Class.forName("com.taosdata.jdbc.TSDBDriver");
        } catch (ClassNotFoundException e) {
            throw new RuntimeException(e);
        }
        SpringApplication.run(VehicleApplication.class, args);
    }
}

不加也能跑,但手动注册一下稳一点。


六、验证

6.1 启动项目

启动 Spring Boot,日志里看到这一条:

MQTT 连接成功

说明代码已经连上 EMQX 了。

6.2 打开 EMQX Dashboard

http://localhost:18083 → 连接管理 → 能看到 vehicle-server-001 这个客户端上线了。

6.3 用 MQTTX 发消息

  • Topic: vehicle/001/data
  • Payload:
{"vin":"VIN001","speed":88,"battery":72,"lat":30.5928,"lng":114.3055}

发送之后,IDEA 控制台输出:

收到消息 topic=vehicle/001/data, payload={"vin":"VIN001",...}
数据入库成功 vin=VIN001

6.4 查 TDengine

打开 http://localhost:6060,执行:

SELECT * FROM vehicle.v_vin001 ORDER BY ts DESC LIMIT 5;

能看到刚才发的 88km/h 那条数据就在第一条。整条链路跑通了。


七、踩坑记录

7.1 连接不上 EMQX

报错 Unable to connect to server——先确认 Docker 里 EMQX 在运行:

docker ps | findstr emqx

如果没起来,重新启动:docker start emqx

7.2 收到消息但入库报错 “Table does not exist”

子表 v_vin001 还没建。回上一篇的 Step 4 先把子表建好。每次接入新车都要先建子表,TDengine 超级表就是这样,或者用代码里的超级表 USING 语法

7.3 重复收到同一条消息

setCleanSession(false) 改成 true 就行。纯 false 的时候 Broker 会把离线期间缓存的消息全部重放一遍。


车联网学习笔记,持续更新。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐