Skip to content

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 用 DriverProtocolDriverReadServiceDriverWriteServiceDriverCustomService 等能力接口隔离协议差异。Driver 启动时通过 DriverRegisterService 调用 Manager 的 gRPC driverRegister 完成业务注册;运行时通过 RabbitMQ 接收位号命令和自定义命令,并上报位号值、状态、事件与执行回执。

基础设施与通信边界

IoT DC3 把关系数据、时序数据和异步消息分别放在可替换边界后。默认开发栈使用 PostgreSQL/TimescaleDB 与 RabbitMQ,Caffeine 提供进程内热点缓存;DC3_DB_TYPEDC3_TSDB_TYPEDC3_MQ_TYPE 分别选择关系方言、时序适配器和消息适配器。RabbitMQ 的 Exchange、队列、TTL、死信和 ack/nack 是默认适配器细节,不应被写成所有 Broker 的共同机制。平台仍没有 Nacos 等独立注册中心;dc3-driver-kafka 是南向数据源驱动,不等于内部 Kafka 适配器。

同步与异步的分工如下:

  1. 外部客户端经 Gateway 同步访问 Auth、Manager、Data、Agentic。
  2. Driver 经 gRPC 同步调用 Manager,完成业务注册与元数据查询。
  3. Data 经 RabbitMQ 异步向目标 Driver 投递位号读写和自定义命令。
  4. Driver 经 RabbitMQ 异步向 Data 上报位号值、状态、事件和命令回执。
图6-6 IoT DC3 模块关系与数据流分层图同步管理与异步设备数据分流,业务数据按中心职责持久化。图6-6 IoT DC3 模块关系与数据流分层图同步管理与异步设备数据分流,业务数据按中心职责持久化。北向接入层Web / 第三方客户端REST / HTTPdc3-gateway固定服务名路由 · 认证平台服务层Auth认证 · 授权 · 租户OAuth / MCPManager设备与驱动元数据gRPC 注册 / 查询Data位号值 · 命令 · 告警命令入口 / 回执处理Agentic模型 · 会话 · ToolsSpring AI @ToolRabbitMQ · 默认消息适配器南向驱动层协议 DriverMQTT / Modbus / OPC UA协议适配与执行现场设备传感器 · 执行器 · PLC协议通信接入基础设施层Caffeine进程内热点缓存(本地)RabbitMQ:默认消息适配器(见中段)PostgreSQL / TimescaleDB业务数据 · 位号历史RESTREST 路由 · 认证gRPC 注册/查询发布命令回执/数据命令投递回执/状态协议通信各中心按职责持久化同步 REST 路由同步 gRPC 管理调用RabbitMQ 异步消息持久化图6-6 默认链路以 RabbitMQ 连接 Data 与 Driver;具体 Broker 由消息适配器选择。
图 6-6 IoT DC3 模块关系与数据流分层图

技术栈选型

以 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。协议能力由细粒度接口组合:

java
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 侧的 DriverReadServiceDriverWriteService 先解析设备、位号和属性元数据,再委托 DriverProtocol 与真实设备通信。协议实现只负责本协议的连接、编解码与读写:MQTT Driver 管理订阅与发布,Modbus Driver 处理寄存器和字节序,OPC UA Driver 处理节点与会话。连接池、心跳和退避策略由各 Driver 按协议特点实现,不能假定存在一套全局固定重连参数。

元数据、位号值与缓存边界

Driver SDK 使用 Caffeine 缓存 Driver、设备、位号及属性等元数据,避免每次采集都跨服务查询。启动时 DriverRegisterService 通过 gRPC 向 Manager 做业务注册和元数据同步;这不是服务注册中心行为。

协议读取成功后,DriverSenderService.pointValueSender 将标准化位号值发布到消息端口。位号值经标准化后进入所选 Broker,缓存与持久化统一由 Data 侧承担。默认数据链路是:

  1. Driver 解析协议数据并生成 PointValue
  2. DriverSenderService 发布到 RabbitMQ 的位号值交换机。
  3. Data 中的 PointValueReceiver 消费消息并显式 ack、reject 或 nack/requeue。
  4. 低于批处理阈值时直接保存;高于阈值时先进入 PointValueJob 的进程内批量缓冲,再异步批量写入。
  5. 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_*_HOSTGATEWAY_ROUTE_*_URI 等环境变量定位中心服务。

Driver 启动后,DriverRegisterService 通过 gRPC 调用 Manager 的 driverRegister 完成业务注册;设备、位号和属性等需要即时返回的元数据也通过 gRPC Facade 查询。这些调用属于同步管理链路,但不表示位号命令会通过 REST 或 gRPC 一路同步执行到物理设备。

异步链路:位号命令、回执与位号值

