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
相关产品推荐
相关产品推荐

