ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

iot-ucy:面向高并发IoT设备的Netty+Redis中间件

iot-ucy:面向高并发IoT设备的Netty+Redis中间件 简介这是一套面向物联网开发工程师与后端架构师的轻量级网络中间件开源实现基于Java技术栈解决设备协议接入、数据路由与边缘服务集成等核心问题适用于智能硬件对接、工业网关开发及边缘计算场景。资源包共631个文件主体为607个Java类涵盖DeviceManagerFactory、MqttClient、ModbusRtuClientProtocol等关键模块辅以14个XML配置、2个SQL建表脚本及properties/factories等基础支撑文件整体仅600KB结构紧凑、开箱即用。已有324人学习下载可直接复用其多协议适配能力——完整支持TCP/UDP、MQTT及其网关、WebSocket、ModbusTCP/RTU、AT指令DTU、西门子/欧姆龙PLC及串口通信并已预置Redis、EMQX、TDengine等主流中间件的快速接入逻辑显著降低物联网平台协议层开发门槛。1. iot-ucy 是什么它不是“又一个 Spring Boot 项目”而是专为高并发设备连接、低延迟指令下发和状态强一致同步设计的物联网网络中间件你手头正跑着几百台 PLC、几十个边缘网关、上千个 NB-IoT 终端Spring Boot 写的 REST API 每秒扛不住 200 个连接请求Netty 手写 TCP 服务又缺设备生命周期管理、离线消息兜底、集群会话同步——这时候iot-ucy 就不是“可选”而是“必须”。它用 Java 语言构建但不走 Web MVC 路线底层用 Netty 实现百万级长连接承载与二进制协议解析支持自定义私有协议或 MQTT over TCP中间层用 Spring Boot 做模块化编排与配置治理非 Web 启动模式禁用 Tomcat数据层深度整合 Redis —— 不只是缓存而是用 Redis Streams 做指令分发队列、用 Redis Hash 存设备影子状态、用 Redis Sorted Set 管理在线会话心跳、用 Redis Lua 原子脚本实现分布式设备锁。它解决的不是“怎么把数据存下来”而是“当 372 台温控器同时上报温度、后台下发 128 条调参指令、其中 4 台正在断网重连”时系统不丢指令、不乱状态、不串设备、不阻塞新连接。适合嵌入式通信协议工程师、IoT 平台后端开发者、以及被“设备在线率波动”“指令下发超时率突增”“Redis 缓存击穿导致设备失联”反复折磨的运维同学。别把它当成 demo 工程——它的核心价值在于把 Netty 的性能、Spring Boot 的可维护性、Redis 的实时协同能力焊死在物联网设备连接这个黑匣子上。2. 从零启动 iot-ucy三步完成最小可运行环境搭建与基础连接验证iot-ucy 的启动逻辑和传统 Spring Boot Web 应用有本质区别它不依赖内嵌 Servlet 容器而是以 Netty Server 为主进程Spring Boot 仅作为 IOC 容器和配置中心。这意味着你不能靠spring-boot-starter-web启动也不能用RestController暴露 HTTP 接口来管理设备——所有设备接入、指令下发、状态查询都走 TCP 或 MQTT 协议通道。下面步骤基于官方推荐的 v1.4.2 版本截至 2024 年 Q2 最稳定分支适配 JDK 17、Spring Boot 2.7.x、Redis 7.x。2.1 下载源码并确认核心模块结构iot-ucy 采用多模块 Maven 结构关键模块如下无需修改即可运行模块名作用是否必需iot-ucy-coreNetty 服务主干、协议编解码器、连接管理器、心跳检测器✅ 必需iot-ucy-spring-boot-starterSpring Boot 自动装配入口注入DeviceSessionManager、CommandDispatcher等 Bean✅ 必需iot-ucy-redis-supportRedis 连接池初始化、Streams 消费组配置、Hash/Sorted Set 操作封装✅ 必需若不用 Redis此模块不可移除iot-ucy-mqtt-adapter可选MQTT 协议适配层兼容 EMQX/Mosquitto⚠️ 按需启用iot-ucy-http-gateway可选HTTP 接口网关仅提供/device/status等只读查询非主通道⚠️ 非必需提示不要 clone 整个仓库后盲目mvn install。先检查pom.xml中parent是否指向spring-boot-starter-parent:2.7.18且netty-all版本为4.1.96.Final低于 4.1.90 有粘包处理缺陷高于 4.1.100 与 Redis Lettuce 6.x 兼容性问题已知。若本地 Maven 仓库无对应版本手动下载 netty-all-4.1.96.Final.jar 放入~/.m2/repository/io/netty/netty-all/4.1.96.Final/目录下再执行构建。2.2 配置 application.yml聚焦 Netty Redis 两大核心参数以下是最小可运行配置删除所有注释保留缩进server: port: 0 # 关键禁用 Tomcat设为 0 表示不启动 Web 容器 spring: profiles: active: prod redis: host: 127.0.0.1 port: 6379 password: database: 0 lettuce: pool: max-active: 50 max-idle: 20 min-idle: 5 max-wait: 3000ms main: allow-bean-definition-overriding: true iot: ucy: netty: boss-thread-count: 1 # Netty Boss 线程数1 足够负责 accept worker-thread-count: 16 # Worker 线程数建议 CPU 核数 × 2处理 read/write tcp: port: 8081 # 设备直连 TCP 端口非 HTTP idle-timeout: 300 # 心跳超时秒数单位秒 max-connections: 100000 # 单节点最大连接数需配合 ulimit -n 调整 mqtt: enabled: false # 默认关闭 MQTT避免端口冲突 device: session-expire: 7200 # 设备会话过期时间秒对应 Redis Sorted Set score shadow-ttl: 86400 # 设备影子状态 TTL秒对应 Redis Hash 过期 command: retry-times: 3 # 指令重试次数失败后自动进 Redis Stream 重试队列 timeout: 5000 # 指令响应超时毫秒数参数说明netty.tcp.port是设备连接的真实入口不是 Spring Boot 的server.portredis.lettuce.pool.max-active必须 ≥netty.worker-thread-count否则 Redis 连接池将成为瓶颈iot.ucy.device.session-expire必须 ≤ Redis 的maxmemory-policy设置建议用allkeys-lru避免影子状态被误删。2.3 启动服务并验证 TCP 连接通路编译完成后进入iot-ucy-core模块目录执行java -jar target/iot-ucy-core-1.4.2.jar --spring.profiles.activeprod启动成功标志日志中出现[INFO] NettyServer started on port 8081, bossGroup threads: 1, workerGroup threads: 16 [INFO] Redis connection established: redis://127.0.0.1:6379 [INFO] DeviceSessionManager initialized, max sessions: 100000 [INFO] CommandDispatcher registered with Redis Stream: iot:command:stream此时用telnet 127.0.0.1 8081测试基础连通性。若返回Connected to 127.0.0.1即表示 Netty TCP 服务已就绪。注意此时无任何协议握手仅验证 TCP 层可达。下一步需用 iot-ucy 提供的DeviceSimulator工具位于tools/目录模拟真实设备注册cd tools java -jar device-simulator-1.4.2.jar --host 127.0.0.1 --port 8081 --device-id DEV_001 --protocol binary成功后观察日志出现[INFO] Device DEV_001 connected via binary protocol, session created [INFO] Redis: HSET iot:shadow:DEV_001 status online, last_seen 1717023456 [INFO] Redis: ZADD iot:sessions:online 1717023456 DEV_001这证明设备已注册进 Redis 影子状态Hash和在线会话Sorted SetNetty 连接、Spring Boot IOC、Redis 三者已形成闭环。3. 设备接入协议定制如何安全扩展私有二进制协议并规避 Netty 粘包/半包问题iot-ucy 默认支持两种协议binary紧凑型私有二进制和json调试用明文。生产环境强烈建议使用binary但必须自行定义帧结构。其核心不在“加个 Decoder”而在于协议设计阶段就规避粘包根源——iot-ucy 的BinaryProtocolDecoder采用“长度域 校验和”双保险机制而非简单换行符或固定长度。3.1 二进制帧格式规范必须严格遵循每个设备报文由 4 部分组成总长度 ≥ 12 字节字段长度字节说明示例十六进制Magic Number2固定0x55AA55 AAPayload Length2后续 payload 总长度不含 magic 和 length 自身00 1824 字节PayloadN实际业务数据含设备 ID、指令类型、参数等...CRC162payload 的 CRC16-CCITT 校验值初始值 0xFFFFA1 B2注意Payload Length字段本身不参与 CRC 计算校验范围仅为Payload区域。若设备端计算 CRC 错误iot-ucy 会在BinaryProtocolDecoder中直接丢弃该帧并记录CRC mismatch日志不触发后续业务逻辑。3.2 自定义 Decoder 实现继承 BaseBinaryDecoder新建类CustomDeviceDecoder extends BaseBinaryDecoder重写decodePayload()方法Override protected void decodePayload(ByteBuf in, DeviceSession session, ListObject out) { // 1. 读取设备ID固定8字节ASCII字符串 byte[] deviceIdBytes new byte[8]; in.readBytes(deviceIdBytes); String deviceId new String(deviceIdBytes).trim(); // 2. 读取指令类型1字节 byte cmdType in.readByte(); // 3. 读取参数长度2字节 int paramLen in.readShort(); // 4. 读取参数paramLen 字节 byte[] params new byte[paramLen]; in.readBytes(params); // 5. 构建业务对象此处仅为示意实际应映射到具体 POJO CustomDeviceMessage msg new CustomDeviceMessage(); msg.setDeviceId(deviceId); msg.setCmdType(cmdType); msg.setParams(params); // 6. 关键将消息放入 pipeline 下一环如 DeviceMessageHandler out.add(msg); }逻辑说明BaseBinaryDecoder已完成 Magic、Length、CRC 校验decodePayload()只需专注业务字段解析in.readBytes()必须严格按协议长度读取不可用readString()等变长方法out.add(msg)后消息将交由DeviceMessageHandler处理该 Handler 会自动更新 Redis 影子状态HSET iot:shadow:{deviceId} ...。3.3 粘包/半包处理的三个血泪经验Netty 粘包是 IoT 场景高频翻车点iot-ucy 的LengthFieldBasedFrameDecoder配置极易出错现象设备连续发送两条指令服务端收到一条 200 字节的混合帧原因LengthFieldBasedFrameDecoder的lengthFieldOffset设为0误以为 Magic 后就是长度域实际应为2Magic 占 2 字节解决在BinaryProtocolDecoder构造函数中显式设置super(65535, 2, 2, 0, 2); // maxFrameLength, lengthFieldOffset, lengthFieldLength, lengthAdjustment, initialBytesToStrip现象设备偶尔失联日志出现io.netty.handler.timeout.ReadTimeoutException原因ReadTimeoutHandler超时时间默认 30 秒小于设备心跳间隔如 60 秒解决在NettyServerConfig中调整pipeline.addLast(new ReadTimeoutHandler(90)); // 设为心跳间隔的 1.5 倍现象同一设备多次重连Redis 中残留多个iot:sessions:online成员原因设备断连时未触发channelInactive()因 Netty ChannelFuture 未正确监听解决在DeviceChannelHandler中确保Override public void channelInactive(ChannelHandlerContext ctx) throws Exception { String deviceId getDeviceId(ctx.channel()); redisTemplate.opsForZSet().remove(iot:sessions:online, deviceId); // 主动清理 super.channelInactive(ctx); }4. Redis 深度集成用 Streams 做指令队列、Hash 做设备影子、Sorted Set 做会话管理iot-ucy 不把 Redis 当缓存用而是当作分布式状态协调中枢。它放弃 Redis Pub/Sub可靠性差改用 Streams 实现指令的有序、可重试、可追溯分发放弃 String 存状态改用 Hash 结构支持字段级原子更新放弃 List 做在线列表改用 Sorted Set 实现按心跳时间排序的在线设备索引。这种设计让 Redis 从“辅助角色”变成“状态事实源”。4.1 指令下发Redis Streams 的消费组模型落地当后台调用CommandService.sendCommand(deviceId, command)时iot-ucy 执行将指令序列化为 JSON 字符串通过XADD iot:command:stream * deviceId DEV_001 type UPDATE params {temp:25} retry 3写入 Streams同时向iot:command:pending:DEV_001List 推入指令 ID用于去重消费者CommandConsumer使用消费组iot-command-group从 Streams 拉取// 初始化消费组仅首次执行 redisTemplate.execute((RedisCallbackObject) connection - { connection.executeCommand(XGROUP, CREATE, iot:command:stream, iot-command-group, $, MKSTREAM); return null; }); // 拉取指令阻塞 100ms ListMapObject, Object messages redisTemplate.opsForStream() .read(Consumer.from(iot-command-group, worker-1), StreamReadOptions.empty().count(1).block(Duration.ofMillis(100)), StreamOffset.create(iot:command:stream, ReadOffset.from()));关键参数ReadOffset.from()表示只读新消息count(1)防止单次拉取过多导致处理延迟消费成功后必须调用XACK失败则XCLAIM转给其他 worker。iot-ucy 默认配置 3 个 consumer worker保证单指令最多重试 3 次由iot.ucy.command.retry-times控制。4.2 设备影子状态Hash 结构的字段级原子更新设备影子Device Shadow存储在iot:shadow:{deviceId}Hash 中字段包括字段名类型说明更新方式statusStringonline/offlineHSET原子更新last_seenLong上次心跳时间戳HINCRBY原子递增避免并发覆盖firmware_versionString固件版本HSETbattery_levelInteger电量百分比HINCRBY支持负值为什么不用 String 存整个 JSON因为HINCRBY可对battery_level做原子减法设备每上报一次电量-1%而 JSON 字符串需先 GET 再解析再 SET存在并发覆盖风险。iot-ucy 的ShadowService.updateField()内部即调用redisTemplate.opsForHash().increment(key, field, delta)。4.3 在线会话管理Sorted Set 的心跳时间排序所有在线设备 ID 存于iot:sessions:onlineSorted Setscore 为最后一次心跳时间戳秒级。常用操作操作Redis 命令iot-ucy 封装方法用途查询在线设备数ZCARD iot:sessions:onlineSessionManager.getOnlineCount()实时监控面板获取最近 100 个上线设备ZREVRANGE iot:sessions:online 0 99 WITHSCORESSessionManager.getRecentOnlineDevices(100)运维排查清理超时设备ZREMRANGEBYSCORE iot:sessions:online 0 (1717020000)SessionManager.expireStaleSessions()定时任务每 5 分钟执行注意ZREMRANGEBYSCORE的(1717020000表示“小于 1717020000”括号为开区间。iot-ucy 的expireStaleSessions()会计算System.currentTimeMillis()/1000 - iot.ucy.device.session-expire作为阈值精准剔除掉线超时设备避免 Sorted Set 无限膨胀。5. 避坑指南生产环境踩过的 5 个致命坑及根治方案iot-ucy 在高并发设备场景下暴露的坑往往不在代码逻辑而在基础设施协同。以下是我在三个不同规模项目5k/50k/500k 设备中反复验证的 5 个必踩坑每条都附带线上已验证的修复命令。5.1 Redis 连接池耗尽Netty Worker 线程卡死在LettuceConnectionProvider.getConnection()现象设备连接数 5000 后新设备无法接入日志频繁出现Unable to acquire connection from pooltop显示 Java 进程 CPU 100%但 Netty 线程无日志输出。原因lettuce-core默认连接池max-active8而 iot-ucy 的worker-thread-count16每个 Worker 在处理指令时需获取 Redis 连接连接数不足导致线程阻塞等待。解决# 修改 application.yml 中 redis.lettuce.pool.max-active 至 ≥ worker-thread-count × 2 redis: lettuce: pool: max-active: 32 # 16 worker × 2 安全系数 max-wait: 1000ms # 降低等待上限快速失败而非卡死补充同时在NettyServerConfig中添加连接池健康检查lettuceClientResources ClientResources.builder() .dnsResolver(new DirContextDnsResolver()) // 防 DNS 缓存失效 .build();5.2 设备 ID 冲突两个不同设备使用相同 ID 导致影子状态错乱现象设备 A 和 B 都用DEV_001连接A 上报温度 25℃B 上报 30℃Redis 中iot:shadow:DEV_001的temperature字段交替变化前端看到温度跳变。原因iot-ucy 默认允许重复 ID 连接旧会话未强制踢出。解决启用强制踢出策略在application.yml中iot: ucy: device: duplicate-id-policy: KICK_OLD # 可选值ALLOW / REJECT / KICK_OLD启用后当DEV_001新连接建立时DeviceSessionManager会查找旧会话session redisTemplate.opsForHash().get(iot:sessions:map, DEV_001)向旧 Channel 发送DISCONNECT指令触发channelInactive删除旧会话ZREM iot:sessions:online {oldSessionId}建立新会话5.3 MQTT 适配器内存泄漏EMQX 断连后 Netty Channel 未释放现象启用iot-ucy-mqtt-adapter后运行 72 小时jstat -gc pid显示OU老年代持续增长jmap -histo pid | grep Channel显示io.netty.channel.DefaultChannelPipeline实例数达 2w。原因MQTT Broker 断连时MqttChannelHandler未正确触发channelInactive()导致 Channel 对象无法 GC。解决重写MqttChannelHandler的断连逻辑Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { if (cause instanceof IOException || cause.getCause() instanceof IOException) { log.warn(MQTT channel exception, force close: {}, ctx.channel().id(), cause); ctx.close(); // 强制关闭触发资源释放 } }同时在application.yml中增加 MQTT 心跳保活iot: ucy: mqtt: keep-alive: 60 # 秒 clean-session: true5.4 Redis Stream 消费积压指令下发延迟从 200ms 涨至 15s现象后台批量下发 1000 条指令XLEN iot:command:stream达 5wXPENDING返回 3w pending 消息XINFO GROUPS显示pel-count持续增长。原因CommandConsumer处理逻辑中存在数据库慢查询如 MyBatis 同步写日志导致单条指令处理 1s消费组积压。解决将耗时操作异步化// 原同步写日志 commandLogMapper.insert(log); // ❌ // 改为异步 CompletableFuture.runAsync(() - commandLogMapper.insert(log), logExecutor); // ✅增加消费组 worker 数量# 创建新 worker redis-cli --raw XGROUP CREATECONSUMER iot:command:stream iot-command-group worker-25.5 Netty 内存泄漏Direct Buffer OOM 导致服务崩溃现象运行 5 天后JVM 报java.lang.OutOfMemoryError: Direct buffer memoryjcmd pid VM.native_memory summary显示direct内存占用 2G。原因PooledByteBufAllocator未配置最大内存Netty 默认使用maxDirectMemory Runtime.getRuntime().maxMemory() * 2在容器环境下易超限。解决# 启动时强制指定 Direct Memory 上限 java -XX:MaxDirectMemorySize512m -jar iot-ucy-core-1.4.2.jar并在NettyServerConfig中显式配置 allocatorEventLoopGroup workerGroup new NioEventLoopGroup( 16, new DefaultThreadFactory(netty-worker), // 关键限制 direct buffer 总量 new PooledByteBufAllocator(true, 1, 1, 8192, 11, 0, 0, PlatformDependent.directMemoryCapacity() / 4) );6. 进阶技巧用 Redis Lua 脚本实现设备级分布式锁彻底解决并发指令冲突在工业控制场景中“同一设备同一时刻只能执行一条指令”是硬性要求。比如温控器正在执行“升温至 28℃”此时若又收到“降温至 22℃”必须拒绝后者而非覆盖。iot-ucy 的DeviceLockService就是为此而生——它不依赖 Redisson 或自研锁框架而是用一行 Lua 脚本实现原子性加锁、自动过期、可重入判断实测 QPS 12w 无锁竞争。6.1 设备锁 Lua 脚本原理与部署脚本device-lock.lua内容如下保存为文件后续加载-- KEYS[1] lock key (e.g., iot:lock:DEV_001) -- ARGV[1] request id (unique per client) -- ARGV[2] expire seconds (e.g., 30) local current redis.call(GET, KEYS[1]) if current false then -- 锁不存在直接设置 redis.call(SET, KEYS[1], ARGV[1], EX, ARGV[2]) return 1 elseif current ARGV[1] then -- 已持有锁延长过期时间可重入 redis.call(EXPIRE, KEYS[1], ARGV[2]) return 1 else -- 锁被他人持有 return 0 end为什么不用SET key value EX seconds NX因为 NX 无法判断是否自己持有锁导致可重入失败。此脚本通过GET判断当前值是否等于本请求 ID实现真正的可重入。6.2 在指令处理器中集成设备锁CommandProcessor的process()方法开头加入String lockKey iot:lock: deviceId; String requestId UUID.randomUUID().toString().replace(-, ); Long result redisTemplate.execute( deviceLockScript, // 上述 Lua 脚本 Collections.singletonList(lockKey), requestId, String.valueOf(30) // 锁过期 30 秒 ); if (result 0L) { log.warn(Device {} locked by another request, reject command, deviceId); throw new DeviceLockedException(Device deviceId is locked); } try { // 执行实际指令逻辑如 Modbus 写寄存器 executeModbusWrite(deviceId, command); } finally { // 释放锁仅当本请求持有锁时才删除 String unlockScript if redis.call(GET, KEYS[1]) ARGV[1] then return redis.call(DEL, KEYS[1]) else return 0 end; redisTemplate.execute( RedisScript.of(unlockScript, Long.class), Collections.singletonList(lockKey), requestId ); }关键细节requestId必须全局唯一UUID且在整个指令生命周期内保持不变finally中的解锁脚本必须用GET判断后再DEL防止误删他人锁锁过期时间30 秒应略大于指令执行最大耗时如 Modbus 超时设为 25 秒。6.3 生产环境锁性能压测数据i7-11800H Redis 7.2并发线程数指令总数加锁成功率平均加锁耗时锁冲突率10010000100%0.18ms0.02%100010000099.998%0.23ms1.8%500050000099.992%0.31ms8.7%数据说明锁冲突率 被拒绝指令数 / 总指令数× 100%当冲突率 5%说明设备指令频次过高需前端限流或合并指令平均耗时稳定在 0.3ms 内证明 Lua 脚本方案远优于网络往返的 Redisson 方案后者平均 1.2ms。我上线第一个 iot-ucy 项目时曾因没加设备锁导致某产线 12 台 PLC 同时收到“急停”和“复位”指令PLC 状态机混乱停机 47 分钟。后来我把device-lock.lua脚本刻在团队 Wiki 首页要求所有指令处理器必须前置调用。现在每次新同事问“为什么指令要加锁”我就打开 Redis CLI现场执行EVAL脚本看着1和0的返回值比讲半小时原理都管用。希望帮到你。本文还有配套的精品资源点击获取
返回列表