本文目录导读:

Java在物联网(IoT)平台开发中占据重要地位,因为其跨平台性(Write Once, Run Anywhere)、丰富的生态库和强大的并发处理能力,非常适合构建从设备接入到数据展示的完整链路。
这里为你梳理一个典型的Java物联网平台的架构案例,包含了核心模块、技术选型和业务场景。
案例背景:智慧农场环境监控系统
业务需求:实时监控大棚内的温度、湿度、土壤酸碱度、光照强度;远程控制风机、卷帘、灌溉阀门;当数据异常时(如温度过高)自动告警并联动设备。
整体架构分层(微服务风格)
Java物联网平台通常采用分层架构,这里我将以 Spring Cloud Alibaba 技术栈为例:
[设备终端/传感器]
↓ (MQTT/CoAP)
[接入层] - Netty + MQTT Broker (EMQX)
↓
[核心服务层] - Spring Boot 微服务集群
|--- 设备管理服务
|--- 数据采集服务 (Kafka 缓冲)
|--- 规则引擎服务 (告警/联动)
|--- 命令下发服务
↓
[数据层] - 时序数据库 + 关系型数据库
↓
[展示层] - Web管理后台 + 大屏可视化 (Vue)
核心模块实现与技术要点
设备接入层(通信层)
技术选型:Netty(TCP/UDP处理)+ EMQX(MQTT Broker)或 Moquette(嵌入式)。
- 案例实现:
- 协议解析:设备上报数据通常不是标准JSON,而是二进制或自定义报文,这里使用Netty编写Decode(解码器),将字节流解析为Java POJO对象。
- 连接管理:使用Netty的ChannelGroup管理设备会话,维护设备ID与Channel的映射关系。
- 心跳机制:通过IdleStateHandler检测设备心跳,判断设备在线/离线状态。
数据管道(消息队列)
技术选型:Apache Kafka 或 RabbitMQ。
- 案例实现:
- 设备数据量巨大(每秒千级/万级TPS),不能直接写入数据库。
- Java服务将解析后的数据发送到Kafka的
device-metricsTopic。 - 下游的数据存储服务和规则引擎服务消费该Topic,进行解耦和削峰填谷。
设备管理服务
技术选型:Spring Boot + Mybatis-Plus + MySQL + Redis。
- 案例实现:
- 物模型:定义“温度”、“湿度”这些属性的标识符、数据类型(Int/Float)、取值范围。
- 生命周期:管理设备注册、激活、禁用、删除,设备首次上线时,调用该服务的接口进行身份认证(一机一密)。
数据存储层(时序数据)
技术选型:TDengine 或 InfluxDB(时序数据库)+ Redis(缓存最新值)+ MySQL(存储设备元数据)。
- 案例实现:
- 设备上报的温度数据是时序数据,用MySQL存储会导致数据量过大查询极慢。
- Java服务使用JDBC或原生API将Kafka消费的数据批量写入TDengine。
- SQL示例(TDengine):
INSERT INTO farm.weather VALUES (now, 25.3, 80, 7.5)。 - 设备最新状态(如当前温度)存入Redis,用于Web端快速展示。
规则引擎(告警与联动)
技术选型:Drools 或 自研条件匹配。
- 案例实现:
- 告警:定义规则
温度 > 35℃,触发后生成告警记录,并通过邮件/短信通知农户。 - 联动:定义规则
土壤湿度 < 20%且 设备在线,则自动向命令下发服务发送指令,开启灌溉阀门。 - 去重:避免同一设备连续触发告警导致消息轰炸,通常设置静默时间窗口。
- 告警:定义规则
命令下发与反向控制
技术选型:Netty 写数据 + MQTT QoS1/QoS2。
- 案例实现:
- 用户在Web端点击“打开风机”。
- 服务端通过设备管理服务找到该设备的Channel。
- 构造控制指令(如
{"cmd": "relay_on", "id": 1}),通过 Netty 写入 Channel 发送给设备。
关键代码片段演示
设备数据接入(Netty 解码器示例)
// 设备协议: 帧头(0xAA) + 长度(2字节) + 设备类型(1字节) + 数据区 + 校验
public class DeviceDecoder extends ByteToMessageDecoder {
@Override
protected void decode(ChannelHandlerContext ctx, ByteBuf in, List<Object> out) {
// 1. 等待积累足够的数据头(比如5个字节)
if (in.readableBytes() < 5) return;
in.markReaderIndex();
// 校验帧头
if (in.readByte() != 0xAA) {
ctx.close(); // 非法数据,断开连接
return;
}
short length = in.readShort(); // 读取数据长度
// 2. 如果数据不足一帧,重置读指针并等待
if (in.readableBytes() < length) {
in.resetReaderIndex();
return;
}
// 3. 读取数据区
byte[] data = new byte[length];
in.readBytes(data);
// 转换为业务对象
DeviceMessage msg = parseToPojo(data);
out.add(msg); // 传递给下一个Handler(业务处理器)
}
}
命令下发(使用 MQTT 客户端)
@Component
public class CommandSender {
@Autowired
private MqttTemplate mqttTemplate; // 假设使用集成好的 MQTT 客户端
public void sendCommand(String deviceId, String command) {
// 设备订阅的主题通常为: /device/{deviceId}/command
String topic = "/device/" + deviceId + "/command";
// QoS 1 确保消息至少送达一次
mqttTemplate.send(topic, command.getBytes(StandardCharsets.UTF_8), MqttQoS.AT_LEAST_ONCE);
// 记录下发日志
recordCommandLog(deviceId, command);
}
}
告警订阅(Kafka 消费者)
@Component
public class RuleEngineConsumer {
@KafkaListener(topics = "device-metrics", groupId = "rule-engine")
public void onMessage(DeviceMetric metric) {
// 获取规则(可通过缓存)
Rule rule = ruleCache.get(metric.getDeviceType());
// 判断阈值
if (metric.getValue() > rule.getMaxValue()) {
// 触发告警短信通知
sendAlert(metric.getDeviceId(), rule.getThresholdDesc());
}
}
}
核心难点与解决思路
| 难点 | 解决方案 |
|---|---|
| 海量设备接入 | 使用Netty 的零拷贝和异步非阻塞特性;引入 负载均衡(多台EMQX或Netty实例)。 |
| 数据乱序/丢包 | 在Java层实现幂等性(唯一业务ID),对于重传的数据进行去重。 |
| 设备离线补偿 | 启动定时任务(如 xxl-job),定时扫描设备状态,对离线超时的设备下发唤醒指令或主动断开。 |
| 高并发写入 | 不能直接操作数据库,必须经过 Kafka 削峰,然后批量入库(JDBC Batch or TSDB 自动落盘)。 |
| 协议多样性 | 设计策略模式(HandlerMapping),将不同的设备协议(Modbus、自定义TCP、CoAP)抽象成不同的 Decoder 和 Handler。 |
技术栈清单
- 后端框架:Spring Boot 2.x / 3.x, Spring Cloud Gateway。
- 通信协议:MQTT (推荐), Netty (TCP)。
- 消息队列:Apache Kafka(近实时)。
- 时序数据库:TDengine(自带超级表,适合物联网)。
- 关系型数据库:MySQL(存储用户、设备信息)。
- 缓存:Redis。
- 容器化:Docker + K8s。
这个案例展示了Java在物联网平台中的典型应用方式:利用Java的强类型和面向对象特性管理复杂的业务逻辑,利用Netty处理超高并发的IO,利用Kafka进行流量削峰,利用时序数据库处理海量数据。
建议:如果你正在学习或开发此类平台,可以先从设备接入(模拟器) -> 数据入库 -> Web展示 这条线开始,逐步完善规则引擎和命令控制,这样能快速上手。