Python Background Jobs

作者 wshobson46891e7e60da無授權條款收錄於 2026年10月8日更新於 2026年10月8日

Python background job patterns including task queues, workers, and event-driven architecture. Use when implementing async task processing, job queues, long-running operations, or decoupling work from request/response cycles.

AI 產生的概覽

指導 Python 背景工作設計,涵蓋任務佇列、工作者、重試、幂等性與任務狀態追蹤。

功能
此技能提供模式與程式碼範例,用於在 Python 應用程式中將長時間執行的工作與請求/回應週期解耦。內容涵蓋使用 Celery 的任務佇列、立即回傳任務 ID、設定重試與逾時、讓任務具備幂等性,以及持久化任務狀態轉換。它也指向一個參考檔案以取得進階模式。
適用情境
適用於在 Python 中實作非同步任務處理、任務佇列或長時間執行操作時。適合寄送電子郵件或 webhook、產生報表、處理上傳檔案,以及與不可靠外部服務整合等情境。
執行需求
不隨附指令碼,僅為說明性內容。範例假定使用 Python、Celery 以及 Redis 等訊息代理,並需要資料庫保存任務狀態,但這些只是示例而非隨附相依項目。

Python Background Jobs & Task Queues

Decouple long-running or unreliable work from request/response cycles. Return immediately to the user while background workers handle the heavy lifting asynchronously.

When to Use This Skill

  • Processing tasks that take longer than a few seconds
  • Sending emails, notifications, or webhooks
  • Generating reports or exporting data
  • Processing uploads or media transformations
  • Integrating with unreliable external services
  • Building event-driven architectures

Core Concepts

1. Task Queue Pattern

API accepts request, enqueues a job, returns immediately with a job ID. Workers process jobs asynchronously.

2. Idempotency

Tasks may be retried on failure. Design for safe re-execution.

3. Job State Machine

Jobs transition through states: pending → running → succeeded/failed.

4. At-Least-Once Delivery

Most queues guarantee at-least-once delivery. Your code must handle duplicates.

Quick Start

This skill uses Celery for examples, a widely adopted task queue. Alternatives like RQ, Dramatiq, and cloud-native solutions (AWS SQS, GCP Tasks) are equally valid choices.

python
from celery import Celery
app = Celery("tasks", broker="redis://localhost:6379")
@app.taskdef send_email(to: str, subject: str, body: str) -> None:    # This runs in a background worker    email_client.send(to, subject, body)
# In your API handlersend_email.delay("[email protected]", "Welcome!", "Thanks for signing up")

Fundamental Patterns

Pattern 1: Return Job ID Immediately

For operations exceeding a few seconds, return a job ID and process asynchronously.

python
from uuid import uuid4from dataclasses import dataclassfrom enum import Enumfrom datetime import datetime
class JobStatus(Enum):    PENDING = "pending"    RUNNING = "running"    SUCCEEDED = "succeeded"    FAILED = "failed"
@dataclassclass Job:    id: str    status: JobStatus    created_at: datetime    started_at: datetime | None = None    completed_at: datetime | None = None    result: dict | None = None    error: str | None = None
# API endpointasync def start_export(request: ExportRequest) -> JobResponse:    """Start export job and return job ID."""    job_id = str(uuid4())
    # Persist job record    await jobs_repo.create(Job(        id=job_id,        status=JobStatus.PENDING,        created_at=datetime.utcnow(),    ))
    # Enqueue task for background processing    await task_queue.enqueue(        "export_data",        job_id=job_id,        params=request.model_dump(),    )
    # Return immediately with job ID    return JobResponse(        job_id=job_id,        status="pending",        poll_url=f"/jobs/{job_id}",    )

Pattern 2: Celery Task Configuration

Configure Celery tasks with proper retry and timeout settings.

python
from celery import Celery
app = Celery("tasks", broker="redis://localhost:6379")
# Global configurationapp.conf.update(    task_time_limit=3600,          # Hard limit: 1 hour    task_soft_time_limit=3000,      # Soft limit: 50 minutes    task_acks_late=True,            # Acknowledge after completion    task_reject_on_worker_lost=True,    worker_prefetch_multiplier=1,   # Don't prefetch too many tasks)
@app.task(    bind=True,    max_retries=3,    default_retry_delay=60,    autoretry_for=(ConnectionError, TimeoutError),)def process_payment(self, payment_id: str) -> dict:    """Process payment with automatic retry on transient errors."""    try:        result = payment_gateway.charge(payment_id)        return {"status": "success", "transaction_id": result.id}    except PaymentDeclinedError as e:        # Don't retry permanent failures        return {"status": "declined", "reason": str(e)}    except TransientError as e:        # Retry with exponential backoff        raise self.retry(exc=e, countdown=2 ** self.request.retries * 60)

Pattern 3: Make Tasks Idempotent

Workers may retry on crash or timeout. Design for safe re-execution.

python
@app.task(bind=True)def process_order(self, order_id: str) -> None:    """Process order idempotently."""    order = orders_repo.get(order_id)
    # Already processed? Return early    if order.status == OrderStatus.COMPLETED:        logger.info("Order already processed", order_id=order_id)        return
    # Already in progress? Check if we should continue    if order.status == OrderStatus.PROCESSING:        # Use idempotency key to avoid double-charging        pass
    # Process with idempotency key    result = payment_provider.charge(        amount=order.total,        idempotency_key=f"order-{order_id}",  # Critical!    )
    orders_repo.update(order_id, status=OrderStatus.COMPLETED)

Idempotency Strategies:

  1. Check-before-write: Verify state before action
  2. Idempotency keys: Use unique tokens with external services
  3. Upsert patterns: INSERT ... ON CONFLICT UPDATE
  4. Deduplication window: Track processed IDs for N hours

Pattern 4: Job State Management

Persist job state transitions for visibility and debugging.

python
class JobRepository:    """Repository for managing job state."""
    async def create(self, job: Job) -> Job:        """Create new job record."""        await self._db.execute(            """INSERT INTO jobs (id, status, created_at)               VALUES ($1, $2, $3)""",            job.id, job.status.value, job.created_at,        )        return job
    async def update_status(        self,        job_id: str,        status: JobStatus,        **fields,    ) -> None:        """Update job status with timestamp."""        updates = {"status": status.value, **fields}
        if status == JobStatus.RUNNING:            updates["started_at"] = datetime.utcnow()        elif status in (JobStatus.SUCCEEDED, JobStatus.FAILED):            updates["completed_at"] = datetime.utcnow()
        await self._db.execute(            "UPDATE jobs SET status = $1, ... WHERE id = $2",            updates, job_id,        )
        logger.info(            "Job status updated",            job_id=job_id,            status=status.value,        )

Detailed worked examples and patterns

Detailed sections (starting with ## Advanced Patterns) live in references/details.md. Read that file when the navigation summary above is insufficient.

Best Practices Summary

  1. Return immediately - Don't block requests for long operations
  2. Persist job state - Enable status polling and debugging
  3. Make tasks idempotent - Safe to retry on any failure
  4. Use idempotency keys - For external service calls
  5. Set timeouts - Both soft and hard limits
  6. Implement DLQ - Capture permanently failed tasks
  7. Log transitions - Track job state changes
  8. Retry appropriately - Exponential backoff for transient errors
  9. Don't retry permanent failures - Validation errors, invalid credentials
  10. Monitor queue depth - Alert on backlog growth

來源與署名

來源:wshobson/agents位於plugins/python-development/skills/python-background-jobs提交46891e7

授權條款: 無授權條款

內容歸原作者所有。SourceWeft 從公開儲存庫中收錄這些內容。

檢舉或申請下架