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

基于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.

1. Production-Grade Choice: Celery + Message Queue (Redis/RabbitMQ)

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

  1. Start Redis (required for Celery broker/backend)
  2. Start Celery worker:
    celery -A celery_app worker --loglevel=info --concurrency=4
    
    (Adjust concurrency to match your CPU core count)
  3. Start FastAPI:
    uvicorn main:app --reload
    
2. Lightweight Alternative: FastAPI + ProcessPoolExecutor

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()
3. Real-Time USB Scan Triggering

To detect USB insertion and trigger scans, build a separate monitoring service (runs alongside your API):

  • Linux: Use pyudev to listen for udev block device events, then call the /api/scan/usb endpoint
  • Windows: Use pywin32 to listen for WM_DEVICECHANGE messages and get the USB drive letter
  • macOS: Use pyobjc to 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})
Key Optimizations for Large Files
  • 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 06:39:58