基于Netty与Spring Boot的物联网设备接入中间件设计实践

基于Netty与Spring Boot的物联网设备接入中间件设计实践 简介iot-ucy是一套使用Java语言开发的物联网网络中间件基于Netty、Spring Boot、Redis等主流开源技术构建面向物联网后端开发者、嵌入式工程师及边缘计算研究人员主要解决多协议设备接入、链路维护与数据汇聚等通用问题。压缩包共631个文件、约600KB以607个Java源码文件为主配有少量XML配置、SQL脚本、Markdown说明和Spring Boot自动装配工厂文件便于按模块阅读和二次开发。目前已支持TCP、UDP、MQTT、MQTT网关、WebSocket、Modbus、DTUAT协议、DTUModbusTCP/RTU、西门子与欧姆龙PLC以及串口等常见物联网协议并可快速对接Redis、EMQX、TDengine等中间件。从内容预览可见MqttClient、ModbusDtuTestHandle、ByteUtil、SocketClient等关键类覆盖协议解析、设备管理、异步通信等核心逻辑适合需要搭建物联网接入层或深入研习中间件源码的开发者参考。已有319人学习与下载可作为实战项目的重要借鉴。1. iot-ucy 是什么面向设备接入的 Java 物联网网络中间件设备接入层最常见的坑是拿 Tomcat 挂 HTTP 接口收设备数据。几千台设备长连接一撑线程池直接占满业务代码里还到处是 Socket。iot-ucy 是这种场景下的常见解法用 Java 语言以 Netty 做网络内核Spring Boot 做装配Redis 保存在线状态。它本质是网络中间件把设备端 TCP 连接、私有二进制协议、心跳保活、会话状态管理从业务代码里剥出来。做 IoT 平台时业务系统只关心设备上报的数据不关心每个连接怎么保活这正是它存在的位置。适合的人群是正在做设备接入网关、需要自研私有 TCP 协议接入的团队。下面按 Netty Spring Boot Redis 这条主线把能落地的代码、参数和坑位讲清楚新手可以直接复刻老手可以重点看鉴权和降级这几节。2. iot-ucy 的技术底座Netty 负责连接Spring Boot 负责装配Redis 负责状态2.1 为什么用 Netty 而不是 Tomcat 或传统 BIO Socket连接型中间件最重要的选型依据是连接模型。Tomcat 的线程池按 HTTP 请求设计来一个请求占一个线程设备长连接一挂就是几小时如果让 Tomcat 接受设备连接几千个设备就把线程池占满重启和超时问题会连锁放大。传统 BIO Socket 更不用提一连接一线程在物联网场景下几乎不可用。Netty 是 NIO 事件驱动模型Boss 线程只负责 acceptWorker 线程处理读写事件一个线程可以管理上千条闲置连接。这个差异直接决定了在同样的 4C8G 机器上Netty 服务端能撑住的设备连接数量比 Tomcat 高一个量级。方案连接模型协议适合场景BIO Socket一连接一线程任意自定义协议几十台设备量的临时服务TomcatServlet 线程池HTTP/HTTPSWeb API、设备走 HTTP 短连接NettyNIO 事件循环TCP/UDP/私有协议高并发长连接、自定义报文在 iot-ucy 里直接用 netty-all 而不是自己写 select 循环原因就是 NIO 的边界情况非常多半包、写缓冲、内存泄漏现成框架已经把主流问题兜住了剩下的粘包和业务协议才是自己的活。2.2 Spring Boot 在中间件里的定位很多人看到 Spring Boot 就先想到 Web MVC但它在 iot-ucy 里不是给设备跑 HTTP 的设备不会从 Controller 进来。Spring Boot 在这里做的是依赖装配把 NettyServer 的启停纳入 Spring 生命周期把 RedisTemplate、事件发布器、配置参数注入到各个 Handler再暴露一组管理 HTTP 接口用于查询在线设备。所以不能照搬 Controller-Service-Dao 的四层架构来做网络中间件。连接层里核心是 ChannelHandler它处理的是二进制帧而不是 RequestBodyService 层处理业务逻辑Redis 和数据库才在更底层。控制层只负责管理接口比如断开某台设备、查看某节点连接数这部分才完整走 Spring MVC。Spring Boot 的另一个好处是配置优先级和配置项绑定。设备端口、心跳超时、Redis key 前缀这些全部可以放进 application.yml用 ConfigurationProperties 映射成 IotUcyProperties这样中间件切换环境时不用改代码只换配置。2.3 Redis 使用的数据类型与 Key 设计中间件跨节点部署后Netty 的 Channel 只存在于本机内存别的节点不知道某台设备到底连在哪个进程上。Redis 在这里承担的是共享状态源每个节点的启停都要以它为准。在 iot-ucy 里Redis 数据类型可以对应到具体职责上数据类型Key 示例说明在线状态Stringiot:online:{deviceId}value 是 channelId带 90 秒过期设备鉴权信息Stringiot:token:{token}登录时生成value 是 deviceId设备元数据Hashiot:device:{deviceId}固件版本、IP、最后上线时间离线指令队列Listiot:offline:{deviceId}设备离线时业务推送的消息先 LPUSH下发锁Stringiot:lock:send防止多节点重复下发同一条指令String 的过期时间就是天然的下线兜底设备宕机来不及发断开包90 秒后 key 消失在线状态自动翻转。Hash 适合存上报的属性快照List 适合做离线消息这两种数据类型的读写都是 O(1)在连接事件高频发生时性能非常关键。2.4 最小工程依赖配置先搭出中间件的最小工程依赖要稳定且不贪多。pom.xml 里通常只保留这几个dependency groupIdio.netty/groupId artifactIdnetty-all/artifactId version4.1.100.Final/version /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependencynetty-all 把传输、编解码、处理器都带齐避免手动拼多个 netty 子模块出现版本冲突。spring-boot-starter-data-redis 提供 StringRedisTemplate 和 RedisTemplate连接池默认用 Lettuce不需要额外引 Jedis。spring-boot-starter-web 只是为了让管理 API 能跑如果这个中间件只想做纯接入也可以去掉它只留 Netty 和 Redis。如果不想让 Redis 故障导致整个中间件启动失败可以设置短连接超时spring: data: redis: timeout: 3s connect-timeout: 3s启动时的连接检查会变快但这只是缓解。真正对 Redis 故障的降级策略放到最后一章说。3. 实现 iot-ucy 核心连接层Netty 粘包处理与心跳超时3.1 私有协议与 LengthFieldBasedFrameDecoder设备端上报的场景TCP 会按照底层缓冲把多个包合并或者把一个包拆成两次 send这就是粘包/拆包。Netty 自己不做业务协议只提供把字节流切成完整帧的组件。iot-ucy 里最常见的协议设计是 3 字节固定头加 4 字节长度字段再跟 payload。字段字节数值magic10x5Aversion20x0001length4payload 长度payloadlength实际业务数据对应的服务端初始化代码ServerBootstrap b new ServerBootstrap(); b.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .childHandler(new ChannelInitializerSocketChannel() { Override protected void initChannel(SocketChannel ch) { ChannelPipeline p ch.pipeline(); p.addLast(new LengthFieldBasedFrameDecoder(65535, 3, 4, 0, 0)); p.addLast(new ByteToMessageDecoder() { Override protected void decode(ChannelHandlerContext ctx, ByteBuf in, ListObject out) { in.skipBytes(7); byte[] payload new byte[in.readableBytes()]; in.readBytes(payload); out.add(new IotMessage(payload)); } }); p.addLast(new IotMessageHandler()); } });LengthFieldBasedFrameDecoder 的 5 个参数是这类代码里最需要说清楚的部分maxFrameLength65535 是单帧上限防止恶意长度头导致内存爆掉lengthFieldOffset3 表示长度字段在 magic 和 version 之后lengthFieldLength4 表示长度本身占 4 字节lengthAdjustment0 表示长度字段的值只指 payload不含前面的 7 字节头部initialBytesToStrip0 表示完整帧保留帧头后面的 decoder 跳过 7 字节。这个配置下如果设备上报 5 个包粘在一起Netty 会还给你 5 个完整 ByteBuf业务 handler 不会看到半包。3.2 心跳处理IdleStateHandler 与下线兜底长连接如果没有心跳设备死机、网线断开这类情况双方可能很久都发现不了。Netty 的 IdleStateHandler 可以在 pipeline 里定时检查读写空闲iot-ucy 里一般放在解码器之前p.addLast(new IdleStateHandler(90, 0, 0, TimeUnit.SECONDS)); p.addLast(new HeartbeatHandler());这段配置表示 90 秒内没有收到设备端任何数据就触发 READER_IDLE 事件。第 2、3 个参数是写空闲和全空闲这里设 0 表示不检查。心跳包本身是业务数据类别服务端只做读空闲判断。心跳 handler 里做两件事删除 Redis 在线状态 key关闭连接。public class HeartbeatHandler extends ChannelInboundHandlerAdapter { Override public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception { if (evt instanceof IdleStateEvent) { IdleStateEvent e (IdleStateEvent) evt; if (e.state() IdleState.READER_IDLE) { SessionManager.remove(ctx.channel()); ctx.close(); } } } }Redis 里的 key 在上线时设置了和读空闲一致的过期时间所以即使这个 handler 来不及执行key 也会自己过期形成双保险。为什么设 90 秒而不是 600 秒很多设备断网时不发 FIN 包读空闲要小于 TCP 的默认超时才有可能感知但也不能太短移动网络下设备偶尔会静默一个周期90 秒是多个项目里测下来的折中值。3.3 会话注册与 Redis 状态同步服务端要下发指令时得根据 deviceId 找到对应的 Channel这个映射关系不能每次遍历所有连接。一般用 ConcurrentHashMap 做本机会话表Redis 做全局会话表。public class SessionManager { private static final ConcurrentHashMapString, Channel CHANNELS new ConcurrentHashMap(); public static void add(String deviceId, Channel ch) { CHANNELS.put(deviceId, ch); } public static Channel get(String deviceId) { return CHANNELS.get(deviceId); } public static void remove(Channel ch) { CHANNELS.entrySet().removeIf(e - e.getValue() ch); } }设备鉴权通过后这样注册SessionManager.add(deviceId, ctx.channel()); StringRedisTemplate redis SpringContextHolder.getBean(StringRedisTemplate.class); redis.opsForValue().set(iot:online: deviceId, ctx.channel().id().asLongText(), 90, TimeUnit.SECONDS);注意 remove 方法用 Channel 实例作为判断条件因为同一 deviceId 可能因为重连被新的 Channel 覆盖用 deviceId 删除有可能误删新连接。本机会话表解决单节点内 O(1) 找 ChannelRedis 解决的是业务系统查询设备是否在线时不连接中间件节点也能直接查 Redis。4. iot-ucy 与 Spring Boot 集成鉴权、事件消息与 Redis 序列化4.1 TCP/WebSocket 首包鉴权物联网设备接入不管是 TCP 还是 WebSocket鉴权都不能等业务报文到了才做。常见方案是连接建立后第一个业务帧必须是鉴权包携带 token。Netty 里放一个专门的鉴权 Handler它处理完后从 pipeline 里把自己删掉后续消息不再走鉴权逻辑。public class AuthHandler extends ChannelInboundHandlerAdapter { private final StringRedisTemplate redis; public AuthHandler(StringRedisTemplate redis) { this.redis redis; } Override public void channelRead(ChannelHandlerContext ctx, Object msg) { IotMessage m (IotMessage) msg; String token new String(m.getPayload(), StandardCharsets.UTF_8); String deviceId redis.opsForValue().get(iot:token: token); if (deviceId null) { ctx.close(); return; } SessionManager.add(deviceId, ctx.channel()); ctx.pipeline().remove(this); ctx.fireChannelRead(msg); } }如果设备是 WebSocket 接入鉴权 token 一般放在握手 URL 参数里在 HttpServerCodec 之后的 HttpRequestHandler 里读取 query校验逻辑完全一样。这个方案在多个中间件节点共享 Redis 时天然有效任何节点都可以处理任意设备的连接请求不需要把连接固定到某个节点。4.2 用 Spring 事件解耦消息上行设备报文经过解码器后如果直接在 Netty 的 handler 里调用业务 Service网络线程会被长耗时的数据库操作卡住。Netty 的 worker 线程很宝贵不能让一个设备上报把其他设备的消息滞后。用 Spring 的 ApplicationEventPublisher 把消息发出去Component public class IotEventPublisher { private final ApplicationEventPublisher publisher; public IotEventPublisher(ApplicationEventPublisher publisher) { this.publisher publisher; } public void publish(String deviceId, IotMessage msg) { publisher.publishEvent(new IotUpstreamEvent(this, deviceId, msg)); } } public class IotUpstreamEvent { private final Object source; private final String deviceId; private final IotMessage message; public IotUpstreamEvent(Object source, String deviceId, IotMessage message) { this.source source; this.deviceId deviceId; this.message message; } // getter }IotEventPublisher 注入到 IotMessageHandler 中在 channelRead0 里调用 publish 后立即返回。业务侧只需要写一个 EventListener 方法监听 IotUpstreamEvent再把消息转换后写入 Kafka 或数据库。同步事件和异步事件的选择可以这样看方式Netty 线程是否阻塞异常处理直接在 Handler 调用 Service是在 EventLoop 里抛异常同步 Spring 事件是监听器异常会影响发布线程Async 异步事件否必须单独处理异常默认情况下 Spring 事件监听器是同步执行的如果想让 Netty 线程彻底不被阻塞监听器方法上要标注 Async并配置一个独立的线程池不要让事件落回 Netty 的 worker 线程。4.3 序列化方式与 key 乱码Redis 操作最常见的问题是用错 RedisTemplate 的序列化器。Spring Boot 自动装配的 RedisTemplate 的 value 默认使用 JdkSerializationRedisSerializer存进去的对象在 Redis Desktop Manager 里显示成 \xAC\xED 开头的一串乱码而且跨语言读取很困难这是 redis 序列化里非常典型的坑。iot-ucy 里保存 token、deviceId、channelId 这类纯字符串建议直接注入 StringRedisTemplate它把 key 和 value 都按 UTF-8 字符串处理。如果一定要用 RedisTemplate 存对象就明确配置序列化器Bean public RedisTemplateString, Object redisTemplate(RedisConnectionFactory factory) { RedisTemplateString, Object template new RedisTemplate(); template.setConnectionFactory(factory); template.setKeySerializer(new StringRedisSerializer()); template.setValueSerializer(new GenericJackson2JsonRedisSerializer()); template.setHashKeySerializer(new StringRedisSerializer()); template.setHashValueSerializer(new GenericJackson2JsonRedisSerializer()); template.afterPropertiesSet(); return template; }注意 Jackson 序列化对象会在 JSON 里带上 class 字段反序列化时依赖 class 信息。如果中间件升级时实体类包名变了老数据会反序列化失败。所以对中间件自身的 Redis 数据尽量用 StringRedisTemplate只有给业务系统复用的扩展数据才走对象序列化。5. iot-ucy 上线前的验证与调优5.1 用 Python 模拟粘包和半包验证解码器是否正确最直接的方式是用原始 socket 发畸形数据。下面这段脚本先发送一个少最后 2 字节的帧间隔 200ms 再发送剩余字节和下一个完整帧用来模拟半包和粘包同时出现import socket, time, struct s socket.socket(socket.AF_INET, socket.SOCK_STREAM) s.connect((127.0.0.1, 8100)) payload bping-001 frame b\x5A\x00\x01 struct.pack(I, len(payload)) payload s.sendall(frame[:-2]) time.sleep(0.2) s.sendall(frame[-2:] frame) time.sleep(1) s.close()如果解码器配置正确服务端日志应该打印两条业务消息而不是一条错乱的帧。如果只打印一条或者客户端直接收到重置优先检查 lengthFieldAdjustment 和 initialBytesToStrip 的计算。5.2 关键参数与压测观察Netty 参数不是越大越好。设备接入网关时最常调整的几项参数建议起始值说明bossGroup 线程数1只做 accept多余线程没有收益workerGroup 线程数CPU 核数 * 2处理读写事件和心跳SO_BACKLOG1024突发连接积压队列WRITE_BUFFER_WATER_MARK64KB / 256KB触发写缓冲高水位防止慢消费者拖死内存压测不用一步到位上万连接。先用 1000 个 socket 每秒发一次心跳观察 worker 线程的 CPU 占比和连接出现 READER_IDLE 的比例。如果设备端大量无法注册先去调 SO_BACKLOG再看 Redis 读写是否成为瓶颈。这种自测脚本能替代一部分 Netty 面试题里的空谈把问题落在实际参数上。5.3 Redis 故障降级与时间窗口设置Redis 不可用时中间件不能直接把设备连接杀掉。给连接层加降级开关Redis 异常时只记录日志设备仍在本机 SessionManager 里保持在线跨节点的在线查询暂时不可用但不影响设备正常上报。等 Redis 恢复后再补偿写入当前节点的在线状态避免业务侧误判。最后留一个细节在读空闲超时和 Redis Key 过期时间设置上两者必须一致且要留出至少 10 秒的设备重传余量。常见做法是把设备心跳周期设 30 秒读空闲设 90 秒Redis 在线状态过期也设 90 秒。这样即使设备在两次心跳之间掉线中间件最多 90 秒就能感知业务侧查在线状态也不会出现 Redis 里还活着但连接已断开的窗口。本文还有配套的精品资源点击获取