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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 09:43:18