◀ 索引angkela_netty 学习指南

angkela_netty 物联网设备接入服务 · 学习指南

昂科拉(Angkela)智能壁挂炉 / 净水设备物联网云平台的设备接入层

设备通过 WiFi/4G 模组以 TCP 长连接接入,实现状态上报、远程控制、用水统计、故障告警、OTA 升级与定时预约。


目录


一、项目是什么

一句话:一个"设备 ↔ 云平台"之间的 TCP 网关,负责跟成百上千台嵌入式设备保持长连接,翻译设备的二进制协议,并给后台管理系统提供 HTTP 接口。

三大职责:

方向内容
上行(设备→云)设备注册、心跳、状态实时上报、版本信息、升级应答、序列号烧录回执
下行(云→设备)远程控制指令、状态全查、时间校准、定时预约下发、OTA 升级分包下发
旁路故障检测 + 短信/微信告警、用水用气统计(小时/日/月/年)、数据落库

技术身份:Spring Boot 1.5.6 + Netty 4.1.25 单体服务,MySQL 持久化,Redis 做状态缓存与定时触发器。

启动入口:src/main/java/com/dafeng/wifi/SpringbootStartApplication.java:21

"k">public "k">static "k">void main("k">String[] args) {
    SpringApplication.run(SpringbootStartApplication."k">class, args);
    "k">new NettyServer().bing("n">8900);   // 端口 "n">8900
}

二、整体架构

