高频OPC Server数据存储咨询:每秒2万点无丢失且支持并行访问
解决方案:高频OPC数据无丢失存储与并行访问方案
核心问题分析
用to_csv丢失90%数据的根本原因:
- CSV是同步阻塞IO,每秒2万条数据的写入速度远超出磁盘单线程处理能力
- 未做数据缓冲与异步解耦,接收线程被IO阻塞导致数据丢包
推荐存储方案
1. 时序数据库(优先推荐)
时序数据库天生适配高频时序数据场景,支持高写入吞吐量、并行读写,完美匹配ANN训练的批量数据读取需求。
- InfluxDB:用
influxdb-client实现异步批量写入from influxdb_client import InfluxDBClient, Point from influxdb_client.client.write_api import ASYNCHRONOUS import asyncio # 初始化异步写入客户端 client = InfluxDBClient(url="http://localhost:8086", token="your-token", org="your-org") write_api = client.write_api(write_options=ASYNCHRONOUS) bucket = "opc_data_bucket" async def write_opc_batch(data_points): # 批量构造时序数据点 points = [ Point("opc_metrics") .tag("device", dp["device"]) .field("value", dp["value"]) .time(dp["timestamp"]) for dp in data_points ] write_api.write(bucket=bucket, record=points) # 接收OPC数据时,每攒1000-2000条批量提交 # 需用异步队列缓冲数据,避免阻塞接收线程 - TimescaleDB:基于PostgreSQL,支持SQL查询,适合需要复杂数据预处理后喂给ANN的场景,用
asyncpg实现异步写入
2. 内存缓冲+Parquet列式存储
无需部署数据库时,采用内存队列缓冲+Parquet文件组合:
- 用
queue.Queue或asyncio.Queue做数据缓冲,接收线程仅负责往队列塞数据,单独开写入线程处理IO - Parquet是列式存储,写入速度快、压缩率高,支持Pandas/PySpark并行加载,适配ANN训练的批量读取需求
import queue import pandas as pd from threading import Thread import time data_queue = queue.Queue(maxsize=100000) # 超大缓冲队列避免溢出 BATCH_SIZE = 2000 # 每批写入条数 def opc_data_receiver(): # 模拟OPC数据接收逻辑 while True: data_point = {"timestamp": time.time(), "device": "sensor_1", "value": 123.45} data_queue.put(data_point) def parquet_writer(): buffer = [] while True: # 攒够批量大小再写入 while len(buffer) < BATCH_SIZE: try: buffer.append(data_queue.get(timeout=1)) except queue.Empty: continue # 写入Parquet文件,按时间分片避免单文件过大 df = pd.DataFrame(buffer) file_name = f"opc_data_{int(time.time())}.parquet" df.to_parquet(file_name, engine="pyarrow") buffer = [] # 启动线程 receiver_thread = Thread(target=opc_data_receiver, daemon=True) writer_thread = Thread(target=parquet_writer, daemon=True) receiver_thread.start() writer_thread.start()
3. 异步二进制文件写入(应急方案)
若必须用文件存储,放弃CSV,采用异步IO+二进制序列化:
- 用
aiofiles做异步写入,配合msgpack序列化数据(比JSON快数倍)
注意:二进制文件需自行处理索引,仅适合临时存储后批量导入到其他存储系统import aiofiles import msgpack import asyncio async def async_write_batch(data_points): async with aiofiles.open("opc_data.msgpack", "ab") as f: packed = msgpack.packb(data_points) await f.write(packed) # 接收数据时批量异步写入
关键注意事项
- 解耦接收与写入逻辑:永远不要让数据接收线程处理IO,用生产者-消费者模型拆分任务
- 批量写入:单条写入效率极低,每次攒1000-10000条再提交,大幅提升吞吐量
- 监控队列状态:若队列持续增长,说明写入速度跟不上,需调整批量大小、增加写入线程或升级存储硬件
内容的提问来源于stack exchange,提问作者BhargavN
相关产品推荐
相关产品推荐

