简介面向工业自动化数据采集场景的OPC数据采集服务项目完整覆盖从PLC/SCADA等设备读取实时数据通过MQTT与Kafka消息管道进行传输最终写入时序数据库InfluxDB的全链路。资源以Java源码为主分为opc-sink-ua-client、opc-sink-mqtt、opc-sink-kafka、opc-sink-influxdb及演示模块清晰展示了OPC UA客户端对接、消息发布订阅、Kafka中转和InfluxDB落库的实现方式适合需要构建实时数据管道或学习工业协议适配的开发者。压缩包共60个文件包括47个Java文件、6个XML配置、5个YAML配置、1个Git忽略文件与1个README说明整体大小仅61KB轻量便携。已有305人学习浏览代码按OPC采集、MQTT发布、Kafka分发、InfluxDB存储分层组织便于快速定位各环节实现。通过阅读源码与配置可掌握OPC UA客户端接入方法、MQTT主题管理、Kafka生产者消费者衔接及InfluxDB批量写入策略无论用于工业现场数据集成还是作为物联网时序数据处理的学习样例都具备很高参考价值可直接作为基础框架进行改造加速工业物联网项目落地。 车间里几十台设备的数据要统一收上来PLC 种类还杂西门子、三菱、罗克韦尔都有老一点的产线只有 OPC 接口数据得先进实时库给上位机看又要同步给后端的分析平台。这套东西听着不难真正落地的时候全是细节OPC 那边怎么连、MQTT 的主题怎么规划、中间为什么要塞一个 Kafka、最后怎么高效写进 InfluxDB。我最近把这条链路完整走了一遍把整个采集服务打包成了一个工程文件今天把这套方案的来龙去脉和实操过程完整拆开讲一遍希望能帮你少走点弯路。1. 链路整体设计与方案选型思路先聊清楚一个问题为什么是 OPC → MQTT → Kafka → InfluxDB 这个走向而不是其他组合。1.1 为什么中间要塞 Kafka有的朋友会问OPC 侧采集到数据直接通过 MQTT 发给 InfluxDB 不就行了吗实测下来不行至少在产线规模下不行。MQTT 虽然轻量但它的持久化能力很弱Broker 挂了消息就没了而且消费者的处理速度要是跟不上生产速度消息堆积在 Broker 内存里最终只能丢。工业数据的特点是频率高、单条小、总量大一台设备一秒可能出几十上百个点位几十台设备并发过来下游的时序库写入压力是瞬间的峰值。Kafka 在这个链路里的作用用一句话概括削峰填谷 多下游解耦。MQTT 把数据推过来之后Kafka 先把这个消息持久化到磁盘下游 InfluxDB 用自己的节奏去消费哪怕 InfluxDB 短暂不可用数据也不会丢恢复之后 Kafka 会继续把积压的数据吐出去。另外一点很实用Kafka 里的数据不只能进 InfluxDB还能同时接流处理引擎做实时计算或者再转发给另一个消息队列一套数据多家用不用再单独采集了。1.2 OPC、MQTT、InfluxDB 各自扮演的角色这三个组件在这个链路里互相补位。OPC 是数据源头。OPC DA 是老协议走 Windows COM/DCOM配置麻烦跨机器访问经常遇到权限问题OPC UA 是新一代标准跨平台、内置安全机制支持证书认证配置起来比 DA 省心很多。现在新项目建议直接上 UA老的 DA 设备通过 Kepware 这类网关软件做转换。MQTT 管设备侧到服务器侧的传输。它是个发布/订阅模型协议非常轻量一个报文头才几个字节在带宽受限的现场环境里很合适。而且 MQTT 支持 QoS 机制重要点位可以设成 QoS 1保证消息至少到达一次。InfluxDB 管数据落地。它就是为时序数据设计的写入和查询都用类 SQL 语法tag 做索引、field 存数值我实测单机写几万点位每秒没什么压力对这条链路来说完全够用。而且它自带数据保留策略可以自动清理老数据不需要手动维护。2. 核心环节实操从 OPC 侧把数据摸出来2.1 OPC UA 与 OPC DA 的接入差异如果你用的是 OPC DA第一关就是 DCOM 配置。我在用 C 调CoCreateInstanceEx连接 OPC DA Server 时遇到过0x80070005拒绝访问这个报错九成是 DCOM 权限问题。解决办法是在 Windows 组件服务里找到对应的 OPC Server 的 DCOM 配置把“启动和激活权限”和“访问权限”都加上匿名用户和交互用户。注意 OPC Core Components Redistributable 必须装x86 和 x64 版本的注册表路径不同32 位客户端连 64 位服务器时经常踩这个坑。OPC UA 就好办很多UA 默认只走 4840 端口通过配置文件或代码里指定 Endpoint 就能连。但要注意 UA 有安全策略None、Basic256Sha256这些要和服务端对齐证书认证的话客户端证书要先放到服务端的信任列表里否则握手直接失败。调试阶段可以直接用 UaExpert 这个免费客户端去试连先确认服务端地址和端口通不通再排查代码里的问题。2.2 点表映射与数据格式约定OPC 服务端的节点是树形结构每个变量节点有一个 NodeId。实际项目中你需要整理一份点表把设备名、点位名、NodeId、数据类型、采集频率列清楚。我通常会在一个配置文件里维护点表结构大概是下面这样devices: - name: welding_robot_01 endpoint: opc.tcp://192.168.1.101:4840 security_policy: Basic256Sha256 points: - name: current node_id: ns2;sRobot1.Current data_type: float interval_ms: 1000 - name: temperature node_id: ns2;sRobot1.Temp data_type: float interval_ms: 5000采集程序启动时读取这个配置订阅这些节点的数据变化。这里有一个很关键的取舍用订阅还是用轮询。OPC UA 的订阅机制是服务端主动推送变化省带宽、实时性好但有些老旧 OPC DA Server 的订阅回调不稳定轮询反而更可靠。我一般优先用订阅对实时性要求不高的点位比如温度、压力这种慢变量用 3~5 秒的轮询兜底避免订阅异常时丢数据。2.3 MQTT 主题规划与 QoS 选择OPC 采集程序拿到数据后要先做一步转换把点位数据整理成统一的 JSON 结构再发布到 MQTT。这一步别偷懒直接在采集器里拼一个符合规范的 payload后面 Kafka 和 InfluxDB 的消费逻辑都会省很多事。我用的 payload 格式是{ device: welding_robot_01, point: current, timestamp: 1712304000000, value: 12.56, quality: 192 }MQTT 主题建议按设备维度拆分这种结构在后面做权限控制和数据过滤时会很顺手比如plant1/line1/welding_robot_01/dataQoS 的设置有个原则不是越高越好。QoS 2 的握手流程繁重吞吐量会下降一大截工业场景一般用 QoS 1 就够保证消息至少到达一次Broker 挂了重连之后也能续传。只有像设备启停、报警这种关键事件才值得用 QoS 2。订阅侧默认 QoS 0 就行因为 Kafka 这一层已经做了持久化不需要 MQTT 再重复保障。3. Kafka 层的设计与实际落地3.1 从 MQTT 到 Kafka 的桥接这一步有现成的轮子EMQX 自带的 Rule Engine 可以直接把 MQTT 消息转发到 Kafka不需要自己写代码。但如果你用的是 Mosquitto 这类轻量 Broker就得自己写桥接程序。我在工程里是写了一个 Go 的小程序做桥接逻辑很简单订阅 MQTT 主题 → 解析 JSON → 规范化消息 → 生产到 Kafka。Go 的好处是部署简单、并发模型适合消息转发场景。核心代码思路func main() { // 连接 MQTT opts : mqtt.NewClientOptions().AddBroker(tcp://localhost:1883) client : mqtt.NewClient(opts) client.Subscribe(plant1/#, 0, func(client mqtt.Client, msg mqtt.Message) { // 生产到 Kafka producer.Send(opc_raw_data, msg.Payload()) }) }这一步值得注意的一点Kafka 的 topic 我建议按数据类型分而不是按设备分。比如opc_raw_data放全部原始数据opc_alarm放告警事件。原因是 Kafka 的分区机制决定了消费并行度一个 topic 只设置一个分区的话下游只能单线程消费吞吐量上不去。多个 topic 可以在消费端独立扩展消费者实例。3.2 分区数、副本数的计算思路Kafka 的性能调优在工业场景里其实不需要太复杂的配置但有两个参数必须提前想好分区数和副本数。分区数决定消费并发度。我的经验公式是分区数 ≥ 预期的最大消费者数。比如你计划开 3 个消费者实例写 InfluxDB分区数至少设 3 个最好留点余量设 6 个。副本数决定容错能力如果只有单台 Kafka Broker副本数设 1 就行多了反而增加同步开销。这里再补充一个经常被忽略的点Kafka 的分区键。如果直接用默认的随机分区同一个设备的数据会分散到不同分区InfluxDB 写入时标签的规整性会受影响而且如果你未来想按设备维度做数据回溯跨分区查询会麻烦。生产者在发送消息时用device字段做分区键保证同一设备的数据进同一分区消费端处理和按设备聚合都更高效。3.3 消费者写 InfluxDB 的批量与缓冲策略从 Kafka 消费写 InfluxDB 的环节最大的坑是写入频率和批量大小。InfluxDB 对批量写入的友好度远高于单条写入我在代码里用了一个简单的批量缓冲策略攒够 5000 条或者 2 秒超时就触发一次批量写入。这样单次写入的吞吐量能从每秒几百条提升到每秒几万条。还要注意 InfluxDB 写入的数据模型。一个 measurement 相当于一张表tag 是索引字段field 是数值字段。设备点位适合做 tag采集到的值适合做 field。下面是我惯用的 schemaopc_data,devicewelding_robot_01,pointcurrent value12.56 1712304000000000000timestamp 这里必须用纳秒单位因为 InfluxDB 的默认时间精度是纳秒。我踩过一个坑从 Kafka 拿到的消息里时间是毫秒单位直接写入后 InfluxDB 把时间戳全当作 1970 年处理排查了半天才发现单位搞错了。4. 一次部署用 Docker Compose 把整条链路拉起来4.1 工程文件里的部署编排我把整套依赖服务打包成了一个docker-compose.yml在工程机上一键启动。EMQX 做 MQTT Broker单节点 Kafka 做消息缓冲InfluxDB 1.8 做时序存储采集器和桥接程序也用容器跑。这份编排文件是这条链路能不能快速重现的关键。version: 3 services: emqx: image: emqx/emqx:5.1.0 container_name: emqx ports: - 1883:1883 - 18083:18083 kafka: image: bitnami/kafka:3.5 container_name: kafka environment: - KAFKA_CFG_NODE_ID0 - KAFKA_CFG_PROCESS_ROLEScontroller,broker - KAFKA_CFG_CONTROLLER_QUORUM_VOTERS0kafka:9093 - KAFKA_CFG_LISTENERSPLAINTEXT://:9092,CONTROLLER://:9093 - KAFKA_CFG_ADVERTISED_LISTENERSPLAINTEXT://kafka:9092 - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAPCONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT ports: - 9092:9092 influxdb: image: influxdb:1.8 container_name: influxdb environment: - INFLUXDB_DBopc - INFLUXDB_ADMIN_USERadmin - INFLUXDB_ADMIN_PASSWORDadmin123 volumes: - ./influxdb-data:/var/lib/influxdb ports: - 8086:8086要注意 Kafka 镜像的版本和配置项在不同镜像仓库之间有差异。我之前用 bitnami 镜像时KRaft 模式的配置项是这种写法换用 wurstmeister 镜像的话配置项又完全是另一套。工程文件里最好锁死镜像版本免得别人拉下来跑不了。4.2 实战一条命令启动全部服务在安装好 Docker 的机器上进入工程目录直接执行docker-compose up -d然后把 OPC 采集器的配置文件和这个组合文件放一起把 OPC Server 的地址改成实际地址启动采集器。整个链路从底层服务到应用层全部容器化后续在别的机器上复现环境非常方便。这里有一个部署细节提醒OPC UA 的采集器如果也要容器化网络模式最好用network_mode: host。因为 OPC UA 的服务端有时候会和客户端协商动态端口常规的桥接网络会导致连接失败。在docker-compose.yml里直接设定宿主机网络最省心。4.3 InfluxDB 的保留策略与写入通道调优说到 InfluxDB 1.8 和 2.x 的区别这个我要单独提一下。搜索热词里很多人问到 token 问题那是 InfluxDB 2.x 的认证方式1.x 用的是用户名密码。2.x 的 token 需要先在 Web UI 里创建之后所有请求的 Header 里都要带Authorization: Token xxx。还遇到过 InfluxDB 2.x 写入报 401 的情况基本都是 token 权限只给了读没给写在 UI 里重新建一个Read/Write权限的 token 就能解决。在 InfluxDB 1.x 里一定要建好数据保留策略Retention Policy不然/var/lib/influxdb会被数据撑爆。我实测过生产环境磁盘写满后 InfluxDB 会直接拒绝写入日志里报engine: error writing wal entry: write /var/lib/influxdb/wal/xxx/autogen: no space left on device这条链路会从最末端开始堵一直反压到 Kafka 消费积压。所以我的工程文件里默认配了一条 30 天的保留策略curl -XPOST http://localhost:8086/query \ --data-urlencode qCREATE RETENTION POLICY rp_30d ON opc DURATION 30d REPLICATION 1 DEFAULT5. 常见问题与排查技巧实录这一套链路牵扯的组件多出现问题后定位周期相对较长。我把自己踩过的坑整理成一个速查表按“现象 → 原因 → 解法”的思路展开方便你直接抄作业。现象大概率原因解决思路OPC DA 连接报0x80070005DCOM 权限不足Windows 组件服务中修改 OPC Server 的 DCOM 启动/访问权限加入匿名用户连接 OPC UA 时证书校验失败客户端证书未加入服务端信任列表导出客户端证书在服务端 UA Server 的信任列表中添加后重启MQTT 订阅不到任何消息主题层级不匹配或 Broker 未开启匿名访问先检查订阅主题是否带通配符再确认 EMQX 的 ACL 是否放行对应主题Kafka 报error while fetching metadataadvertised.listeners地址不对容器内必须配置为其他服务能访问到的地址本机调试改为localhost:9092写入 InfluxDB 速率极低批量太小或时间精度错误改为批量写入每次至少上千条确认时间戳精度是纳秒且单位对齐InfluxDB 重启后数据全部丢失未挂载数据卷检查容器volumes是否映射了宿主机目录到/var/lib/influxdbKafka 消息延迟持续升高消费者处理能力不足或分区数过少提高消费批量大小、增加分区数需在创建 topic 前确定或增加消费者实例还有一个高频问题要单独展开一下Kafka 延迟消费。很多场景希望数据延迟一段时间再处理比如告警去抖我见过有同事用Thread.sleep硬延迟这会让消费者线程被占用积压会越来越严重。更稳的做法是利用 Kafka 本身的时间戳机制消费者读到消息后判断时间差未到预定时间的消息先发回一个延迟队列或直接塞进本地时间轮等时间到了再处理。关于 MQTT 的 Broker 选型实测下来 EMQX 对工业场景的适配度最高和 Kafka 的桥接也最顺畅但它的内存占用和配置复杂度会高于 Mosquitto。如果只是测试环境、点数不超过百个用一个轻量的 Mosquitto 单机就能扛住上了规模再换 EMQX 不迟。最后再分享一个个人经验这条链路第一次搭建的时候别急着把全流程串起来一定分阶段验证。先单独测 OPC 采集确认点位数据能读出来再单独测 MQTT 收发用 MQTT X 这类客户端看主题消息是否正常然后测 Kafka 的生产消费最后才拼装到 InfluxDB 写入。如果一次全通了多半是运气好出了问题你根本不知道是哪一层挂的。先把每一环都按住再谈整条链路的稳定性。这样排查问题的时候每个环节都能自证清白省下的时间足够你把整条链路再优化两个来回。这套采集服务打包之后我把它放在工程目录里复用了好几个项目每次部署新的产线环境时都不用再从零开始。如果你正在做类似的设备数据接入项目照这个链路搭一遍框架把点位表换成自己的设备基本就能跑通。后面你也可以在这个基础上继续扩展比如把数据从 InfluxDB 接到 Grafana 做实时看板或者从 Kafka 接 Flink 做实时计算都是顺理成章的事情。本文还有配套的精品资源点击获取