┌──────────────────────────────────────────────────────────────┐
│                    后台管理系统 (angkela_backend)              │
│          HTTP 调用 ──>  DeviceController /device/*            │
└──────────────────────────────┬───────────────────────────────┘
                               │ 组装下行协议帧
┌──────────────────────────────▼───────────────────────────────┐
│                     Netty 服务端 (端口 8900)                   │
│   ServerBootstrap (boss/child 两个 NIO EventLoopGroup)         │
│   Pipeline:                                                   │
│   ┌────────────────────────────────────────────────────────┐  │
│   │ ServerDecodeHandler  拆包/组帧 → byte[]                 │  │
│   │ ServerEncodeHandler  byte[] → ByteBuf                  │  │
│   │ IdleStateHandler(180s)  读空闲心跳检测                  │  │
│   │ MyServerHandler       业务总入口                        │  │
│   └────────────────────────────────────────────────────────┘  │
└───────────────┬───────────────────────────────┬──────────────┘
                │                               │
        ┌───────▼────────┐            ┌─────────▼─────────┐
        │  协议层 utils/   │            │   Redis           │
        │  Encode/Decode  │            │  状态/去重/定时器   │
        └───────┬────────┘            └───────────────────┘
                │
        ┌───────▼────────┐
        │  业务层 Service │  → MyBatis → MySQL
        │  设备/告警/统计/升级 │  → 阿里云短信/微信/邮件
        └────────────────┘

连接管理模型


三、技术栈与依赖详解

来自 pom.xml,按"真实作用"分类:

1. 核心运行框架

依赖版本作用为什么选它
spring-boot-starter-web1.5.6HTTP/JSON 接口、内嵌容器快速搭建管理系统 API
spring-boot-starter-jetty1.5.6Web 容器替代 Tomcat(Jetty 更轻,适合嵌入式)
netty-all4.1.25TCP 长连接服务器物联网接入层事实标准,NIO 事件驱动,单线程管上万连接
spring-boot-starter-data-redis1.5.6Redis 客户端状态缓存、去重、key 过期事件
mybatis-spring-boot-starter1.3.2ORMSQL 可控、复杂报表 SQL 好写
mysql-connector-javaruntimeMySQL 驱动数据库连接

2. 序列化 / JSON

依赖版本状态说明
fastjson1.2.47✅ 在用Redis 存设备状态(JSON 字符串)、HTTP 返回 JSON
protostuff-core / runtime + objenesis1.0.10⚠️ 仅声明未使用二进制序列化(可能是早期方案遗留,实际用自定义二进制协议)

3. 外部服务集成

依赖版本作用
aliyun-java-sdk-core / ecs4.x阿里云短信(AliSMS)、IP 定位
spring-boot-starter-mail-邮件发送(MailUtil
fluent-hc (apache http)4.5.6HTTP 工具(HttpUtils,调用微信 API 等)

4. 其他(重点注意:遗留/未用)

依赖版本状态说明
shiro-spring + shiro-redis1.4.0⚠️ 未使用安全框架,本模块没有拦截器
mina-core2.0.16⚠️ 遗留只残留 IoSession import(ConnectionCache/DeviceServiceImpl/TimeServiceImpl/SessionCache),说明项目从 MINA 迁移到了 Netty
log4j / logback-日志(logback 实际生效,见 configs/logback-spring.xml
commons-lang3 / commons-pool2-⚠️ 少用通用工具 / Redis 连接池

阅读技巧:看到一个项目 pom 里"声明但没用"的依赖,往往是技术演进的历史痕迹。本项目就是 MINA → Netty 迁移过的典型案例。


四、目录结构总览

src/main/java/com/dafeng/wifi/
├── SpringbootStartApplication.java   启动类(拉起 Netty)
├── netty/                            Netty 服务端(接入层)
│   ├── NettyServer.java              ServerBootstrap 配置与启动
│   ├── MyChannelInitializer.java     pipeline 装配
│   ├── ServerDecodeHandler.java      拆包解码(粘包/半包处理)
│   ├── ServerEncodeHandler.java      编码写出
│   └── MyServerHandler.java          ★业务总入口(775 行核心)
├── client/                           客户端(联调/测试)
│   ├── NettyClient.java / MyClientChannelInitializer.java
│   ├── ClientDecodeHandler / ClientEncodeHandler
│   └── MyClientHandler.java
├── controllers/                      HTTP 接口(给后台)
│   ├── DeviceController.java         设备控制/查询/升级/烧录
│   └── LoadBalanceController.java    负载均衡(空壳)
├── services/ + services/impls/       业务层
│   ├── DeviceServiceImpl.java        ★设备控制/查询/OTA
│   ├── FaultWarnServiceImpl.java     故障告警处理
│   ├── SmsServiceImpl.java           短信/微信推送
│   ├── TimeServiceImpl.java          ★定时统计/预约下发
│   └── ...
├── mappers/                          MyBatis DAO(XML 在 resources/mappers)
├── pojos/                            数据库实体(Device/DeviceState/...)
├── models/                           协议结构体(InFrame/OutFrame/DeviceState/Stru*)
├── constants/                        协议常量、故障码、Redis key 名
├── config/                           Redis 工具/过期监听/启动初始化
└── utils/                            ★协议编解码、Hex、CRC、连接缓存等工具

五、通信协议详解(核心)

协议原始定义见项目根目录 文档/云平台5.昂科拉/通讯协议/显示板与云端通讯协议*.doc

5.1 三种帧类型(靠帧头区分)

帧头类型说明
0xFE数据上报/控制帧设备上报状态、服务器下发控制
0xFA注册帧/模块信息帧设备首次连接上报 ICCID
0xF5心跳帧保活

常量见 constants/Constant.java:13

5.2 上报帧格式(0xFE

FE | 帧长度(1B) | MAC(15B ASCII) | 数据域(nB) | 校验和(1B) | AA

5.3 注册帧格式(0xFA,40 字节)

FA | 00 | MAC(15B) | CMD(1B=02) | ICCID(20B) | checksum(1B) | AA

服务端解析(MyServerHandler.java:316):取 MAC、CMD、ICCID 后,把设备登记上线并回发"全查 + 时间校准"。

5.4 心跳帧格式(0xF5,20 字节)

F5 | 00 | MAC(15B) | 00 | checksum(1B) | AA

服务端收到心跳后回一个响应帧(MyServerHandler.java:331)。

5.5 数据域:TLV 风格(协议灵魂)

数据域是从起始 key 开始,key 连续递增,后面每个字节是 value

[start_key][value_for_key][value_for_key+1][value_for_key+2] ...

clazz.getMethod("setData" + key, Integer.class) 把 value 设进 DeviceStatedata00~data81 属性——这就是 DeviceState 为什么有几十个 data 字段(pojos/DeviceState.java

关键 key 含义(对照 test/DeviceSimulator.java:222 的模拟数据域):

key含义
0x00系统故障码
0x04燃烧状态/散热模式(位标志
0x05开关机/ECO/定时/冬夏模式(位标志
0x06/0x07累积用水量高低字节
0x08瞬时流量
0x0b卫浴水温目标
0x0c供暖水温目标
0x12当前采暖水温

5.6 位标志拆分

data04/data05 是位标志字节,一个字节承载多个开关状态:

// DecodeFrame.java:"n">246
"k">int data04_0 = unsignedValue & 0x1;        // bit0 燃烧状态
"k">int data04_1 = (unsignedValue >> "n">1) & 0x1; // bit1 散热模式(0散热器/1地板)
...
"k">int data05_1 = (unsignedValue >> "n">1) & 0x1; // bit1 ECO模式
"k">int data05_2 = (unsignedValue >> "n">2) & 0x1; // bit2 定时状态
"k">int data05_3 = (unsignedValue >> "n">3) & 0x1; // bit3 冬夏模式

控制下发时也通过位运算设置/清除某一位(DeviceController.stateCommandDeviceController.java:129):

dataCode = ("k">byte) (dataCode + ("n">1 << position));  // 置1
dataCode = ("k">byte) (dataCode - ("n">1 << position));  // 清0

5.7 下行帧(服务器→设备)

5.8 校验算法


六、核心业务流程详解

6.1 设备注册上线(MyServerHandler.channelRead / connNull

设备连上 8900
  → 发注册帧 FA(40B)  或  直接发上报帧 FE
  → 服务端 connNull():
     1. 从 channel 拿 IP,用 AddressUtil 做 IP 归属地定位(省/市)
     2. 按 MAC 查 DeviceModel(机型绑定表)
        - 已绑定:新增/更新 t_device 记录(省份/城市/IP/机型),标识机型已使用
        - 未绑定:写入临时表 + Redis 记 `mac:serialNumber` 时间戳(30s 后重试绑定)
     3. deviceService.online(): 建 Connection → 注册进 ConnectionCache → 设备状态置在线
     4. 下发「全查(READ_A=0xEE)」+「时间校准」+「当天预约时段」

6.2 数据上报处理(handlerFirst

收到 FE 上报帧
  → 从缓存取上次 devSta(JSON)和 devStaHex(原始 hex)
  → 若本次 hex 与上次相同,则跳过历史存储(去重,避免高频率重复落库)
  → 否则逐字节解析 TLV 数据域 → 填充 DeviceState
     · data00 != 0 → handlerForAll → 系统故障处理
     · data04 == 0x4 → handlerTwoForGourami → 燃烧用水/气统计
  → saveHistory(): 写 t_device_state_history_*(按表索引,超 5000万 行自动建新表)
  → 更新 t_device_state(不存在插入,存在更新)
  → 回发响应帧(页面开着回 UP_RESP_A=0xF0,否则 WRITE_A=0xF1)
  → 更新 Redis: mac+devSta(最新状态 JSON)、mac+devStaHex

6.3 心跳与离线(userEventTriggered / channelInactive

6.4 远程控制下发(DeviceController.controlDeviceServiceImpl.control

后台 HTTP: GET /device/control?strMac=xx&strData=0501...
  → 校验设备在线(ConnectionCache)
  → EncodeFrame4GModel.controllerEncodeFrame 组装下行控制帧
  → channel.writeAndFlush(bytes)
  → 本地模拟解析 data 期望的新状态(DecodeFrame.dataDecode)
  → 轮询 Redis mac+devSta(每 200ms,最多 5 次),比对期望值是否生效
  → 返回 resCode

⚠️ 这里是「sleep + 轮询 Redis」实现的同步确认,属于老代码典型写法,可优化(见第十一节)。

6.5 故障告警(FaultWarnServiceImpl + RedisKeyExpirationListener

流程(以严重故障为例):

上报 data00 != 0
  → faultHandler: 首次故障写 t_fault_warn + t_fault_notice
    → Redis 设 mac:seriousFaultFaultWarnId 与 5 秒过期 key
  → Redis key 过期事件触发 RedisKeyExpirationListener
    → 查故障分类缓存,写 fault_notice,调 smsService.faultSendWeiXin(微信公众号模板推送)

其他告警类型:软水用尽、漏水、长时间出水、滤芯剩余(4F/50/51)、盐桥/盐不足(163/16)。

6.6 OTA 固件升级(DeviceServiceImpl

后台触发升级(/device/upgradeBoard 或 autoUp/enforcement)
  → 校验设备在线、空闲(data12==0)、是否已在升级
  → 从磁盘读固件文件(Constant.FIRMWARE_PATH + location)
  → FileUtil.getStruUpdateProgInfo 生成升级信息(文件长度/CRC/总包数)
  → UpgradeEncodeFrame.encode 组「待升级程序信息」帧下发
  → Redis 记 boardAddress / struUpdateProgBoardInfo(6 秒过期 = 超时判失败)
  → 设备回 ACK 后,逐包下发分片(CmdWriteEncodeFrame)
  → 每包设备回执,服务端记录进度(uwPackIndex → 百分比)

6.7 用水/用气统计

statistics 用反射调用 DeviceData.setClock0~23 把用量填进对应小时列(DeviceServiceImpl.java:1260)。

6.8 定时预约下发(TimeServiceImpl.setDeviceReservation

每小时 15 分 + 设备上线时:

  1. 校准设备时间
  2. 查当天启用的预约时段(t_device_reservation,按星期)
  3. 组帧 50 + 各时段(start_h,start_m,end_h,end_m,temp) 下发
  4. 按当天是否有时段,设置 data05 的"定时图标"位(bit2)

七、Redis 的玩法

7.1 三种角色

角色示例 key说明
状态仓库mac+devSta(JSON)、mac+devStaHex最新设备状态,页面/控制校验直接读
去重/防抖mac:seriousFaultFaultWarnIdmac+valve同一故障/同一状态变化只处理一次
定时器mac:openWebSetTimemac:seriousFault:*利用 key 过期事件 当闹钟

7.2 Key 过期事件(进阶技巧)

RedisKeyExpirationListenerconfig/RedisKeyExpirationListener.java)继承 KeyExpirationEventMessageListener,订阅 Redis 的 __keyevent@0__:expired

前提:redis.conf 需开启 notify-keyspace-events Ex,否则收不到事件。

7.3 常用 key 约定

统一格式:<mac>:<业务名><mac><业务名>(历史混用)。


八、数据库模型

核心表(Mapper 名即表语义):

表/实体说明
t_device设备档案(MAC、机型、序列号、省市 IP、固件版本、升级状态)
t_device_model机型绑定表(MAC↔modelCode↔devCode)
t_device_state设备最新状态(data00~data81)
t_device_state_history_*状态历史(按表索引分表,超 5000 万行自动建新表)
t_device_data按小时用水量(clock0~23 列)
t_data_hour / t_data_daily / t_data_month时/日/月汇总
t_water_using_duration单次用水记录
t_fault_warn / t_fault_notice故障告警与通知记录
t_fault_code / t_customer_fault故障码字典 / 用户通知设置
t_device_reservation定时预约时段
t_device_burn_log序列号烧录日志
t_gourami_*燃热/RO 专用记录

九、设计动机分析

  1. 为什么 Netty 不用 Tomcat/BIO?

设备是海量 TCP 长连接、低频大连接数。Netty 的 NIO 事件循环(一个线程管成千上万个 channel)与零拷贝,是物联网接入层的标配。

  1. 为什么自定义二进制协议?

MCU 内存和带宽有限,定长二进制帧省流量、解析快;TLV 数据域让"加字段不改协议",只改设备固件和服务端解析即可。

  1. 为什么设备状态放 Redis?

设备几十秒上报一次,全写 MySQL 扛不住。实时状态放 Redis,历史数据低频落库;还顺带利用过期事件做定时器。

  1. 为什么在线表放内存 ConcurrentHashMap?

下行控制需要 O(1) 找到设备 channel,内存比查库快几个数量级;单机部署够用。

  1. 为什么 handler 用 @PostConstruct 存静态实例?

Netty 的 handler 是 new 出来的,不在 Spring 容器里。项目用静态引用绕开 IoC 注入 Mapper/Service(MyServerHandler.java:89)。这是"能用但别扭"的历史写法。

  1. 为什么有 MINA 依赖却用 Netty?

项目经历过 MINA → Netty 迁移,IoSession import 是漏删的残留;Protostuff/Shiro 同理是早期方案的遗留。


十、学习路线

阶段 1:Netty 基本功(必读,1 周)

  1. 读《Netty in Action》或官方用户指南:Reactor 模型、EventLoop、Pipeline、ByteBuf
  2. 对照本代码:
    • NettyServer.java — ServerBootstrap / boss-child 线程组
    • MyChannelInitializer.java — pipeline 装配顺序
    • ServerDecodeHandler.java — ByteToMessageDecoder 拆包
    • ServerEncodeHandler.java — MessageToByteEncoder
    • MyServerHandler.java — 入站事件处理/心跳/异常

阶段 2:把协议吃透(本项目的灵魂)

  1. 先读 文档/云平台5.昂科拉/通讯协议/显示板与云端通讯协议*.doc
  2. 再读代码:Constant(帧头帧尾)→ HexUtils(字节↔hex)→ DecodeFrame(解析)→ EncodeFrame4GModel(组帧)→ CRC16Util(校验)
  3. 重点:粘包/半包位标志拆分反射填充属性无符号数与大小端

阶段 3:Redis 在 IoT 里的套路

  1. 学习 Redis key 过期通知(keyspace notifications)
  2. 理解 RedisCache 封装;思考"为什么实时状态在 Redis、历史落 MySQL"

阶段 4:把业务闭环跑一遍

  1. test/DeviceSimulator.java 的服务器地址,模拟设备注册/上报/心跳
  2. 用后台调 /device/control,在 DeviceServiceImpl.control 断点看组帧 → 下发 → 轮询校验
  3. 追一次故障链路:data00≠0 → faultHandler → Redis 过期 → 微信推送

阶段 5:工程化进阶

  1. 学 Netty 官方拆包器(LengthFieldBasedFrameDecoder)替换手写解码
  2. 学异步控制(CompletableFuture / ChannelPromise)替换 sleep 轮询
  3. 学 Spring 事件/消息队列做故障通知,替换 Redis 过期事件
  4. 学连接级读写限流、黑白名单(当前代码未做安全防护)

十一、问题与改进建议

🔴 安全隐患(重要)

→ 应移到环境变量/配置中心并加密

🟠 性能与并发

🟡 健壮性

🟢 可维护性


十二、附录:关键常量速查

帧常量(Constant.java

常量含义
OUTER_HEAD0xFE数据/控制帧头
BURN_OUTER_HEAD0xFA注册帧头
HEART_HEAD0xF5心跳帧头
OUTER_END0xAA帧尾
INSIDE_HEADDF FD内帧帧头
INSIDE_END0xDE内帧帧尾

功能码(InConCode.java

常量含义
READ_A0xEE下行全查
WRITE_A0xF1下行写 A 类参数 / 上报 ACK
UP_RESP_A0xF0上行页面查询回应
READ_VERSION0xF0读版本
WRITE_VERSION0xF1写版本
DOWNLOAD_UPINFO0xF2下发待升级信息
UPBOARDINFO_BURST0xF3主板分片数据
UPEND0xF4升级结束
UPFONTINFO_BURST0xF5字库分片数据
REBOOT0xF9重启
BURN_READ_SEQ 系列0x01~0x04序列号读写

常用 Redis key

key含义
mac+devSta最新设备状态 JSON
mac+devStaHex最近一帧原始 hex(去重用)
mac+valve / mac+open / mac+openUsing用水开闭状态与累计
mac:seriousFaultFaultWarnId严重故障去重
mac:seriousFault:*故障推送定时 key
mac:serialNumber / mac:serialNumberType序列号绑定状态/类型
mac+boardAddress / mac:struUpdateProgBoardInfoOTA 文件地址与进度
weixin_access_token微信 access_token 缓存

参考资源