车联网学习笔记:Spring Boot 整合 MQTT + TDengine,用代码接管消息收发
上一篇用 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/data、vehicle/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 会把离线期间缓存的消息全部重放一遍。
车联网学习笔记,持续更新。
更多推荐

所有评论(0)