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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 02:34:58