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

如何在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=True flag 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:

  1. Install dependencies:

    pip install rq redis
    
  2. Define 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)
    
  3. 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"}
    
  4. 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:

  1. Install Celery and a broker (e.g., Redis):

    pip install celery redis
    
  2. Configure 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
        pass
    
  3. Start the worker with concurrency=1:

    celery -A celery_app worker --concurrency=1 --loglevel=info
    

Key Considerations for Production

  • Task Persistence: Built-in queue.Queue is 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 queue module 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 09:28:26