位号读写入口位于 Data。Data 根据目标 Driver 服务名将命令发布到 RabbitMQ,Driver 的 PointCommandReceiver 消费后调用 DriverReadServiceDriverWriteService,再由 DriverSenderService 发布执行结果。自定义命令由 CommandReceiver 走同类路径。

上行方向同样使用 RabbitMQ:Driver 将位号值、设备状态、Driver 状态、事件和告警发布到对应交换机,Data 或 Manager 的消费者按职责处理。因此设备命令的真实语义是“提交—异步执行—结果回执”,而不是“HTTP 请求阻塞直到设备执行完成”。

java
@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 消息适配器不是同一层。

图6-7 IoT DC3 服务间通信架构图同步链路用于管理查询,异步链路承载设备命令、数据、状态和回执。图6-7 IoT DC3 服务间通信架构图同步链路用于管理查询,异步链路承载设备命令、数据、状态和回执。同步管理链路双向异步设备链路北向接入层Web / 第三方客户端REST / HTTPSpring Cloud Gateway固定服务名路由 · 认证平台服务层Auth认证 · 授权 · 租户Manager设备与驱动元数据gRPC 注册 / 查询Data点位读写与命令入口Agentic模型 · 会话 · ToolsRabbitMQ · 唯一消息总线TTL · DLX · ack / nack南向 DriverDriver执行 · 去重 · 回执MQTT / Modbus / OPC UA现场设备传感器 · 执行器 · PLC协议通信接入基础设施层Caffeine进程内热点缓存RabbitMQ:消息总线(见中段)PostgreSQL业务数据 · 位号历史RESTREST 路由 · 认证gRPC 注册/查询命令数据 / 回执驱动队列上报队列协议通信各中心按职责持久化同步 REST 路由同步 gRPC 管理调用RabbitMQ 异步消息图6-7 Driver → Manager 使用同步 gRPC;Data ↔ RabbitMQ ↔ Driver 使用双向异步消息。
图 6-7 IoT DC3 服务间通信架构图

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 接口扩展业务健康检查:

java
@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 就绪探针立刻可以据此将流量切走。

指标可视化与告警流

指标数据需要聚合层才能发挥作用。推荐的做法是:

  1. 指标暴露:每个微服务在 application.yml 中启用 management.endpoints.web.exposure.include=health,info,metrics,prometheus
  2. 数据采集:Prometheus 以 pull 模式定期拉取各节点的 /actuator/prometheus 端点。
  3. 可视化:Grafana 对接 Prometheus 数据源,配置设备接入数、消息队列积压、API 响应百分位等仪表盘。
  4. 告警:设定阈值触发告警通知(如接入 Prometheus Alertmanager)。

这套链路的核心在于业务指标的选取。常见的物联网指标包括:设备注册成功率、消息发布 QPS、位号查询 P99 延迟、驱动连接断开频次。对这些指标设定基线值之后,才算真正拥有了对系统异常的“可观测性”。

工程检查清单中的关键判断

清单中有几条容易在项目初期被忽视:

  • 日志 traceId 必须贯穿端到端:设备数据从驱动到消息队列再到数据服务,如果每一跳都切断 traceId,调试时只能翻三四个日志文件去拼时间戳。统一注入 traceId 的成本很低,收益极高。
  • 自定义健康检查不要只是“UP/DOWN”:返回详细的状态键值对,让运维人员一眼看出“哪个设备离线”“哪个数据库连接池满了”。
  • 告警规则要有分级:设备心跳超时可触发 WARNING 告警;核心位号数据连续缺失要触发 CRITICAL 告警,通知值班工程师。

以下是用分层架构图形式总结的监控体系设计,每种类型的指标对应不同的采集与存储路径。

图6-8 微服务可观测性分层架构指标、日志与告警各走明确链路,Grafana 从两类存储查询展示。图6-8 微服务可观测性分层架构指标、日志与告警各走明确链路,Grafana 从两类存储查询展示。服务暴露层采集与存储层展示与告警层微服务 / Driver/actuator/prometheus暴露指标端点日志服务结构化日志 · traceId时间戳 · 租户 · 错误码Prometheus拉取指标 · 规则评估pull 模式日志采集器采集 · 解析Elasticsearch日志索引存储Alertmanager分组 · 路由 · 抑制接收 Prometheus 告警钉钉 / 邮件 / 值班外部通知通道Grafana指标与日志查询Prometheus pull(拉取指标)日志采集写入索引告警规则通知路由查询指标查询日志指标拉取(pull)日志链路告警与通知Grafana 查询(虚线)图6-8 Prometheus 主动拉取指标并驱动告警;日志经采集器入 ES,Grafana 分别查询 Prometheus 与 ES。
图 6-8 微服务可观测性分层架构

三支柱可观测性:从设备命令到最终状态

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 版本区间,平台升级前先核对兼容矩阵,再排驱动升级批次。

从工业软件到 AI 智能体 · 构建面向智能体演进的多协议、云原生、开源工业物联网平台