OPC UA多节点数据变更通知同步写入CSV时数据丢失问题排查
问题:OPC UA多节点订阅数据写入CSV丢失问题
我同时订阅多个OPC UA节点,仅在节点值变化时接收数据变更通知,需要将时间戳+所有节点值写入CSV文件。使用asyncio.Queue处理并发写入,但同一时间收到多个通知时,仅第一个通知的数据写入CSV,其余数据丢失。
程序输出
sv_writer_task started Queue list: <Queue maxsize=0> Datachange notification received. 29032024 14:10:16Data change notification was received and queued: ns=2;s=Unit_SB.UC3.TEMP.TEMP_v 19.731399536132812 Datachange notification received. 29032024 14:10:16Data change notification was received and queued: ns=2;s=Unit_SB.UC3.HUMI.HUMI_v 36.523399353027344 Dequeue value: 19.731399536132812 val name: TEMP Prefix UC1 is no match for node: Unit_SB.UC3.TEMP.TEMP_v. Skipping. Prefix UC2 is no match for node: Unit_SB.UC3.TEMP.TEMP_v. Skipping. csv file: UC3.csv Appending new row Data written to file UC3.csv succesfully Dequeue value: 36.523399353027344 val name: HUMI Prefix UC1 is no match for node: Unit_SB.UC3.HUMI.HUMI_v. Skipping. Prefix UC2 is no match for node: Unit_SB.UC3.HUMI.HUMI_v. Skipping. csv file: UC3.csv Updating existing row Data written to file UC3.csv succesfully
CSV文件结果
Timestamp,OUT_G,OUT_U,HUMI,TEMP 29032024 14:08:42,True,True,47.38769912719727,15.043899536132812 29032024 14:10:16,True,True,47.38769912719727,19.731399536132812
原始代码
# 从队列消费数据的任务 async def csv_writer_task(queue): print("csv_writer_task started") print('Queue list: ') print(queue) print('') while True: try: dte, node, val = await queue.get() print('Dequeue value: ', val) except Exception as e: print(f"csv_writer_task从队列取数据时出错: {e}") node_id_str = str(node.nodeid.Identifier) node_parts = node_id_str[len("Unit_SB."):].split('.') val_name = node_parts[-1].replace('_v', '') print('val name: ', val_name) for key, header_row in prefix_mapping.items(): if f"Unit_SB.{key}" in node_id_str: csv_file = f"{key}.csv" print('csv file: ', csv_file) break else: print(f"前缀{key}与节点{node_id_str}不匹配,跳过。") df = pd.read_csv(csv_file) last_row = df.iloc[-1].copy() if last_row['Timestamp'] == dte: print("更新现有行") last_row[val_name] = val else: print("追加新行") new_row = last_row.copy() new_row['Timestamp'] = dte new_row[val_name] = val df = pd.concat([df, new_row.to_frame().T], ignore_index=True) df.to_csv(csv_file, index=False) print(f'数据成功写入文件{csv_file}') queue.task_done() class SubscriptionHandler(object): def __init__(self, wsconn): self.wsconn = wsconn self.q = asyncio.Queue() self.queuwriter = asyncio.create_task(csv_writer_task(self.q)) # 生产者:将数据放入队列 async def datachange_notification(self, node, val, data): print("收到数据变更通知。") dte = data.monitored_item.Value.ServerTimestamp.strftime("%d%m%Y %H:%M:%S") await self.q.put((dte, node, val)) print(dte + "数据变更通知已接收并放入队列: ", node, val) class SBConnection(): def __init__(self): self.listOfWSNode = [] self.dpsList = ... async def connectAndSubscribeToServer(self): self.csv_file = '' async with Client(url=self.url) as self.client: for element in self.dpsList: node = "ns=" + element["NS"] + ";s=" + element["Name"] var = self.client.get_node(node) self.listOfWSNode.append(var) print("订阅节点: ", self.listOfWSNode) handler = SubscriptionHandler(self) sub = await self.client.create_subscription(period=10, handler=handler) await sub.subscribe_data_change(self.listOfWSNode) print('订阅已创建') # 保持程序运行 while True: await asyncio.sleep(0.1) async def main(): uc = SBConnection() await uc.connectAndSubscribeToServer() if __name__ == '__main__': asyncio.run(main())
修正后代码
async def csv_writer_task(queue): print("csv_writer_task started") print('Queue list: ') print(queue) print('') while True: try: dte, node, val = await queue.get() print('Dequeue value: ', val) except Exception as e: print(f"csv_writer_task从队列取数据时出错: {e}") node_id_str = str(node.nodeid.Identifier) node_parts = node_id_str[len("Unit_SB."):].split('.') val_name = node_parts[-1].replace('_v', '') print('val name: ', val_name) for key, header_row in prefix_mapping.items(): if f"Unit_SB.{key}" in node_id_str: csv_file = f"{key}.csv" print('csv file: ', csv_file) break else: print(f"前缀{key}与节点{node_id_str}不匹配,跳过。") df = pd.read_csv(csv_file) if df.iloc[-1]['Timestamp'] == dte: print("更新现有行") df.loc[df.index[-1], val_name] = val # 直接修改DataFrame中的最后一行 else: print("追加新行") new_row = df.iloc[-1].copy() new_row['Timestamp'] = dte new_row[val_name] = val df = pd.concat([df, new_row.to_frame().T], ignore_index=True) df.to_csv(csv_file, index=False) print(f'数据成功写入文件{csv_file}') queue.task_done()
问题原因与修复说明
原始代码的问题
原始代码中,当检测到时间戳一致需要更新现有行时,仅修改了从DataFrame复制出来的last_row对象,但未将修改同步回原DataFrame。最终写入CSV的还是未修改的原DataFrame,导致HUMI的新值丢失。
last_row = df.iloc[-1].copy() if last_row['Timestamp'] == dte: print("更新现有行") last_row[val_name] = val # 仅修改了副本,原df未变化
修正方案
修正后直接操作原DataFrame的最后一行,通过df.loc[df.index[-1], val_name] = val修改对应单元格,确保修改会被同步到后续的CSV写入操作中。这样同一时间戳的第二个节点数据修改会直接作用于最新的DataFrame,HUMI的新值就能被正确保存。
内容的提问来源于stack exchange,提问作者Marm90
相关产品推荐
相关产品推荐

