☰
SpringBoot物联网数据采集服务端源码解析:从MQTT接入到InfluxDB存储
2026/10/2 1:08:09 网站建设 项目流程

简介:基于SpringBoot+MyBatis构建的物联网数据采集系统服务器端源码,适合熟悉Java Web与物联网基础、希望掌握企业级架构的开发者。项目大幅减少xml配置,仅需在application.yml中做少量设置,并集成Redis缓存:单点查询结果缓存、传感器Data写入缓存队列、登录信息存Redis实现分布式session共享,同时通过线程池异步将缓存数据落库,降低数据库写入压力。内置Tomcat便于集群部署,附带IP/端口查看API,可配合nginx反向代理与负载均衡测试。资源共94个文件,以48个Java源码为核心,辅以25个HTML页面、8个XML配置、5个JS脚本及SQL等,压缩包仅644KB,结构清晰便于学习。已有466人学习,适合作为SpringBoot综合项目实践参考,可从中获得缓存设计、异步任务、集群会话共享等关键思路。

1. SpringBoot 物联网数据采集服务器端源码值不值得下:先看完这张图再动手

车间里几十台注塑机、plc 或传感器盒子定时上报温度、压力和运行状态,后端要稳定接收、解析、存储并提供查询和指令下发,这活儿看着简单,真做起来坑不少。这份基于 SpringBoot 框架的物联网数据采集系统服务器端源码,解决的就是设备接入、数据解析、时序存储、在线状态和指令下发这一条完整链路。它不是那种只跑通的 demo,而是把 MQTT 接入、InfluxDB 存储、设备鉴权、命令下发这些常用模块都放好了,适合刚接手物联网后端、拿它做毕业设计,或者想把自己那套 TCP 长连接协议改造成标准 MQTT 方案的开发者。下文按“框架怎么立—怎么跑起来—核心链路实现—踩坑记录—进阶改造”的顺序拆,读之前建议先把 MySQL、Redis、InfluxDB 和 EMQX 这几个组件准备好。

2. 先立住框架:SpringBoot 物联网服务端的模块划分与三层架构落地

2.1 设备接入层:从 MQTT Broker 到 Handler 的必经链路

物联网服务端和普通 Web CRUD 最大的区别是入口不是 HTTP,而是 MQTT。我一般不会自己用 Netty 去怼 TCP,因为设备断线重连、消息超时重发、主题订阅这些底层能力,MQTT Broker 已经做得很成熟。源码里选的是 Eclipse Paho + Spring Integration MQTT,先把依赖加进 pom.xml:

<dependency> <groupId>org.springframework.integration</groupId> <artifactId>spring-integration-mqtt</artifactId> <version>5.5.15</version> </dependency> <dependency> <groupId>org.eclipse.paho</groupId> <artifactId>org.eclipse.paho.client.mqttv3</artifactId> <version>1.2.5</version> </dependency>

依赖只是第一步,真正的接入逻辑在 MqttConfig 里。这里控制着客户端怎么连 broker、订阅哪些主题、断线后怎么办:

@Configuration public class MqttConfig { @Value("${mqtt.broker.url}") private String brokerUrl; @Value("${mqtt.client.id}") private String clientId; @Value("${mqtt.topic.filter}") private String topicFilter; @Bean public MqttPahoClientFactory mqttClientFactory() { DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory(); MqttConnectOptions options = new MqttConnectOptions(); options.setCleanSession(false); options.setAutomaticReconnect(true); options.setMaxReconnectDelay(30000); options.setKeepAliveInterval(30); factory.setConnectionOptions(options); return factory; } @Bean public MessageProducer inbound() { MqttPahoMessageDrivenChannelAdapter adapter = new MqttPahoMessageDrivenChannelAdapter(clientId, mqttClientFactory(), topicFilter); adapter.setCompletionTimeout(5000); adapter.setQos(1); adapter.setOutputChannel(mqttInputChannel()); return adapter; } }

这段配置里最值得关注的是setCleanSession(false)和setAutomaticReconnect(true)。前者保证服务端短暂重启时,broker 还会为这个 clientId 保留未消费的消息,设备侧不用感知服务端重启;后者是断线自动重连,省掉了自己写重连循环。keepAliveInterval=30表示 30 秒一次心跳,如果设备侧心跳间隔设置得更短,这里要跟着改,否则 broker 会误判离线。

