异步追加CSV文件出现损坏与重复表头的问题排查
问题分析与修复方案
你的代码出现CSV内容错位、表头重复的核心原因是锁的使用完全错误,再加上几个逻辑漏洞,具体问题和修复方法如下:
核心问题点
- 局部锁无法实现跨任务互斥:你在每次写入表头或数据时新建
asyncio.Lock(),每个任务的锁都是独立实例,相当于完全没加锁——多个任务可以同时写入同一个文件,直接导致内容错乱。 - 表头写入的竞态条件:检查文件大小和写入表头是两个分离的操作,多个任务会同时检测到文件为空,进而各自写入表头,造成重复。
- random.shuffle误用:
random.shuffle()是原地修改列表,返回值为None。当prefix == "file_one"时,item会被赋值为None,后续遍历会抛出TypeError。 - 文件模式冗余:
a+模式允许读写,但你只需要追加写入权限,额外的读权限可能引入指针位置的潜在问题。
修复后的代码
import asyncio import aiofiles import aiofiles.os import aiocsv import uuid import random import json from pathlib import Path from datetime import datetime, timezone # 为每个文件维护全局锁,key为文件路径的字符串形式 file_locks = {} async def get_file_lock(file_path: Path) -> asyncio.Lock: """获取指定文件的全局锁,不存在则创建""" key = str(file_path.resolve()) if key not in file_locks: file_locks[key] = asyncio.Lock() return file_locks[key] async def write_csv(item: list, load_id: str, prefix: str) -> None: Path("./test_files").mkdir(parents=True, exist_ok=True) file_path = Path("./test_files").joinpath(f"{prefix}_{load_id}.csv") file_lock = await get_file_lock(file_path) # 整个文件写入过程持有全局锁,确保同一时间只有一个任务操作该文件 async with file_lock: # 改用a模式(仅追加写入),无需a+的读权限 async with aiofiles.open(file_path, mode="a", newline="") as f: w: aiocsv.AsyncWriter = aiocsv.AsyncWriter(f) file_size = await aiofiles.os.path.getsize(file_path) print(f"INFO: writing file: {file_path.resolve()}, current size: {file_size}") # 原子性完成「检查空文件+写入表头」操作 if file_size == 0: print(f"file {file_path.name} was empty! writing header") await w.writerow([ "response", "load_id", "last_updated_timestamp_utc" ]) # 修复shuffle逻辑:直接原地修改列表,不要赋值回item if prefix == "file_one": random.shuffle(item) # 写入数据,此时已持有锁,无需额外加锁 for chunk in item: await w.writerow([ json.dumps(chunk), load_id, datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S") ]) async def main() -> None: items: list[str] = [["hello", "world"], ["asyncio", "test"]] * 500 prefixes: list[str] = ["file_one", "file_two"] tasks: list = [] load_id = str(uuid.uuid4()) for i in items: task = asyncio.create_task(write_csv(i, load_id, random.choice(prefixes))) tasks.append(task) errors = await asyncio.gather(*tasks, return_exceptions=True) # 可选:打印错误排查问题 # for err in errors: # if err: # print(f"Error occurred: {err}") if __name__ == "__main__": loop = asyncio.new_event_loop() loop.run_until_complete(main())
关键修复说明
- 全局文件锁:通过
file_locks字典为每个文件维护唯一的锁实例,确保同一文件的所有写入任务共用一把锁,真正实现互斥访问。 - 原子性表头操作:将「检查文件大小+写入表头」放在同一个锁块内,保证只有第一个任务能写入表头,后续任务不会重复执行。
- 修正shuffle逻辑:直接对
item原地执行random.shuffle(),避免因赋值None导致的遍历错误。 - 简化文件模式:改用
a模式打开文件,只保留追加写入权限,减少不必要的指针操作风险。
内容的提问来源于stack exchange,提问作者Coldchain9
相关产品推荐
相关产品推荐

