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

异步追加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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 05:44:51