Inbound 适配器订阅的主题是+/+/data这种通配符格式,两个加号分别匹配产品标识和设备标识,比如/suzhu-machine/dev001-up. 后面所有设备上报的数据都会进入mqttInputChannel(),再转交给下游的 MessageHandler 做解析。到这里接入层就算立住了,接下来要解决的是数据往哪存、怎么存得高效。

2.2 服务端核心:InfluxDB 存储时序数据,Redis 缓存设备状态

设备上报的数据是典型的时序数据:同一设备同一指标,每秒钟或每分钟产生一条记录。我之前见过有人用 MySQL 一张大表存所有设备数据,三个月后单表两千万行,查一条曲线要七八秒,索引再优化也救不回来。所以源码里实测路径是:MySQL 只存设备元信息、用户、指令记录这些关系型数据,真正的采集数据写进 InfluxDB。

InfluxDB 的写入用官方 client,Measurement 命名device_data,Tags 放deviceId和metric,Field 放value,时间戳用毫秒。核心写入代码如下:

@Service public class InfluxService { @Value("${influxdb.bucket}") private String bucket; private final InfluxDBClient influxDBClient; public void writePoint(String deviceId, String metric, double value, long timestampMs) { Point point = Point.measurement("device_data") .addTag("deviceId", deviceId) .addTag("metric", metric) .addField("value", value) .time(timestampMs, WritePrecision.MS); influxDBClient.getWriteApiBlocking().writePoint(bucket, "iot-rp", point); } }

写入时有两个参数要留意:WritePrecision.MS表示时间戳用毫秒精度,设备上报时原始时间单位如果是秒,要先乘以 1000,否则时间轴会乱;bucket 对应 InfluxDB 2.x 的存储桶,1.x 里则是 database + retention policy,源码里这个iot-rp就是保留策略,默认 30 天。真实项目里我一般把降精度查询交给连续查询,不能全量留存原始数据。

设备在线状态不适合频繁读写 MySQL,Redis 是最省事的选择。设备每上报一条消息,接入层就刷新一次 key:

stringRedisTemplate.opsForValue().set("device:online:" + deviceId, "1", 90, TimeUnit.SECONDS);

90秒是一个保守的过期时间,只要设备还在上报,key 就会续期;一旦设备断电,90 秒后 key 自动消失,接口查询在线状态时读到null就判定离线。这个时间要和心跳周期匹配,设备 60 秒一次心跳,过期时间设 90 秒比较合理,太短会造成误判离线,太长则离线感知太慢。

2.3 数据解析:JSON 上报与二进制协议的自适应处理

接入层收进来的都是 MQTT 的字节数组,但不同设备上报格式天差地别。便宜的 DTU 可能发 JSON,PLC 网关可能发十六进制帧。如果每个设备型号都往 Handler 里塞一套 if-else,代码很快就烂了。源码里的做法是定义统一的 DataParser 接口,再把不同解析器塞进工厂:

public interface DataParser { ParserResult parse(byte[] payload); }

工厂类根据设备型号取对应解析器:

@Component public class ParserFactory { private Map<String, DataParser> parserMap = new HashMap<>(); public DataParser getParser(String deviceModel) { DataParser parser = parserMap.get(deviceModel); if (parser == null) { throw new IllegalArgumentException("不支持的设备型号: " + deviceModel); } return parser; } }

JSON 解析器处理{"temperature":26.5,"pressure":1.2,"ts":1700000000000}这类消息,二进制解析器则先读帧头、长度字段和 CRC 校验,再按协议字段偏移量取值。在 MessageHandler 里先拿到 Topic 里的设备型号,再交给工厂选择解析器,这样新增一种设备协议时,只需要新写一个实现类并注册进工厂,原有代码完全不动。

这里最容易忽略的坑是:解析器不能只关注 payload,还要把 topic 里的 deviceId 和 productKey 一并塞进 ParserResult。因为很多协议本身的 payload 里不带设备号,全靠 topic 路由校验,解析完再手动覆盖时间戳和数据点。我把这个链路称为“topic 是设备身份证,payload 是数据本体”,两者必须在进入存储前合并成一条完整记录。

3. 把源码跑起来:建库建表、改配置、启动的三步复现

3.1 环境准备:JDK、Maven、MySQL、InfluxDB 与 EMQX 的选型

复现这份源码不用最新的服务,我用的是稳定组合:JDK 8 + Spring Boot 2.7.x,Maven 3.6.x,MySQL 8.0,Redis 5.x,InfluxDB 2.6,EMQX 4.4。选型理由很简单:这套组合网上排错资料最多,设备接入层的第三方库兼容性也最好。如果你机器上已经有 Docker,直接一条命令起中间件:

docker run -d --name emqx -p 1883:1883 -p 18083:18083 emqx/emqx:4.4.3 docker run -d --name influxdb -p 8086:8086 influxdb:2.6 docker run -d --name mysql -e MYSQL_ROOT_PASSWORD=123456 -p 3306:3306 mysql:8.0 docker run -d --name redis -p 6379:6379 redis:5-alpine

注意 EMQX 的18083是后台管理端口,1883才是 MQTT 端口。很多新手只映射了 18083 就跑去连 1883,结果是管理后台能打开,服务端却一直报连接超时。MySQL 启动后执行源码里的init.sql,这里面建了device_info、command_record、user_info三张表,另外记得把 root 密码换成你自己的。

3.2 核心配置:application.yml 里必须改的七个地方

中间件就绪后,配置是复现成功的关键。源码的application.yml长这样:

server: port: 8080 spring: datasource: url: jdbc:mysql://127.0.0.1:3306/iot_server?useUnicode=true&characterEncoding=utf8 username: root password: 123456 redis: host: 127.0.0.1 port: 6379 mqtt: broker: url: tcp://127.0.0.1:1883 client-id: iot-server-001 username: iot_user password: iot_pass topic: filter: +/+/data influxdb: url: http://127.0.0.1:8086 token: my-token org: iot bucket: iot/autogen device: secret-expire-hours: 24

按我的习惯,每次要修改的值正好是七个地方:MySQL 的 url 里的iot_server数据库名、username、password,Redis 的host,MQTT 的url和client-id,InfluxDB 的token。其中client-id必须要全局唯一,不能多套服务共用同一个,否则 EMQX 会把后连的踢掉。topic.filter保持+/+/data,如果你的设备上报主题是/factoryA/device001/upload/data,那这里就要改成+/+/+/data,对应的数据解析前取主题段位也要调整偏移量。

mqtt.username和mqtt.password是服务端连接 broker 的凭证,不是设备连接凭证。EMQX 4.x 默认关闭认证,这俩配了也不验证;但生产环境开了认证后,这组账号要在 EMQX 里单独创建,并只授予订阅权限。InfluxDB 的token是在初始化时生成的,遗漏会导致启动报 401。

3.3 启动与验证:用模拟设备压测数据链路

配置改完,启动 SpringBoot 主类。没有报错只代表启动成功,不代表数据链路通了。我习惯用一段 Python 脚本模拟设备,每秒上报一条数据,验证全链路是否闭合:

import paho.mqtt.client as mqtt import json import time client = mqtt.Client("sim_device_001") client.username_pw_set("device_001", "password") def on_connect(c, u, f, rc): print("connected:", rc) client.on_connect = on_connect client.connect("127.0.0.1", 1883, 60) client.loop_start() for i in range(10): payload = json.dumps({ "temperature": 25 + i, "pressure": 1.2, "ts": int(time.time() * 1000) }) client.publish("/sim-device/dev001/data", payload, qos=1) time.sleep(2) client.disconnect()

脚本里client.publish("/sim-device/dev001/data", payload, qos=1)的主题格式是/{productKey}/{deviceName}/data,和服务端的+/+/data通配符完全匹配。服务端收到后按 JSON 解析,再写 InfluxDB。此时到 MySQL 的device_info表里确认设备存在,再到 InfluxDB 的执行窗口输入from(bucket: "iot/autogen") |> range(start: -5m)查数据点。如果查不到,优先看 SpringBoot 日志里有没有 “message arrived” 的打点。这个验证动作,我从第一次跑物联网服务端开始就再没跳过。

4. 数据采集与下发:从设备注册、指令下发到分表查询的实现要点

4.1 设备注册与鉴权:token 过期与重连避坑

任何设备要上报数据,都得先在平台注册,拿到设备 ID 和密钥。源码里的注册接口会为设备生成deviceSecret,并基于 HMAC 算出一串动态密码,供设备连接 MQTT 时使用:

public String buildMqttPassword(String deviceId, String deviceSecret, long timestamp) { String raw = deviceId + timestamp + deviceSecret; return DigestUtils.md5Hex(raw); }

