昂科拉(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
│ 设备/告警/统计/升级 │ → 阿里云短信/微信/邮件
└────────────────┘
连接管理模型:
ChannelConnectionCache.connections:ConcurrentHashMap<mac, Connection>,内存中维护"设备在线表",下行控制时 O(1) 找到设备的 ChannelConnection 同时把自身绑定到 Channel 的 AttributeKey 上(AttributeMapConstant.connection),让 handler 能从 channel 反查设备来自 pom.xml,按"真实作用"分类:
| 依赖 | 版本 | 作用 | 为什么选它 |
|---|---|---|---|
| spring-boot-starter-web | 1.5.6 | HTTP/JSON 接口、内嵌容器 | 快速搭建管理系统 API |
| spring-boot-starter-jetty | 1.5.6 | Web 容器 | 替代 Tomcat(Jetty 更轻,适合嵌入式) |
| netty-all | 4.1.25 | TCP 长连接服务器 | 物联网接入层事实标准,NIO 事件驱动,单线程管上万连接 |
| spring-boot-starter-data-redis | 1.5.6 | Redis 客户端 | 状态缓存、去重、key 过期事件 |
| mybatis-spring-boot-starter | 1.3.2 | ORM | SQL 可控、复杂报表 SQL 好写 |
| mysql-connector-java | runtime | MySQL 驱动 | 数据库连接 |
| 依赖 | 版本 | 状态 | 说明 |
|---|---|---|---|
| fastjson | 1.2.47 | ✅ 在用 | Redis 存设备状态(JSON 字符串)、HTTP 返回 JSON |
| protostuff-core / runtime + objenesis | 1.0.10 | ⚠️ 仅声明未使用 | 二进制序列化(可能是早期方案遗留,实际用自定义二进制协议) |
| 依赖 | 版本 | 作用 |
|---|---|---|
| aliyun-java-sdk-core / ecs | 4.x | 阿里云短信(AliSMS)、IP 定位 |
| spring-boot-starter-mail | - | 邮件发送(MailUtil) |
| fluent-hc (apache http) | 4.5.6 | HTTP 工具(HttpUtils,调用微信 API 等) |
| 依赖 | 版本 | 状态 | 说明 |
|---|---|---|---|
| shiro-spring + shiro-redis | 1.4.0 | ⚠️ 未使用 | 安全框架,本模块没有拦截器 |
| mina-core | 2.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。
| 帧头 | 类型 | 说明 |
|---|---|---|
0xFE | 数据上报/控制帧 | 设备上报状态、服务器下发控制 |
0xFA | 注册帧/模块信息帧 | 设备首次连接上报 ICCID |
0xF5 | 心跳帧 | 保活 |
常量见 constants/Constant.java:13。
0xFE)FE | 帧长度(1B) | MAC(15B ASCII) | 数据域(nB) | 校验和(1B) | AA
frame[1])frame[2..17))CRC16Util.computeChecksum,CRC16Util.java:62)0xAA0xFA,40 字节)FA | 00 | MAC(15B) | CMD(1B=02) | ICCID(20B) | checksum(1B) | AA
服务端解析(MyServerHandler.java:316):取 MAC、CMD、ICCID 后,把设备登记上线并回发"全查 + 时间校准"。
0xF5,20 字节)F5 | 00 | MAC(15B) | 00 | checksum(1B) | AA
服务端收到心跳后回一个响应帧(MyServerHandler.java:331)。
数据域是从起始 key 开始,key 连续递增,后面每个字节是 value:
[start_key][value_for_key][value_for_key+1][value_for_key+2] ...
0x00),后续每个字节依次是 key=0x00, 0x01, 0x02... 的 valueMyServerHandler.handlerFirst(MyServerHandler.java:608):从 data[0] 读起始 key,temp 每循环一次 +1,DecodeFrame.dataDecode 负责把 key/value 写入 DeviceStateDecodeFrame.dataDecode(DecodeFrame.java:232)通过clazz.getMethod("setData" + key, Integer.class) 把 value 设进 DeviceState 的 data00~data81 属性——这就是 DeviceState 为什么有几十个 data 字段(pojos/DeviceState.java)
关键 key 含义(对照 test/DeviceSimulator.java:222 的模拟数据域):
| key | 含义 |
|---|---|
| 0x00 | 系统故障码 |
| 0x04 | 燃烧状态/散热模式(位标志) |
| 0x05 | 开关机/ECO/定时/冬夏模式(位标志) |
| 0x06/0x07 | 累积用水量高低字节 |
| 0x08 | 瞬时流量 |
| 0x0b | 卫浴水温目标 |
| 0x0c | 供暖水温目标 |
| 0x12 | 当前采暖水温 |
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.stateCommand,DeviceController.java:129):
dataCode = ("k">byte) (dataCode + ("n">1 << position)); // 置1
dataCode = ("k">byte) (dataCode - ("n">1 << position)); // 清0
FE | len(0x14) | MAC(15B) | 功能码(1B) | checksum | AA,功能码在 byte[17](EncodeFrame4GModel.getResponseFrame,EncodeFrame4GModel.java:99)FE | len | MAC(15B) | 数据域(TLV) | checksum | AA,数据域从 byte[17] 开始(getConFrame)data=3C + 时 + 分 + 秒 组帧下发(adjustTime)CRC16Util.computeChecksum)CRC16Util.getCRC)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)」+「时间校准」+「当天预约时段」
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
userEventTriggered / channelInactive)IdleStateHandler(180,180,180) 注册读/写/全空闲事件MyServerHandler.java:166),防误判updateStateByMac(mac,0))、删除该设备的 Redis 缓存、移除 ConnectionCache 并关闭 channelDeviceController.control → DeviceServiceImpl.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」实现的同步确认,属于老代码典型写法,可优化(见第十一节)。
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)。
DeviceServiceImpl)后台触发升级(/device/upgradeBoard 或 autoUp/enforcement)
→ 校验设备在线、空闲(data12==0)、是否已在升级
→ 从磁盘读固件文件(Constant.FIRMWARE_PATH + location)
→ FileUtil.getStruUpdateProgInfo 生成升级信息(文件长度/CRC/总包数)
→ UpgradeEncodeFrame.encode 组「待升级程序信息」帧下发
→ Redis 记 boardAddress / struUpdateProgBoardInfo(6 秒过期 = 超时判失败)
→ 设备回 ACK 后,逐包下发分片(CmdWriteEncodeFrame)
→ 每包设备回执,服务端记录进度(uwPackIndex → 百分比)
enforcement:遍历目标设备,分批推送autoUp:内存队列 upQueue 逐台处理 + Redis 计数分页/device/upgradeProcess:从 Redis 读 StruUpdateProgInfo 计算百分比waterUsePeriod):瞬时流量>0 记录开始时间;流量变 0 计算时长与用量 → 写 t_water_using_duration,并累加到 t_device_databurningRecord):记录燃烧开始/结束的累积水/气量差值TimeServiceImpl.hourCount,每 55 分):遍历在线设备,用当前累积量 - 上次累积量 → t_data_hourdailyCount,凌晨 1 点)→ t_data_dailymonthCount,每月 1 日 3 点)→ t_data_monthstatistics 用反射调用 DeviceData.setClock0~23 把用量填进对应小时列(DeviceServiceImpl.java:1260)。
TimeServiceImpl.setDeviceReservation)每小时 15 分 + 设备上线时:
t_device_reservation,按星期)50 + 各时段(start_h,start_m,end_h,end_m,temp) 下发| 角色 | 示例 key | 说明 |
|---|---|---|
| 状态仓库 | mac+devSta(JSON)、mac+devStaHex | 最新设备状态,页面/控制校验直接读 |
| 去重/防抖 | mac:seriousFaultFaultWarnId、mac+valve | 同一故障/同一状态变化只处理一次 |
| 定时器 | mac:openWebSetTime、mac:seriousFault:* | 利用 key 过期事件 当闹钟 |
RedisKeyExpirationListener(config/RedisKeyExpirationListener.java)继承 KeyExpirationEventMessageListener,订阅 Redis 的 __keyevent@0__:expired:
mac:openWebSetTime 过期 → 页面监控超时清理mac:adjustCity:ip 过期 → 重新定位设备城市mac:seriousFault:* 过期 → 触发故障微信推送weixin_access_token 过期 → 自动刷新 access_token前提:redis.conf 需开启
notify-keyspace-events Ex,否则收不到事件。
统一格式:<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 专用记录 |
设备是海量 TCP 长连接、低频大连接数。Netty 的 NIO 事件循环(一个线程管成千上万个 channel)与零拷贝,是物联网接入层的标配。
MCU 内存和带宽有限,定长二进制帧省流量、解析快;TLV 数据域让"加字段不改协议",只改设备固件和服务端解析即可。
设备几十秒上报一次,全写 MySQL 扛不住。实时状态放 Redis,历史数据低频落库;还顺带利用过期事件做定时器。
下行控制需要 O(1) 找到设备 channel,内存比查库快几个数量级;单机部署够用。
Netty 的 handler 是 new 出来的,不在 Spring 容器里。项目用静态引用绕开 IoC 注入 Mapper/Service(MyServerHandler.java:89)。这是"能用但别扭"的历史写法。
项目经历过 MINA → Netty 迁移,IoSession import 是漏删的残留;Protostuff/Shiro 同理是早期方案的遗留。
NettyServer.java — ServerBootstrap / boss-child 线程组MyChannelInitializer.java — pipeline 装配顺序ServerDecodeHandler.java — ByteToMessageDecoder 拆包ServerEncodeHandler.java — MessageToByteEncoderMyServerHandler.java — 入站事件处理/心跳/异常文档/云平台5.昂科拉/通讯协议/显示板与云端通讯协议*.docConstant(帧头帧尾)→ HexUtils(字节↔hex)→ DecodeFrame(解析)→ EncodeFrame4GModel(组帧)→ CRC16Util(校验)RedisCache 封装;思考"为什么实时状态在 Redis、历史落 MySQL"test/DeviceSimulator.java 的服务器地址,模拟设备注册/上报/心跳/device/control,在 DeviceServiceImpl.control 断点看组帧 → 下发 → 轮询校验data00≠0 → faultHandler → Redis 过期 → 微信推送Constant.java:74-76、application.properties:15 硬编码了阿里云 AccessKey、微信模板、数据库密码→ 应移到环境变量/配置中心并加密
Thread.sleep + 轮询 Redis 做控制确认(DeviceServiceImpl.control)→ 改异步回调Thread.sleep(300) 串行下发 3 个帧 → 用顺序化 promisenew Thread() 裸线程(waterAdjust 等)→ 用线程池System.out.println 大量调试输出 → 统一走 logbacksaveHistory 里同步执行 → 异步化 + 分库分表MyApplicationRunner 把所有设备置离线(MyApplicationRunner.java:31)是兜底方案bytes[17]/bytes[22] 偏移量 → 抽象成注解式协议定义/编解码器DeviceState 130+ 个字段 getter/setter 手动维护 → 按协议自动生成RedisCache/HttpUtils 等工具散落 → 模块化Constant.java)| 常量 | 值 | 含义 |
|---|---|---|
OUTER_HEAD | 0xFE | 数据/控制帧头 |
BURN_OUTER_HEAD | 0xFA | 注册帧头 |
HEART_HEAD | 0xF5 | 心跳帧头 |
OUTER_END | 0xAA | 帧尾 |
INSIDE_HEAD | DF FD | 内帧帧头 |
INSIDE_END | 0xDE | 内帧帧尾 |
InConCode.java)| 常量 | 值 | 含义 |
|---|---|---|
READ_A | 0xEE | 下行全查 |
WRITE_A | 0xF1 | 下行写 A 类参数 / 上报 ACK |
UP_RESP_A | 0xF0 | 上行页面查询回应 |
READ_VERSION | 0xF0 | 读版本 |
WRITE_VERSION | 0xF1 | 写版本 |
DOWNLOAD_UPINFO | 0xF2 | 下发待升级信息 |
UPBOARDINFO_BURST | 0xF3 | 主板分片数据 |
UPEND | 0xF4 | 升级结束 |
UPFONTINFO_BURST | 0xF5 | 字库分片数据 |
REBOOT | 0xF9 | 重启 |
BURN_READ_SEQ 系列 | 0x01~0x04 | 序列号读写 |
| 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:struUpdateProgBoardInfo | OTA 文件地址与进度 |
weixin_access_token | 微信 access_token 缓存 |
文档/云平台5.昂科拉/通讯协议/src/test/java/com/dafeng/test/DeviceSimulator.java