
简介基于SpringBootMyBatis构建的物联网数据采集系统服务器端源码适合熟悉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.xmldependency groupIdorg.springframework.integration/groupId artifactIdspring-integration-mqtt/artifactId version5.5.15/version /dependency dependency groupIdorg.eclipse.paho/groupId artifactIdorg.eclipse.paho.client.mqttv3/artifactId version1.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 保留未消费的消息设备侧不用感知服务端重启后者是断线自动重连省掉了自己写重连循环。keepAliveInterval30表示 30 秒一次心跳如果设备侧心跳间隔设置得更短这里要跟着改否则 broker 会误判离线。Inbound 适配器订阅的主题是//data这种通配符格式两个加号分别匹配产品标识和设备标识比如/suzhu-machine/dev001-up. 后面所有设备上报的数据都会进入mqttInputChannel()再转交给下游的 MessageHandler 做解析。到这里接入层就算立住了接下来要解决的是数据往哪存、怎么存得高效。2.2 服务端核心InfluxDB 存储时序数据Redis 缓存设备状态设备上报的数据是典型的时序数据同一设备同一指标每秒钟或每分钟产生一条记录。我之前见过有人用 MySQL 一张大表存所有设备数据三个月后单表两千万行查一条曲线要七八秒索引再优化也救不回来。所以源码里实测路径是MySQL 只存设备元信息、用户、指令记录这些关系型数据真正的采集数据写进 InfluxDB。InfluxDB 的写入用官方 clientMeasurement 命名device_dataTags 放deviceId和metricField 放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 天。真实项目里我一般把降精度查询交给连续查询不能全量留存原始数据。设备在线状态不适合频繁读写 MySQLRedis 是最省事的选择。设备每上报一条消息接入层就刷新一次 keystringRedisTemplate.opsForValue().set(device:online: deviceId, 1, 90, TimeUnit.SECONDS);90秒是一个保守的过期时间只要设备还在上报key 就会续期一旦设备断电90 秒后 key 自动消失接口查询在线状态时读到null就判定离线。这个时间要和心跳周期匹配设备 60 秒一次心跳过期时间设 90 秒比较合理太短会造成误判离线太长则离线感知太慢。2.3 数据解析JSON 上报与二进制协议的自适应处理接入层收进来的都是 MQTT 的字节数组但不同设备上报格式天差地别。便宜的 DTU 可能发 JSONPLC 网关可能发十六进制帧。如果每个设备型号都往 Handler 里塞一套 if-else代码很快就烂了。源码里的做法是定义统一的 DataParser 接口再把不同解析器塞进工厂public interface DataParser { ParserResult parse(byte[] payload); }工厂类根据设备型号取对应解析器Component public class ParserFactory { private MapString, 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.xMaven 3.6.xMySQL 8.0Redis 5.xInfluxDB 2.6EMQX 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_PASSWORD123456 -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?useUnicodetruecharacterEncodingutf8 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、passwordRedis 的hostMQTT 的url和client-idInfluxDB 的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, qos1) time.sleep(2) client.disconnect()脚本里client.publish(/sim-device/dev001/data, payload, qos1)的主题格式是/{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、当前时间戳和密钥拼接后做 MD5broker 侧鉴权插件用同样的方式计算并比对。这样做的好处是密钥本身不直接出现在 MQTT 报文里截获链表数据也没法拿去伪造登录。timestamp参与计算后服务端要校验时间戳与当前时间的偏差我一般允许前后 5 分钟超过就拒绝连接。而device.secret-expire-hours24控制的是密钥本身的有效期过期后设备必须调注册接口重新获取以此应对密钥泄露。这里最大的坑是设备本地时钟不准。如果设备 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并在接本文还有配套的精品资源点击获取