这段逻辑的原理是:设备把deviceId、当前时间戳和密钥拼接后做 MD5,broker 侧鉴权插件用同样的方式计算并比对。这样做的好处是密钥本身不直接出现在 MQTT 报文里,截获链表数据也没法拿去伪造登录。timestamp参与计算后,服务端要校验时间戳与当前时间的偏差,我一般允许前后 5 分钟,超过就拒绝连接。而device.secret-expire-hours=24控制的是密钥本身的有效期,过期后设备必须调注册接口重新获取,以此应对密钥泄露。

这里最大的坑是设备本地时钟不准。如果设备 RTC 快了 10 分钟,算出来的动态密码和服务端对不上,会出现“能上线,但每隔几小时掉线一次”的诡异现象。排查时先对比两端时间差,别一上来就怀疑鉴权逻辑。

4.2 指令下发:QoS 1 与应答超时重试机制

数据采集是上行,服务端还要能下行控制设备,比如远程开关、调整参数。指令下发不能想当然地直接mqttGateway.sendToMqtt,因为设备不在线时消息会直接丢失。源码里的做法是先查在线状态,不在线就落库等上线补发:

public void sendCommand(String deviceId, String command) { Device device = deviceMapper.selectById(deviceId); if (!isOnline(deviceId)) { commandService.saveWaiting(deviceId, command); return; } String topic = "/" + device.getProductKey() + "/" + deviceId + "/cmd"; mqttGateway.sendToMqtt(command, topic); commandService.markSent(deviceId, command); }

sendToMqtt默认走 QoS 1,保证消息至少送达一次。但“至少一次”不代表设备一定执行,设备收到后可能处理失败。所以源码里还维护了一张command_record表,记录每次下发的状态:SENT、ACKED、FAILED。服务端发出去后启动一个 30 秒定时任务,如果设备没回 ACK 主题,就把状态改成FAILED并重发,最多重试 3 次。这个重试次数要克制,否则设备反复收到重复指令,可能出现双重开启之类的故障。真实设备处理指令时通常要做去重,即根据指令里的消息 ID 判断是否已经执行过。

4.3 数据查询:按设备、时间范围分页与聚合的接口写法

数据采好了要能查。查询接口如果直接“SELECT * FROM device_data”这种思维去套 MySQL,那 InfluxDB 的优势就废了。源码里的 Flux 查询方式是这样:

@GetMapping("/api/v1/devices/{deviceId}/data") public Result listData(@PathVariable String deviceId, @RequestParam long timeStart, @RequestParam long timeEnd, @RequestParam int page, @RequestParam int size) { String flux = "from(bucket: \"iot/autogen\")" + " |> range(start: " + timeStart + ", stop: " + timeEnd + ")" + " |> filter(fn: (r) => r._measurement == \"device_data\" and r.deviceId == \"" + deviceId + "\")" + " |> sort(columns: [\"_time\"], desc: true)" + " |> limit(n: " + size + ", offset: " + (page - 1) * size + ")"; return success(influxDBClient.getQueryApi().query(flux)); }

range必须传毫秒时间戳,而且接口层就要强制校验timeEnd - timeStart不能超过 7 天。原因很简单:没有时间范围的 InfluxDB 查询会扫全库,每次页面刷新都触发一次全量扫描,服务端内存吃不住。limit(n, offset)实现了分页,但 offset 太大时效率下降,所以真实业务里我建议改为按时间游标分页,客户端传上次最后一条数据的时间,而不是页数。这个接口的参数校验逻辑直接决定了服务端能不能撑过三个月,不要省。

5. 避坑与排查:SpringBoot 物联网服务端最常见的五个翻车现场

5.1 连接与鉴权:频繁掉线和消息丢失

现象一:设备每隔几十秒就掉线重连,服务端日志反复出现 “Client sim_device_001 already connected”。

原因:模拟设备脚本和服务端 MQTT 客户端用了同一个 clientId。EMQX 对重复 clientId 的处理是后连接踢掉前连接,于是两台客户端不停地互踢。

解决:把设备 clientId 设为sim_device_001的 MAC 地址后缀,服务端 clientId 设为iot-server-001,保证全局唯一。如果设备数量超过几千,clientId 还要加上设备型号前缀,避免不同厂商设备 ID 撞车。

现象二:设备明明上报了数据,服务端却一条消息都没收到。

原因:服务端订阅的 topic filter 写成了/sim-device/dev001/data,精确匹配了单台设备,新产品上线时忘了扩通配符。

解决:订阅改成+/+/data或+/+/+/data,并在接

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询