如何在asyncio协程间共享Pandas DataFrame?pd.concat副本问题
解决asyncio协程并发修改同一Pandas DataFrame的问题
核心问题分析
pd.concat会返回新的DataFrame对象,如果你在协程中执行df = pd.concat([df, new_row_df]),只是修改了当前协程内df变量的引用,原DataFrame对象并未被修改,其他协程仍指向旧对象,这就是内存地址变化的原因。而append是原地修改(虽然已弃用),所以能保持对象一致。
可行解决方案
方案1:使用df.loc原地添加行 + asyncio锁保证原子性
asyncio是单线程并发,但协程切换发生在await点,若多个协程同时执行df.loc添加行,可能因同时获取len(df)导致行覆盖。用asyncio.Lock确保每次修改是原子操作,所有协程操作的都是同一个DataFrame对象。
示例代码:
import asyncio import pandas as pd async def add_row(df, lock, row_data): async with lock: # 原地添加行,修改原DataFrame对象 df.loc[len(df)] = row_data async def main(): # 初始化共享DataFrame shared_df = pd.DataFrame(columns=['col1', 'col2']) lock = asyncio.Lock() # 创建多个协程任务 tasks = [ add_row(shared_df, lock, {'col1': 1, 'col2': 'a'}), add_row(shared_df, lock, {'col1': 2, 'col2': 'b'}), add_row(shared_df, lock, {'col1': 3, 'col2': 'c'}) ] await asyncio.gather(*tasks) print(shared_df) if __name__ == "__main__": asyncio.run(main())
方案2:先收集所有协程的行数据,最后一次性合并
这种方式更高效,尤其是数据量较大时,避免多次原地修改的性能损耗。协程只负责生成行数据,最后统一合并到原DataFrame。
示例代码:
import asyncio import pandas as pd async def generate_row(row_data): # 模拟IO操作,比如从API/数据库获取数据 await asyncio.sleep(0.1) return row_data async def main(): shared_df = pd.DataFrame(columns=['col1', 'col2']) row_datas = [{'col1': 1, 'col2': 'a'}, {'col1': 2, 'col2': 'b'}, {'col1': 3, 'col2': 'c'}] # 并发生成所有行数据 tasks = [generate_row(data) for data in row_datas] results = await asyncio.gather(*tasks) # 一次性合并到原DataFrame(直接赋值更新引用即可) shared_df = pd.concat([shared_df, pd.DataFrame(results)], ignore_index=True) print(shared_df) if __name__ == "__main__": asyncio.run(main())
注意事项
- 方案1的原地修改适合小批量数据,频繁修改可能导致Pandas性能下降,因为DataFrame底层的数组结构是不可变的,扩容时会重新分配内存,但对象引用始终保持一致。
- 即使是单线程asyncio,修改共享状态时也要用锁,避免协程切换导致的竞态条件(比如两个协程同时计算
len(df),然后写入同一索引行)。
内容的提问来源于stack exchange,提问作者Tim Cassidy
相关产品推荐
相关产品推荐

