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

FastAPI实时文件扫描最优架构咨询:Celery还是多进程?

FastAPI 下杀毒工具后台扫描与实时监控的最佳架构方案

核心问题分析

当前扫描接口阻塞、速度慢的核心原因:

  • 普通循环/基础异步函数无法突破GIL限制,无法充分利用多核CPU(文件扫描属于CPU+IO混合密集型任务)
  • 同步接口会占用FastAPI事件循环,导致其他请求被阻塞
  • 大型文件系统扫描需要分批次后台处理,不能占用API主线程

推荐架构选型

1. Celery + Redis/RabbitMQ(首选,适配复杂后台任务场景)

适合场景:后台批量扫描、USB实时扫描触发的异步任务、任务状态追踪
优势:

  • 完全解耦API与扫描任务,FastAPI仅负责接收请求、触发任务、返回任务ID
  • 支持任务队列、优先级调度、任务重试、状态查询
  • 可横向扩展worker节点,应对大规模文件扫描需求

示例代码

1.1 Celery配置
# celery_config.py
from celery import Celery

celery_app = Celery(
    "antivirus_scanner",
    broker="redis://localhost:6379/0",
    backend="redis://localhost:6379/0"
)

celery_app.conf.task_serializer = 'json'
celery_app.conf.result_serializer = 'json'
celery_app.conf.accept_content = ['json']
1.2 扫描任务定义
# tasks.py
from celery_config import celery_app
import os
import hashlib

def calculate_sha256(file_path, chunk_size=4096):
    sha256 = hashlib.sha256()
    with open(file_path, 'rb') as f:
        while chunk := f.read(chunk_size):
            sha256.update(chunk)
    return sha256.hexdigest()

def scan_single_file(file_path):
    # 替换为你的恶意文件检测逻辑
    sha256_hash = calculate_sha256(file_path)
    malicious_hashes = ["abc123...", "def456..."]  # 模拟恶意哈希库
    return {
        "file_path": file_path,
        "sha256": sha256_hash,
        "is_malicious": sha256_hash in malicious_hashes
    }

@celery_app.task(bind=True)
def scan_directory_task(self, dir_path):
    scan_results = []
    total_files = sum(len(files) for _, _, files in os.walk(dir_path))
    processed = 0
    
    for root, _, files in os.walk(dir_path):
        for file in files:
            file_path = os.path.join(root, file)
            try:
                result = scan_single_file(file_path)
                scan_results.append(result)
                processed += 1
                # 更新任务进度
                self.update_state(
                    state='PROGRESS',
                    meta={'current': processed, 'total': total_files}
                )
            except Exception as e:
                scan_results.append({"file_path": file_path, "error": str(e)})
                processed += 1
    return {"results": scan_results}
1.3 FastAPI接口集成
# main.py
from fastapi import FastAPI
from celery_config import celery_app
from tasks import scan_directory_task
from celery.result import AsyncResult

app = FastAPI()

@app.post("/scan/directory")
async def start_directory_scan(dir_path: str):
    task = scan_directory_task.delay(dir_path)
    return {"task_id": task.id, "status": "started"}

@app.get("/scan/status/{task_id}")
async def get_scan_status(task_id: str):
    result = AsyncResult(task_id, app=celery_app)
    if result.state == 'PENDING':
        return {"status": result.state}
    elif result.state == 'PROGRESS':
        return {"status": result.state, "progress": result.info}
    else:
        return {"status": result.state, "results": result.info}

2. 多进程(适合轻量场景,快速落地)

适合场景:小型文件扫描、无需复杂任务调度的情况
优势:

  • 无需额外依赖消息队列,代码简洁
  • 利用multiprocessing绕过GIL,充分利用多核CPU

示例代码

# main.py
from fastapi import FastAPI
from multiprocessing import Pool, cpu_count
import os
import hashlib

app = FastAPI()

def scan_single_file(file_path):
    sha256 = hashlib.sha256()
    try:
        with open(file_path, 'rb') as f:
            while chunk := f.read(4096):
                sha256.update(chunk)
        malicious_hashes = ["abc123...", "def456..."]
        return {
            "file_path": file_path,
            "sha256": sha256.hexdigest(),
            "is_malicious": sha256.hexdigest() in malicious_hashes
        }
    except Exception as e:
        return {"file_path": file_path, "error": str(e)}

@app.post("/scan/directory")
async def scan_directory(dir_path: str):
    file_list = []
    for root, _, files in os.walk(dir_path):
        for file in files:
            file_list.append(os.path.join(root, file))
    
    # 使用CPU核心数创建进程池
    with Pool(cpu_count()) as pool:
        results = pool.map(scan_single_file, file_list)
    return {"results": results}

注意:该方式下接口仍会等待所有进程完成后返回,适合小型目录扫描。若要实现非阻塞,可结合后台线程+进程池,或使用FastAPI的BackgroundTasks(仅适合轻量短任务,不适合长时间扫描)。

3. USB实时扫描实现思路

结合系统API监听USB插入事件,自动触发后台扫描:

  • Linux:使用pyudev库监听udev设备事件
  • Windows:使用pywin32监听WM_DEVICECHANGE消息
  • macOS:使用pyobjc监听IOKit设备通知

检测到USB挂载后,自动调用Celery任务或多进程扫描挂载目录,通过WebSocket向前端实时推送任务状态。

性能优化建议

  • 分块处理大文件:计算SHA256时用固定大小块读取,避免一次性加载大文件到内存
  • 任务优先级调度:给USB扫描任务设置更高优先级,优先处理
  • 扫描结果缓存:记录已扫描文件的SHA256结果,避免重复扫描
  • 增量扫描:通过文件mtime判断,仅扫描新增或修改的文件

内容的提问来源于stack exchange,提问作者Ahsan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.01 13:24:57