基于FastAPI的杀毒类安全工具实时文件扫描最优架构选型:Celery、多进程还是其他方案?
Great question! Let's break down the best architectures for your antivirus tool, focusing on solving the blocking API issue, enabling real-time USB scans, and handling large file systems efficiently.
For your use case—long-running scan tasks, background execution, and USB-triggered scans—Celery is the most robust solution. It fully decouples your FastAPI API from scan workloads, avoids blocking the event loop, and supports scalable, tracked task execution.
Why Celery?
- Complete separation of API and scan logic: FastAPI only submits tasks and returns task IDs, while Celery workers handle the heavy lifting
- Built-in task status tracking, retries, and priority management
- Horizontal scalability: Add more workers to handle large file systems or high scan volumes
- Easy integration with USB monitoring services (they can directly submit tasks to the Celery queue)
Example Implementation
Step 1: Celery Setup (celery_app.py)
from celery import Celery import hashlib import os # Initialize Celery with Redis as broker/backend celery_app = Celery( "antivirus_scanner", broker="redis://localhost:6379/0", backend="redis://localhost:6379/0" ) def calculate_sha256(file_path): """Calculate SHA256 in chunks to avoid memory overload with large files""" sha256_hash = hashlib.sha256() with open(file_path, "rb") as f: for byte_block in iter(lambda: f.read(4096), b""): sha256_hash.update(byte_block) return sha256_hash.hexdigest() @celery_app.task(bind=True) def scan_directory_task(self, directory_path): scanned_files = [] malicious_count = 0 total_files = sum(len(files) for _, _, files in os.walk(directory_path)) for root, _, files in os.walk(directory_path): for file in files: file_path = os.path.join(root, file) try: file_hash = calculate_sha256(file_path) # Replace with your actual malicious hash database check is_malicious = file_hash in {"fake_malicious_hash_1", "fake_malicious_hash_2"} scanned_files.append({ "path": file_path, "sha256": file_hash, "is_malicious": is_malicious }) if is_malicious: malicious_count += 1 # Update task progress for frontend polling self.update_state( state="PROGRESS", meta={ "current": len(scanned_files), "total": total_files, "malicious_found": malicious_count } ) except Exception as e: scanned_files.append({"path": file_path, "error": str(e)}) return { "total_scanned": len(scanned_files), "malicious_found": malicious_count, "details": scanned_files }
Step 2: FastAPI Integration (main.py)
from fastapi import FastAPI from celery_app import celery_app, scan_directory_task from celery.result import AsyncResult app = FastAPI(title="Antivirus Scanner API") # Submit a directory scan task @app.post("/api/scan/directory") async def submit_directory_scan(directory_path: str): task = scan_directory_task.delay(directory_path) return {"task_id": task.id, "status": "Task submitted"} # Check scan task status/progress @app.get("/api/scan/status/{task_id}") async def get_scan_status(task_id: str): task_result = AsyncResult(task_id, app=celery_app) match task_result.state: case "PENDING": return {"status": "Pending"} case "PROGRESS": return {"status": "In progress", "progress": task_result.info} case "SUCCESS": return {"status": "Completed", "result": task_result.info} case _: return {"status": "Failed", "error": str(task_result.info)} # USB-triggered scan endpoint (called by your USB monitoring service) @app.post("/api/scan/usb") async def submit_usb_scan(usb_mount_path: str): task = scan_directory_task.delay(usb_mount_path) return {"task_id": task.id, "status": "USB scan initiated"}
Step 3: Run Services
- Start Redis (required for Celery broker/backend)
- Start Celery worker:
(Adjustcelery -A celery_app worker --loglevel=info --concurrency=4concurrencyto match your CPU core count) - Start FastAPI:
uvicorn main:app --reload
If you want to avoid adding Celery/message queue complexity (good for small-scale deployments), use concurrent.futures.ProcessPoolExecutor to offload CPU-intensive scan tasks to separate processes (bypasses Python's GIL, critical for SHA256 calculations).
Example Implementation
from fastapi import FastAPI from concurrent.futures import ProcessPoolExecutor import hashlib import os from typing import Dict, List app = FastAPI(title="Lightweight Antivirus Scanner") app.state.executor = ProcessPoolExecutor(max_workers=4) app.state.tasks = {} # In-memory task store (use Redis for production) def scan_single_file(file_path: str) -> Dict: try: sha256_hash = hashlib.sha256() with open(file_path, "rb") as f: for byte_block in iter(lambda: f.read(4096), b""): sha256_hash.update(byte_block) file_hash = sha256_hash.hexdigest() return { "path": file_path, "sha256": file_hash, "is_malicious": file_hash in {"fake_malicious_hash_1"} } except Exception as e: return {"path": file_path, "error": str(e)} @app.post("/api/scan/directory") async def scan_directory(directory_path: str): file_paths = [os.path.join(root, f) for root, _, files in os.walk(directory_path) for f in files] # Submit scan to process pool task = app.state.executor.submit(lambda: list(map(scan_single_file, file_paths))) task_id = id(task) app.state.tasks[task_id] = task return {"task_id": task_id, "status": "Started"} @app.get("/api/scan/status/{task_id}") async def get_status(task_id: int): task = app.state.tasks.get(task_id) if not task: return {"error": "Task not found"} if task.done(): return {"status": "Completed", "result": task.result()} return {"status": "In progress"} @app.on_event("shutdown") async def shutdown(): app.state.executor.shutdown()
To detect USB insertion and trigger scans, build a separate monitoring service (runs alongside your API):
- Linux: Use
pyudevto listen for udev block device events, then call the/api/scan/usbendpoint - Windows: Use
pywin32to listen forWM_DEVICECHANGEmessages and get the USB drive letter - macOS: Use
pyobjcto monitor IOKit device connection events
Example Linux USB Monitor:
import pyudev import requests context = pyudev.Context() monitor = pyudev.Monitor.from_netlink(context) monitor.filter_by(subsystem="block", device_type="disk") for device in iter(monitor.poll, None): if device.action == "add" and "ID_FS_TYPE" in device.properties: # Adjust mount path based on your system's USB mounting logic mount_path = f"/media/{device.properties['ID_FS_LABEL']}" requests.post("http://localhost:8000/api/scan/usb", json={"usb_mount_path": mount_path})
- Chunked file reading: Never load entire large files into memory (use 4KB-8KB chunks for SHA256 calculations)
- Parallel processing: Scan multiple files at once with Celery workers or process pools
- Error handling: Skip unreadable files (permission issues, system files) to avoid task stalls
- Progress tracking: Return real-time progress to your React frontend to improve UX
内容的提问来源于stack exchange,提问作者Ahsan

