如何使用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
实现步骤
- 配置连接参数:定义OPC UA服务器地址、目标节点ID,以及InfluxDB的URL、token、组织名、桶名
- 建立客户端连接:分别初始化OPC UA和InfluxDB的客户端实例并建立连接
- 定时采集数据:按设定间隔读取指定OPC UA节点的数值与时间戳
- 格式化数据:将采集到的数据转换为InfluxDB兼容的点结构
- 写入InfluxDB:通过InfluxDB客户端的写入API提交数据
- 资源清理:程序退出时关闭所有连接,处理异常场景
完整示例代码
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
相关产品推荐
相关产品推荐

