如何在Python中实现单任务串行的文件处理工作队列?
Hey there! Let's break down how to implement a single-worker task queue for your file-processing Python app—perfect for handling concurrent API submissions while ensuring only one file is processed at a time. I'll cover a few approaches, from lightweight built-in tools to production-ready third-party libraries:
Lightweight Built-in Solution: queue.Queue + Background Thread
If you want to avoid external dependencies and keep things simple, Python's built-in queue.Queue (thread-safe out of the box) paired with a background thread works great. Here's a quick example using FastAPI for your API layer:
from fastapi import FastAPI, UploadFile, File import queue import threading import time from typing import Tuple # Initialize a thread-safe queue task_queue = queue.Queue() app = FastAPI() # Simulate your file processing and DB storage logic def process_file_task(file: UploadFile, metadata: dict): print(f"Starting processing for file: {file.filename}") # Replace with your actual file processing code time.sleep(5) # Simulate long-running process # Extract conclusion and save to DB conclusion = f"Processed {file.filename}: sample conclusion" print(f"Finished processing {file.filename}, saved conclusion to DB") return conclusion # Background worker thread that runs indefinitely def worker(): while True: # Block until a task is available task = task_queue.get() try: file, metadata = task process_file_task(file, metadata) except Exception as e: print(f"Error processing task: {str(e)}") finally: # Mark task as done (important for queue.join()) task_queue.task_done() # Start the worker thread when the app launches threading.Thread(target=worker, daemon=True).start() # API endpoint to submit files @app.post("/submit-file") async def submit_file(file: UploadFile = File(...), user_id: str = None): # Package file and related metadata into a task task = (file, {"user_id": user_id}) task_queue.put(task) return {"status": "success", "message": f"File {file.filename} added to queue. Position: {task_queue.qsize()}"}
Key notes here:
- The
daemon=Trueflag ensures the worker thread exits when the main app does. task_queue.get()blocks until a task is available, so the worker won't waste resources polling.- Always wrap processing in a try/except to handle failures without crashing the worker.
Production-Grade Solution: RQ (Redis Queue)
For a more robust setup (especially if you need task persistence, retries, or distributed workers later), RQ (Redis Queue) is a fantastic choice. It's lightweight, easy to set up, and natively supports limiting worker concurrency.
Steps:
Install dependencies:
pip install rq redisDefine your task and queue setup (e.g.,
tasks.py):import redis from rq import Queue from my_file_processor import process_file_task # Your actual processing function # Connect to Redis (default local setup; adjust for production) redis_conn = redis.Redis() q = Queue(connection=redis_conn)Update your API to enqueue tasks:
from fastapi import FastAPI, UploadFile, File from tasks import q, process_file_task app = FastAPI() @app.post("/submit-file") async def submit_file(file: UploadFile = File(...), user_id: str = None): # Enqueue the task (RQ handles serializing arguments) job = q.enqueue(process_file_task, file, {"user_id": user_id}) return {"status": "success", "job_id": job.id, "message": "File added to queue"}Start a single worker to process tasks one at a time:
rq worker --workers 1
RQ gives you extra perks like:
- Task persistence (tasks survive app restarts)
- Job status tracking (check if a task is pending, running, or finished)
- Retry logic for failed tasks
Alternative: Celery (For Complex Workflows)
If you anticipate growing into more complex workflows (like scheduled tasks, multiple queue types), Celery is the industry standard. To enforce single-task processing, just set worker concurrency to 1:
Install Celery and a broker (e.g., Redis):
pip install celery redisConfigure Celery (e.g.,
celery_app.py):from celery import Celery app = Celery('file_processor', broker='redis://localhost:6379/0') @app.task def process_file_task(file, metadata): # Your processing logic here passStart the worker with concurrency=1:
celery -A celery_app worker --concurrency=1 --loglevel=info
Key Considerations for Production
- Task Persistence: Built-in
queue.Queueis in-memory, so tasks are lost if the app restarts. Use RQ/Celery with a persistent broker (Redis/RabbitMQ) for production. - Error Handling: Add retry logic for transient failures (e.g., DB connection issues) and dead-letter queues for tasks that repeatedly fail.
- Status Tracking: Let users query the status of their submitted tasks (RQ/Celery have built-in methods for this).
- Thread Safety: Ensure your DB operations are thread-safe (use connection pools instead of shared connections).
Reference Resources
- Python's official
queuemodule docs: Covers thread-safe queue behavior and best practices. - RQ official docs: Details on worker configuration, job management, and advanced features.
- Celery official docs: Guides on setting up single-worker queues, task routing, and error handling.
内容的提问来源于stack exchange,提问作者Pau Muñoz

