6.3 IoT DC3工程实践
6.3.1 IoT DC3项目架构概览:模块划分与核心组件
第 2 章给出了物联网平台的分层蓝图,本节用 IoT DC3 将它落到可编译、可部署的模块。理解这一架构时,最重要的是区分三条边界:北向请求怎样进入中心服务,Driver 怎样与 Manager 同步元数据,位号命令与数据怎样通过 RabbitMQ 异步流转。
模块划分:北向统一入口、四中心协作、南向协议适配
北向接入层由 dc3-gateway 提供统一入口。Gateway 基于 Spring Cloud Gateway,将 /api/v3/auth/**、/api/v3/manager/**、/api/v3/data/**、/api/v3/agentic/** 分别路由到对应中心,并在受保护路由上执行 Authentic 过滤器。服务目标使用固定服务名和环境变量,不依赖独立注册中心。
平台服务层包含四个当前实际存在的中心:
dc3-center-auth:认证、授权、租户与 OAuth/MCP(Model Context Protocol)管理。dc3-center-manager:Driver、设备、模板、位号及属性等元数据管理,并向 Driver 提供 gRPC 业务注册与查询接口。dc3-center-data:位号值接收、最新值与历史查询、位号命令和自定义命令提交、执行回执处理及告警数据能力。dc3-center-agentic:模型配置、会话管理和 Spring AI@Tool工具调用。
当前架构中不存在独立的“Command Service”。位号读写入口属于 Data,Data 把命令发布到 RabbitMQ,Driver 异步消费并回传结果。
南向协议层由多个独立 Driver 服务组成,例如 MQTT、Modbus TCP/RTU、OPC UA、S7、IEC 104 等。Driver SDK 用 DriverProtocol、DriverReadService、DriverWriteService、DriverCustomService 等能力接口隔离协议差异。Driver 启动时通过 DriverRegisterService 调用 Manager 的 gRPC driverRegister 完成业务注册;运行时通过 RabbitMQ 接收位号命令和自定义命令,并上报位号值、状态、事件与执行回执。
基础设施与通信边界
IoT DC3 把关系数据、时序数据和异步消息分别放在可替换边界后。默认开发栈使用 PostgreSQL/TimescaleDB 与 RabbitMQ,Caffeine 提供进程内热点缓存;DC3_DB_TYPE、DC3_TSDB_TYPE、DC3_MQ_TYPE 分别选择关系方言、时序适配器和消息适配器。RabbitMQ 的 Exchange、队列、TTL、死信和 ack/nack 是默认适配器细节,不应被写成所有 Broker 的共同机制。平台仍没有 Nacos 等独立注册中心;dc3-driver-kafka 是南向数据源驱动,不等于内部 Kafka 适配器。
同步与异步的分工如下:
- 外部客户端经 Gateway 同步访问 Auth、Manager、Data、Agentic。
- Driver 经 gRPC 同步调用 Manager,完成业务注册与元数据查询。
- Data 经 RabbitMQ 异步向目标 Driver 投递位号读写和自定义命令。
- Driver 经 RabbitMQ 异步向 Data 上报位号值、状态、事件和命令回执。
技术栈选型
以 2026-08-29 的 987c96d50 快照为准,主干使用 Java 21、Spring Boot 4.0.6、Spring Cloud 2025.1.1 和 Spring AI 2.0.0。北向使用 REST/HTTP,中心与 Driver 的管理契约使用 gRPC + Protobuf,设备侧通信由各协议 Driver 选择相应客户端。数据与消息层通过端口适配器隔离具体产品;版本号和适配器清单属于易变事实,升级时应重新核对构建文件与官方能力矩阵。
6.3.2 设备数据采集与协议适配层实现
采集层负责把异构现场报文转换为平台统一的位号值。它需要处理协议连接、编解码、设备与位号元数据、读写语义以及异常恢复,但不应把告警规则、历史查询等平台业务塞进 Driver。IoT DC3 通过独立 Driver 服务与 Driver SDK 把这条边界固定下来。
Driver SDK 的真实能力接口
IoT DC3 当前没有一个所有驱动共同实现的 DeviceDriver 抽象,也没有 SDK 统一提供的全局 ConnectionManager。协议能力由细粒度接口组合:
public interface DriverCustomService extends DriverLifecycle,
DriverMetadataListener, DriverHealth, DeviceHealth,
DriverProtocol, DriverCommand, DriverValidator {
}
public interface DriverProtocol {
ReadPointValue read(Map<String, AttributeBO> driverConfig,
Map<String, AttributeBO> pointConfig,
DeviceBO device, PointBO point);
Boolean write(Map<String, AttributeBO> driverConfig,
Map<String, AttributeBO> pointConfig,
DeviceBO device, PointBO point,
WritePointValue writePointValue);
}SDK 侧的 DriverReadService、DriverWriteService 先解析设备、位号和属性元数据,再委托 DriverProtocol 与真实设备通信。协议实现只负责本协议的连接、编解码与读写:MQTT Driver 管理订阅与发布,Modbus Driver 处理寄存器和字节序,OPC UA Driver 处理节点与会话。连接池、心跳和退避策略由各 Driver 按协议特点实现,不能假定存在一套全局固定重连参数。
元数据、位号值与缓存边界
Driver SDK 使用 Caffeine 缓存 Driver、设备、位号及属性等元数据,避免每次采集都跨服务查询。启动时 DriverRegisterService 通过 gRPC 向 Manager 做业务注册和元数据同步;这不是服务注册中心行为。
协议读取成功后,DriverSenderService.pointValueSender 将标准化位号值发布到消息端口。位号值经标准化后进入所选 Broker,缓存与持久化统一由 Data 侧承担。默认数据链路是:
- Driver 解析协议数据并生成
PointValue。 DriverSenderService发布到 RabbitMQ 的位号值交换机。- Data 中的
PointValueReceiver消费消息并显式 ack、reject 或 nack/requeue。 - 低于批处理阈值时直接保存;高于阈值时先进入
PointValueJob的进程内批量缓冲,再异步批量写入。 - Data 将最新值写入本地 Caffeine 热点缓存,同时经
TsdbStore持久化;缓存未命中时回查所选时序存储。
这里有两类容易混淆的 Caffeine:Driver 侧缓存的是元数据,Data 侧缓存的是最新位号值。项目已用本地缓存替代旧的 Redis Repository 层,当前 Compose 也没有 Redis 服务。
主动轮询与被动上报
MQTT、TCP 等驱动可以在回调中接收设备主动上报;Modbus RTU、串口等协议通常由 Driver 的调度任务主动轮询。无论数据来自订阅回调还是定时读取,最终都应进入同一 DriverSenderService → RabbitMQ → PointValueReceiver 链路。串口驱动的调度结构由各 Driver 按协议特点自行设计,具体实现以对应 Driver 源码为准。
采集层调优也应沿真实瓶颈进行:Driver 侧关注连接数量、轮询周期和协议超时;所选 Broker 关注路由、积压和确认;Data 侧关注消费速度、批量间隔、缓存命中与时序写入。把这些参数误写成一套“Driver 两级缓存方案”,会让排障对象和责任边界全部错位。
6.3.3 微服务间通信:从REST到异步消息
IoT DC3 同时使用 REST、gRPC 与 RabbitMQ,但三者不是随意混用。REST 负责北向接口,gRPC 负责需要即时返回的中心与 Driver 管理契约,RabbitMQ 负责位号命令、执行回执和上行数据。判断某条链路是否准确,关键不是看它叫“控制面”还是“数据面”,而是回到实际生产者、消费者和确认语义。
同步链路:Gateway 路由与 Driver 管理契约
外部请求先由 Gateway 路由到 Auth、Manager、Data 或 Agentic。Gateway 使用固定服务名和 CENTER_*_HOST、GATEWAY_ROUTE_*_URI 等环境变量定位中心服务。
Driver 启动后,DriverRegisterService 通过 gRPC 调用 Manager 的 driverRegister 完成业务注册;设备、位号和属性等需要即时返回的元数据也通过 gRPC Facade 查询。这些调用属于同步管理链路,但不表示位号命令会通过 REST 或 gRPC 一路同步执行到物理设备。
异步链路:位号命令、回执与位号值
位号读写入口位于 Data。Data 根据目标 Driver 服务名将命令发布到 RabbitMQ,Driver 的 PointCommandReceiver 消费后调用 DriverReadService 或 DriverWriteService,再由 DriverSenderService 发布执行结果。自定义命令由 CommandReceiver 走同类路径。
上行方向同样使用 RabbitMQ:Driver 将位号值、设备状态、Driver 状态、事件和告警发布到对应交换机,Data 或 Manager 的消费者按职责处理。因此设备命令的真实语义是“提交—异步执行—结果回执”,而不是“HTTP 请求阻塞直到设备执行完成”。
@RabbitHandler
@RabbitListener(queues = "#{pointCommandQueue.name}")
public void pointCommandReceive(
Channel channel, Message message, PointCommandDTO command) {
// 校验 expireAt 与 commandId,按设备串行执行 read/write,
// 发送结果回执后再 ack;失败时按条件 reject 或 nack/requeue。
}PointCommandReceiver 在执行前检查 expireAt,以 commandId 去重,并用设备级锁避免同一设备的协议操作交错。Driver 专属命令队列还配置 TTL 和死信交换机。这里的幂等依据是命令 DTO 校验与本地去重缓存,属于 Driver 进程内的轻量去重;如需跨实例的严格幂等,应在更上层另行设计统一机制。
RabbitMQ 是当前唯一消息中间件
当前消息端口提供 RabbitMQ、Kafka、RocketMQ、Pulsar、ActiveMQ 与 MQTT 5 等适配器,由 DC3_MQ_TYPE 选择且同一部署只激活一个。RabbitMQ 仍是默认值。选用其他 Broker 不是改一个名称就完成迁移:应按官方能力矩阵和契约测试核对延迟消息、死信、顺序、确认、重试与可观测性。仓库中的 dc3-driver-kafka 是南向数据源驱动,与内部 Kafka 消息适配器不是同一层。
IoT DC3 的通信取舍可以归纳为一句话:同步链路解决“马上拿到管理结果”,异步链路解决“可靠穿过设备网络和服务速率差异”。这条边界与当前源码和部署清单一致。
6.3.4 工程检查清单:编码规范、日志与监控
微服务架构的代码一旦拆分运行,原来单体应用里容易察觉的问题会变得难以追踪。一个空指针异常只在某台节点上冒出来,一条设备上线日志散落在不同容器中,这些分散的碎片让人很难拼出完整的系统状态。本节给出四层工程检查清单,覆盖代码规范、日志体系、健康检查和指标监控——这几项是微服务从“能跑”到“能运维”的分水岭。
检查清单总览
表6-4 从四个维度列出必须覆盖的实践项。每一条都有对应的可操作验证手段,不依赖直觉判断。
表6-4 物联网微服务工程检查清单
| 维度 | 检查项 | 验证方式 | 说明 |
|---|---|---|---|
| 代码规范 | 静态检查工具集成 | 构建阶段强制通过 | 如 SonarQube / Checkstyle / SpotBugs,配置文件纳入版本库 |
| 代码规范 | 统一异常处理 | Handler 类全覆盖 | 使用 @ControllerAdvice 或自定义拦截器,避免 try-catch 污染业务逻辑 |
| 日志体系 | 日志分级标准化 | 按 ERROR/WARN/INFO/DEBUG 输出 | 禁止直接 System.out,日志格式统一含时间戳、线程、traceId |
| 日志体系 | 链路追踪 ID 注入 | 每个请求携带 traceId | 使用 Micrometer Tracing 或 MDC 手动注入,设备事件日志同样带 traceId |
| 健康检查 | Actuator 自定义端点 | /actuator/health 返回业务状态 | 至少检查数据库连接、消息队列状态、驱动心跳 |
| 健康检查 | 启动/存活/就绪探针 | Kubernetes 就绪探针可配置 | /actuator/health/liveness 和 /actuator/health/readiness 分离 |
| 指标监控 | Prometheus 端点暴露 | 采集器能拉取 /actuator/prometheus | 注册 Micrometer 指标,设备采集数、消息处理耗时、位号读写计数等业务指标 |
| 指标监控 | Grafana 告警规则 | 告警阈值配置后测试触发 | 如“设备心跳超时 > 30 秒”触发告警,通过钉钉/邮件通知 |
每一项的实际配置可以参考 Spring Boot Actuator 的官方文档。Actuator 提供了数十个内置端点,其中 /health、/info、/metrics、/prometheus 对微服务运维最为关键。在物联网场景里,设备的心跳超时判定常常不是简单的节点存活检查,需要自定义健康端点来聚合设备级状态。
自定义健康端点示例
假设一个协议驱动组件需要上报它所连接的设备是否在线。默认的 /actuator/health 只检查 Spring 容器和数据库,无法体现“驱动与 PLC 的 TCP 连接是否正常”。以下代码展示如何用 Spring Boot Actuator 的 HealthIndicator 接口扩展业务健康检查:
@Component
public class DeviceDriverHealthIndicator implements HealthIndicator {
private final List<DeviceConnection> connections;
public DeviceDriverHealthIndicator(List<DeviceConnection> connections) {
this.connections = connections;
}
@Override
public Health health() {
long offlineCount = connections.stream().filter(c -> !c.isAlive()).count();
if (offlineCount == 0) {
return Health.up()
.withDetail("totalConnections", connections.size())
.withDetail("status", "all devices online")
.build();
}
return Health.down()
.withDetail("totalConnections", connections.size())
.withDetail("offlineCount", offlineCount)
.withDetail("status", offlineCount + " device(s) offline")
.build();
}
}这段代码将设备驱动的连接状态暴露为健康检查指标。当 offlineCount>0 时整体标记为 DOWN,Kubernetes 就绪探针立刻可以据此将流量切走。
指标可视化与告警流
指标数据需要聚合层才能发挥作用。推荐的做法是:
- 指标暴露:每个微服务在
application.yml中启用management.endpoints.web.exposure.include=health,info,metrics,prometheus。 - 数据采集:Prometheus 以 pull 模式定期拉取各节点的
/actuator/prometheus端点。 - 可视化:Grafana 对接 Prometheus 数据源,配置设备接入数、消息队列积压、API 响应百分位等仪表盘。
- 告警:设定阈值触发告警通知(如接入 Prometheus Alertmanager)。
这套链路的核心在于业务指标的选取。常见的物联网指标包括:设备注册成功率、消息发布 QPS、位号查询 P99 延迟、驱动连接断开频次。对这些指标设定基线值之后,才算真正拥有了对系统异常的“可观测性”。
工程检查清单中的关键判断
清单中有几条容易在项目初期被忽视:
- 日志 traceId 必须贯穿端到端:设备数据从驱动到消息队列再到数据服务,如果每一跳都切断 traceId,调试时只能翻三四个日志文件去拼时间戳。统一注入 traceId 的成本很低,收益极高。
- 自定义健康检查不要只是“UP/DOWN”:返回详细的状态键值对,让运维人员一眼看出“哪个设备离线”“哪个数据库连接池满了”。
- 告警规则要有分级:设备心跳超时可触发 WARNING 告警;核心位号数据连续缺失要触发 CRITICAL 告警,通知值班工程师。
以下是用分层架构图形式总结的监控体系设计,每种类型的指标对应不同的采集与存储路径。
三支柱可观测性:从设备命令到最终状态
AIoT 系统的可观测性不能只回答“进程是否活着”,而要能沿着一次业务动作从设备走到最终状态。建议围绕日志、指标、Trace 三支柱构建统一模型:
- Trace:为每次“API → Gateway → Data → Driver → 设备回执”生成同一个
traceId,可用 OpenTelemetry 语义约定描述 span 名称、属性和状态。 - 指标:设备接入率、消息接收/重复/乱序率、命令成功率、确认时延、告警数;每个指标必须明确分母、窗口和聚合方式,避免同名指标含义漂移。
- 日志:结构化输出,字段至少包含时间戳、traceId、spanId、租户、用户、设备、Tool、审批 ID、错误码。审批、命令回执、模型决策等安全事件独立标签,用于合规审计。
三者之间的绑定比工具本身更重要:Trace 与日志共享 ID,指标与告警共享标签,人工审批与设备回执可回连到原始请求。没有统一 ID,事后回放就只能靠人工拼日志。本章只固定这套通用骨架;当链路中出现模型与工具调用时,LLM/Tool 子 span、token 成本指标和模型/Prompt 版本标签如何纳入三支柱,展开见 7.4.3。
灰度发布与回滚
投产前的部署实践应把“灰度 + 独立回滚”作为默认能力,而不是发生事故后临时补救:
- 每次发布关联一个 manifest:镜像 digest、Compose/K3s 配置、依赖版本、配置项;
- 变更先经过影子流量或 shadow-writes(读真实请求,不产生外部副作用);
- 进入生产按租户/设备维度灰度,观察数据链路、命令回执和业务指标;
- 出现回归时按组件回退:镜像回退、配置回退、依赖回退相互独立;
- 回滚后仍保留 traces,便于复盘失败原因与漂移边界;
- 高风险 OTA、驱动升级和边缘节点变更须走独立审批与批次;不允许一次全量升级所有网关。
灰度和回滚都不是“流程仪式”,其价值是把“看起来更好”变成有证据的变更管理:谁批准、改了什么、观测到什么、下一步如何撤销。若发布单元中还包含模型、Prompt 这类非代码资产,版本登记与按组件回退的要求会更细一层,专门讨论见 7.4.3。
工程协作与多仓库版本对齐
微服务落地后的第一道工程问题往往不是技术,而是协作。IoT DC3 把中心服务与各协议 Driver 放在同一仓库内,模块即边界;当驱动由不同团队甚至不同组织维护时,Driver、Driver SDK 与部署清单常常拆成多个 Git 仓库独立发版。多仓库换来解耦的自由,代价是“线上跑的到底是哪份代码”变得难以回答,需要三条纪律来对齐:仓库按变更节奏切分,跨仓库流动的只有接口契约;镜像 tag 必须能反查到源码 commit,用语义化版本或 commit 短哈希作 tag,禁止只认 latest;Driver SDK 的接口演进保持向后兼容,主版本与平台契约对齐,各 Driver 在依赖清单中声明可用的 SDK 版本区间,平台升级前先核对兼容矩阵,再排驱动升级批次。