You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用influxdb-python从OPC UA服务器采集数据至InfluxDB?

从OPC UA服务器采集数据写入InfluxDB的Python实现

所需依赖库

  • opcua:用于连接OPC UA服务器并读取节点数据
  • influxdb_client:InfluxDB官方推荐的客户端(比旧版influxdb库更稳定)
  • python-dotenv(可选):用于管理配置变量,避免硬编码敏感信息

安装命令:

pip install opcua influxdb-client python-dotenv

实现步骤

  1. 配置连接参数:定义OPC UA服务器地址、目标节点ID,以及InfluxDB的URL、token、组织名、桶名
  2. 建立客户端连接:分别初始化OPC UA和InfluxDB的客户端实例并建立连接
  3. 定时采集数据:按设定间隔读取指定OPC UA节点的数值与时间戳
  4. 格式化数据:将采集到的数据转换为InfluxDB兼容的点结构
  5. 写入InfluxDB:通过InfluxDB客户端的写入API提交数据
  6. 资源清理:程序退出时关闭所有连接,处理异常场景

完整示例代码

import time
from opcua import Client
from influxdb_client import InfluxDBClient, Point
from influxdb_client.client.write_api import SYNCHRONOUS
from dotenv import load_dotenv
import os

# 加载环境变量(使用dotenv时需在项目根目录创建.env文件)
load_dotenv()

# 核心配置参数
OPCUA_SERVER_URL = os.getenv("OPCUA_SERVER_URL", "opc.tcp://localhost:4840/freeopcua/server/")
INFLUXDB_URL = os.getenv("INFLUXDB_URL", "http://localhost:8086")
INFLUXDB_TOKEN = os.getenv("INFLUXDB_TOKEN", "your-influxdb-token")
INFLUXDB_ORG = os.getenv("INFLUXDB_ORG", "your-org")
INFLUXDB_BUCKET = os.getenv("INFLUXDB_BUCKET", "your-bucket")
# 替换为实际需要采集的OPC UA节点ID
OPCUA_NODE_IDS = [
    "ns=2;i=2",
    "ns=2;i=3"
]
POLL_INTERVAL = 5  # 数据采集间隔(单位:秒)

def main():
    # 初始化客户端实例
    opcua_client = Client(OPCUA_SERVER_URL)
    influxdb_client = InfluxDBClient(url=INFLUXDB_URL, token=INFLUXDB_TOKEN, org=INFLUXDB_ORG)
    write_api = influxdb_client.write_api(write_options=SYNCHRONOUS)

    try:
        # 连接OPC UA服务器
        opcua_client.connect()
        print("已连接到OPC UA服务器")

        # 预获取目标节点对象
        nodes = [opcua_client.get_node(node_id) for node_id in OPCUA_NODE_IDS]

        while True:
            # 遍历采集每个节点数据
            for node in nodes:
                try:
                    value = node.get_value()
                    timestamp = node.get_data_value().SourceTimestamp
                    node_name = node.get_browse_name().Name

                    # 构造InfluxDB数据点
                    point = (
                        Point("opcua_data")
                        .tag("node_id", node.nodeid.to_string())
                        .tag("node_name", node_name)
                        .field("value", value)
                        .time(timestamp)
                    )

                    # 写入InfluxDB
                    write_api.write(bucket=INFLUXDB_BUCKET, org=INFLUXDB_ORG, record=point)
                    print(f"已写入数据: {node_name} = {value} @ {timestamp}")
                except Exception as e:
                    print(f"读取节点 {node.nodeid} 失败: {str(e)}")

            time.sleep(POLL_INTERVAL)

    except Exception as e:
        print(f"连接或采集过程出错: {str(e)}")
    finally:
        # 关闭所有连接
        opcua_client.disconnect()
        influxdb_client.close()
        print("已断开所有连接")

if __name__ == "__main__":
    main()

关键说明

  • OPC UA节点ID:需替换为实际要采集的节点ID,可通过UaExpert等OPC UA客户端工具查看目标节点的ID信息
  • InfluxDB配置:确保已在InfluxDB中创建对应的桶和组织,且token拥有该桶的写入权限
  • 异常扩展:示例仅包含基础异常捕获,生产环境可添加重试机制、日志记录等增强稳定性
  • 采集间隔:根据业务场景调整POLL_INTERVAL,避免过于频繁的请求给服务器造成压力

内容的提问来源于stack exchange,提问作者Simeon Ezekiel

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 01:17:15