FastAPI异步场景下文件并发读写竞态条件处理方案咨询
解决FastAPI文件读写并发问题的方案与实现
1. 避免多管理员同时修改的竞态条件
核心是通过独占锁确保同一时间仅一个写操作执行,根据部署场景选择不同锁实现:
单进程部署(单FastAPI实例)
利用asyncio.Lock实现单进程内的异步独占锁,FastAPI的异步事件循环会自动调度锁的获取与释放:
import asyncio from fastapi import FastAPI import aiofiles app = FastAPI() # 全局写锁,单进程内所有写请求共享 write_lock = asyncio.Lock() @app.post("/admin/update-file") async def update_file(content: str): async with write_lock: # 异步写入避免阻塞事件循环 async with aiofiles.open("data.txt", "w", encoding="utf-8") as f: await f.write(content) return {"status": "success"}
多进程/多实例部署
asyncio.Lock仅在单进程有效,需用跨进程锁:
文件系统锁(跨进程)
import fcntl import asyncio from fastapi import FastAPI import aiofiles app = FastAPI() @app.post("/admin/update-file") async def update_file(content: str): async with aiofiles.open("data.txt", "r+", encoding="utf-8") as f: file_fd = f.fileno() # 异步调用系统独占锁(Linux/macOS用fcntl,Windows替换为msvcrt.locking) await asyncio.to_thread(fcntl.flock, file_fd, fcntl.LOCK_EX) try: await f.seek(0) await f.write(content) await f.truncate() finally: await asyncio.to_thread(fcntl.flock, file_fd, fcntl.LOCK_UN) return {"status": "success"}
Redis分布式锁(跨机器多实例)
import asyncio from fastapi import FastAPI import aiofiles import aioredis app = FastAPI() redis = aioredis.from_url("redis://localhost") async def acquire_lock(lock_key: str, timeout: int = 10): # SETNX实现分布式锁,带过期时间防止死锁 return await redis.set(lock_key, "locked", ex=timeout, nx=True) is not None async def release_lock(lock_key: str): await redis.delete(lock_key) @app.post("/admin/update-file") async def update_file(content: str): lock_key = "file_write_lock" # 循环尝试获取锁,避免立即失败 while not await acquire_lock(lock_key): await asyncio.sleep(0.1) try: async with aiofiles.open("data.txt", "w", encoding="utf-8") as f: await f.write(content) finally: await release_lock(lock_key) return {"status": "success"}
2. 读写并发时的数据一致性
采用读写锁(Read-Write Lock):允许多个读请求并行,写请求独占;或用乐观锁实现冲突检测。
异步读写锁(单进程)
使用第三方库aiorwlock简化实现:
from aiorwlock import RWLock import aiofiles from fastapi import FastAPI app = FastAPI() rw_lock = RWLock() @app.get("/user/read-file") async def read_file(): # 读锁允许多个并发读 async with rw_lock.reader: async with aiofiles.open("data.txt", "r", encoding="utf-8") as f: content = await f.read() return {"content": content} @app.post("/admin/update-file") async def update_file(content: str): # 写锁独占资源 async with rw_lock.writer: async with aiofiles.open("data.txt", "w", encoding="utf-8") as f: await f.write(content) return {"status": "success"}
乐观锁(多实例场景)
给文件添加版本号,修改前校验版本避免覆盖:
import json import aiofiles from fastapi import FastAPI, HTTPException app = FastAPI() @app.post("/admin/update-file") async def update_file(content: str, current_version: int): async with aiofiles.open("data.json", "r+", encoding="utf-8") as f: data = json.loads(await f.read()) if data["version"] != current_version: raise HTTPException(status_code=409, detail="文件已被修改,请刷新后重试") data["content"] = content data["version"] += 1 await f.seek(0) await f.write(json.dumps(data)) await f.truncate() return {"status": "success", "new_version": data["version"]}
前端需先读取版本号,提交修改时携带该版本,冲突时提示用户刷新。
3. FastAPI中处理并发的高效技术工具
异步文件操作:aiofiles
FastAPI是异步框架,同步open会阻塞事件循环,必须用aiofiles实现异步IO,保证并发性能。
分布式锁:Redis
多实例部署时,Redis是最常用的分布式锁方案,配合aioredis异步客户端不阻塞事件循环。
替代文件:数据库
若文件操作逻辑复杂,建议迁移到关系型数据库(如PostgreSQL):
- 数据库天然支持事务、行级锁,解决并发问题更可靠;
- 示例(SQLAlchemy异步操作):
from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine from sqlalchemy.orm import declarative_base, sessionmaker from fastapi import FastAPI, Depends, HTTPException from sqlalchemy import Column, Integer, String app = FastAPI() engine = create_async_engine("postgresql+asyncpg://user:pass@localhost/db") AsyncSessionLocal = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False) Base = declarative_base() class FileData(Base): __tablename__ = "file_data" id = Column(Integer, primary_key=True) content = Column(String) version = Column(Integer, default=1) async def get_db(): async with AsyncSessionLocal() as session: yield session @app.get("/user/read-file") async def read_file(db: AsyncSession = Depends(get_db)): data = await db.get(FileData, 1) return {"content": data.content, "version": data.version} @app.post("/admin/update-file") async def update_file(content: str, current_version: int, db: AsyncSession = Depends(get_db)): data = await db.get(FileData, 1, with_for_update=True) if data.version != current_version: raise HTTPException(status_code=409, detail="文件已被修改") data.content = content data.version += 1 await db.commit() return {"status": "success", "new_version": data.version}
消息队列:Celery/RabbitMQ
若写操作耗时,可将任务放入消息队列异步执行,管理员提交后立即返回,后台通过锁保证单任务执行。
内容的提问来源于stack exchange,提问作者Br0k3nS0u1
相关产品推荐
相关产品推荐

