Phase 5: Production & Deployment

Long-running AI tasks with queues & webhooks

Intermediate ~14 min read
Think of it this way A friendly analogy. Read this if the technical version feels dense. Show Hide

Imagine you want to grow a very special, giant sunflower. It's not like just watering a small houseplant; this takes a lot of care and time, maybe days or even weeks! If you just stand there watching the seed, waiting for it to grow, you'd be stuck. You couldn't play, do homework, or anything else. And what if you tried to watch it for a really long time, but then had to leave? You'd miss it, or worse, the seed might get forgotten. Computers face a similar problem with big jobs, especially with smart AI tasks that need to think hard or crunch lots of information. They can't just stop everything and wait for one huge task to finish.

This is where a clever system comes in, like a super-organized garden center. When you want your giant sunflower, you don't stand there watching the seed. Instead, you go to the garden center's counter and tell them what you want. They write down your request on a long list of orders – this list is like a "queue" in computer talk. Then, they immediately give you a little receipt with an order number. You're free to go play! The actual gardeners, who are like the "worker processes" in a computer, will get to your request when they finish the plants before yours on the list. This means the garden center can take lots of orders without getting overwhelmed, and the gardeners can work steadily without anyone bugging them.

But how do you know your sunflower is ready? You don't want to keep walking back to the garden center every hour to ask, "Is it ready yet?" That would be a lot of wasted trips! So, when you place your order, you also tell the garden center your phone number. Once your sunflower is big and beautiful, the gardener sends you a text message or calls you. This is what we call a "webhook" – it’s the garden sending you a message, instead of you constantly checking them.

This whole system makes sure that even when you ask a computer to do something really big, like make a super complex AI drawing or analyze a huge book, it doesn't get stuck or crash. You get your "receipt" right away and can move on, trusting that the computer will do the hard work in the background. Then, when it’s all done, it sends you a notification, so you always get your finished "sunflower" without ever having to wait around or wonder if it's ready. It keeps everything running smoothly and reliably, no matter how long the AI task takes!

The core mental model is a producer/consumer split. Your API server is the producer: it validates the request, serializes the job payload (prompt, model ID, user context, callback URL), and pushes it onto a durable queue. The queue is your buffer and reliability guarantee. Even if every worker crashes, the message sits there until a worker comes back up and picks it up. Workers are the consumers: they pull jobs off the queue, call the AI provider (or run local inference), and write results back to a job store (a Postgres table, Redis hash, or DynamoDB item). Neither side knows or cares about the other's scale or health, which is exactly the point.

Consider a real scenario: a legal-tech SaaS that lets law firms upload contracts and get a structured risk analysis via an LLM. A medium-sized firm uploads 40 contracts at 9 AM Monday. Each analysis takes 15 to 45 seconds. If you handled this synchronously, you'd need 40 concurrent open HTTP connections, your load balancer would likely timeout half of them, and the user would see a broken experience. With a queue, the upload endpoint accepts all 40 files in under a second each, hands back job IDs, and a pool of 8 workers drains the queue over the next few minutes. The law firm's integration code either polls a status endpoint or receives webhook POSTs as each contract finishes. No timeouts, no dropped requests, and you can scale workers independently of the API tier during business hours.

Compared to alternatives: long-polling keeps a connection alive but is only viable under roughly 30 seconds; server-sent events (SSE) and WebSockets work well for streaming token-by-token output but are covered in the streaming lesson and don't fit batch or multi-step agent workloads cleanly. Simple polling by the client works and is easy to implement, but it burns requests and adds latency inversely proportional to your polling interval. Webhooks win when the client can host an endpoint; polling wins as a fallback when they can't (think CLI tools, scripts, or third-party integrations where the caller has no public IP). Designing for both is the professional move: support a callback_url field in your API and also expose a GET /jobs/{id} status endpoint.

At 10 users, a single Celery worker on the same machine as your API is fine. At 10k users with bursty load, you want autoscaling workers: on AWS, that means SQS as the broker, ECS tasks or Lambda as workers, and a CloudWatch alarm on queue depth to trigger scaling. At 10M users, the queue itself becomes a design decision: SQS is managed and cheap, but has a 256KB message size limit (put large payloads in S3 and pass only the S3 key in the message). You also need to think about dead-letter queues (DLQs): messages that fail N times get moved to a DLQ for inspection rather than looping forever. Redis-backed Celery is excellent up to moderate scale but adds operational complexity (Redis Sentinel or Cluster for HA); SQS removes that operational burden at the cost of AWS lock-in.

Webhook delivery has its own reliability surface. Your worker fires a POST to the client's callback URL, but that URL might be down, slow, or returning 500s. Best practice: retry with exponential backoff (2s, 4s, 8s up to N attempts), log every delivery attempt with HTTP status and response body, and sign every payload with an HMAC-SHA256 signature so the receiver can verify it came from you. The client should respond with 200 as fast as possible (queue the payload internally if they need to do more work). On your side, record webhook delivery state per job so you can re-trigger delivery from an admin panel. Cost and latency implications: queuing adds a few hundred milliseconds of overhead in the happy path, which is negligible compared to a 10-second LLM call. The real cost benefit is in retry economics: failed jobs don't lose work, they just wait for a retry, so you're not paying for repeated full-pipeline reruns.

Key Takeaways

  • Return a job ID immediately; never hold HTTP connections open for multi-second AI tasks.
  • Use a durable queue (SQS, Celery/Redis) to decouple API acceptance from worker execution.
  • Webhooks push results to callers; design them with HMAC signatures and idempotency keys.
  • Expose a polling fallback endpoint alongside webhooks for clients that can't receive callbacks.

Pro tips

  • Set task_acks_late = True in Celery so a message is only acknowledged after the task succeeds. Without it, a worker crash mid-LLM-call silently drops the job with no retry.
  • Put large inputs (PDFs, embeddings payloads) in S3 or blob storage and pass only the object key in the queue message. SQS has a 256KB hard limit and even Redis will thank you for it at scale.
  • Always include an idempotency key in your webhook payload. Clients will occasionally receive the same webhook twice due to network retries; a stable job_id in the payload lets them deduplicate without storing state.
  • Track queue depth as your primary scaling signal, not CPU. LLM workers are almost always IO-bound waiting on the provider. If depth climbs, add workers; if time_in_queue spikes, check for provider rate limit errors surfacing as slow retries.

Common pitfalls

  • Mistake: Storing the full result only in Celery's result backend and never persisting it elsewhere. Fix: Write job results to your own database (Postgres, DynamoDB). Celery backends have configurable TTLs and are not a source of truth.
  • Mistake: Firing webhooks with no signature and no retry logic. Fix: Sign every payload with HMAC-SHA256 and retry with exponential backoff; log each delivery attempt with HTTP status for debugging.
  • Mistake: Setting no time_limit on tasks, letting a hung LLM call hold a worker forever. Fix: Set both soft_time_limit and time_limit in Celery (or equivalent in your queue system) so workers are freed and the job re-queued.
  • Mistake: Returning 202 Accepted with a job ID but providing no polling endpoint. Fix: Always expose GET /jobs/{id} alongside webhook support; many clients (CLI tools, third-party scripts) cannot receive inbound HTTP.

When to use webhooks vs polling vs SSE for async AI task results

Option Use when Avoid when
Webhooks (push) Client can host a public HTTPS endpoint and you want minimal latency on completion notification. Client is a CLI tool, script, or browser SPA with no static public URL.
Polling (GET /jobs/{id}) Client cannot receive inbound connections (scripts, mobile apps behind NAT, third-party integrations). Task completes in under 5 seconds; polling overhead becomes noise relative to work done.
Long-polling Task usually finishes in under 30 seconds and you want lower perceived latency without SSE infrastructure. Tasks routinely exceed 30s; connection timeouts at load balancers and proxies become a serious problem.
SSE / WebSocket You want to stream incremental output (tokens, partial results) as the task progresses, not just a final result. You need a simple fire-and-forget model; SSE/WebSocket adds connection management complexity with no benefit.
Webhooks + polling fallback Building a public API used by diverse third-party developers with unknown client architectures. Internal service-to-service integration where you control both sides and can pick one pattern cleanly.

Code Example

python
# celery==5.3, redis==5.0, openai==1.x
from celery import Celery
import openai, os

app = Celery("ai_tasks", broker="redis://localhost:6379/0", backend="redis://localhost:6379/1")

@app.task(bind=True)
def run_llm_task(self, prompt: str, model: str = "gpt-4o-mini") -> dict:
    client = openai.OpenAI(api_key=os.environ["OPENAI_API_KEY"])
    response = client.chat.completions.create(
        model=model,
        messages=[{"role": "user", "content": prompt}],
    )
    return {"text": response.choices[0].message.content, "job_id": self.request.id}

# In your FastAPI endpoint:
# task = run_llm_task.delay(prompt="Summarize this document...")
# return {"job_id": task.id, "status": "pending"}

How this code works

This code defines a background task for processing AI requests, specifically calling an OpenAI model. Its purpose is to offload time-consuming AI operations from a main application (like a web server), allowing the server to respond quickly while the AI work happens in parallel. This prevents an application from freezing or timing out while waiting for a large language model to generate a response.

The setup begins by configuring Celery with redis for both sending tasks (broker) and storing results (backend). The @app.task(bind=True) decorator transforms run_llm_task into a Celery task that can be executed asynchronously. Crucially, bind=True allows the task to access its own metadata via self, specifically self.request.id, which is its unique identifier. Inside this task, an openai.OpenAI client is initialized using an OPENAI_API_KEY from environment variables, enhancing security. The client.chat.completions.create method then sends the prompt to the specified model (defaulting to gpt-4o-mini). Finally, the task returns the AI's text response and its job_id, making it easy for the calling application to track the outcome later. The commented lines illustrate how run_llm_task.delay() initiates this process, immediately returning the job_id to the user, without waiting for the AI result.

Production-grade example

Adds HMAC signing, structured logging, token cost tracking, per-exception retries with backoff, and hard/soft timeouts.

python
# celery==5.3, redis==5.0, openai==1.x, httpx==0.27, structlog==24.x
import os, time, hmac, hashlib, json, logging
import httpx
from celery import Celery
from celery.utils.log import get_task_logger
import openai
import structlog

logger = structlog.get_logger()

app = Celery(
    "ai_tasks",
    broker=os.environ["CELERY_BROKER_URL"],      # e.g. redis://redis:6379/0 or SQS URL
    backend=os.environ["CELERY_RESULT_BACKEND"],
)
app.conf.task_acks_late = True          # ack only after success, so crashes don't lose jobs
app.conf.task_reject_on_worker_lost = True
app.conf.task_max_retries = 5

WEBHOOK_SECRET = os.environ["WEBHOOK_HMAC_SECRET"].encode()

def _sign_payload(payload_bytes: bytes) -> str:
    return hmac.new(WEBHOOK_SECRET, payload_bytes, hashlib.sha256).hexdigest()

@app.task(
    bind=True,
    autoretry_for=(openai.RateLimitError, openai.APITimeoutError),
    retry_backoff=True,
    retry_backoff_max=60,
    max_retries=5,
    time_limit=300,          # hard kill after 5 min
    soft_time_limit=270,     # raises SoftTimeLimitExceeded at 4m30s for graceful cleanup
)
def run_llm_task(self, *, prompt: str, model: str, callback_url: str | None, job_id: str) -> dict:
    log = logger.bind(job_id=job_id, attempt=self.request.retries)
    start = time.monotonic()

    client = openai.OpenAI(api_key=os.environ["OPENAI_API_KEY"], timeout=60.0)
    try:
        response = client.chat.completions.create(
            model=model,
            messages=[{"role": "user", "content": prompt}],
        )
    except openai.OpenAIError as exc:
        log.error("llm_call_failed", error=str(exc))
        raise  # let Celery retry based on autoretry_for

    elapsed = time.monotonic() - start
    usage = response.usage
    log.info(
        "llm_call_succeeded",
        latency_s=round(elapsed, 3),
        prompt_tokens=usage.prompt_tokens,
        completion_tokens=usage.completion_tokens,
        model=model,
    )

    result = {
        "job_id": job_id,
        "status": "complete",
        "text": response.choices[0].message.content,
        "usage": {"prompt": usage.prompt_tokens, "completion": usage.completion_tokens},
    }

    if callback_url:
        _deliver_webhook(callback_url, result, log)

    return result

def _deliver_webhook(url: str, payload: dict, log) -> None:
    body = json.dumps(payload).encode()
    sig = _sign_payload(body)
    for attempt in range(4):
        wait = 2 ** attempt
        try:
            with httpx.Client(timeout=10.0) as client:
                r = client.post(
                    url,
                    content=body,
                    headers={"Content-Type": "application/json", "X-Signature-SHA256": sig},
                )
                r.raise_for_status()
                log.info("webhook_delivered", url=url, status=r.status_code, attempt=attempt)
                return
        except httpx.HTTPError as exc:
            log.warning("webhook_attempt_failed", url=url, attempt=attempt, error=str(exc))
            if attempt < 3:
                time.sleep(wait)
    log.error("webhook_all_attempts_failed", url=url)

How this code works

This code establishes a robust system for performing long-running AI tasks and asynchronously reporting their outcomes. It leverages Celery to manage background jobs, specifically making calls to OpenAI's large language models, and then uses webhooks to notify other services upon task completion. This architecture is vital for executing AI computations that may take significant time without blocking the main application flow, ensuring a responsive user experience.

The Celery application is configured for resilience with task_acks_late and task_reject_on_worker_lost, which prevent job loss during worker crashes, and task_max_retries for general retries. The run_llm_task function, marked with @app.task, is the core AI processing unit. It includes autoretry_for specific OpenAI errors (like rate limits) and retry_backoff to handle transient API issues gracefully, retrying with increasing delays. A subtle but important detail is bind=True, which allows the task to access its own context via self, enabling logging of self.request.retries. After successfully calling openai.OpenAI().chat.completions.create, the result is prepared. If a callback_url is provided, the _deliver_webhook function dispatches the result. This function securely signs the payload using hmac.new and attempts delivery with httpx.Client up to four times, employing an exponential time.sleep backoff to ensure reliable notification even if the receiving service is temporarily unavailable.

Practice & master

Try the exercise, check your understanding, then mark this lesson mastered to track your path to pro.

Exercise

Build a minimal FastAPI app that accepts a summarization request, enqueues it as a Celery task (using a local Redis broker), and returns a job ID. Then add a GET /jobs/{job_id} endpoint that returns the task status and result when complete. Test by submitting a job and polling until it shows 'complete'.

python
# Requirements: fastapi, celery, redis, openai, uvicorn
# Make sure Redis is running locally: docker run -p 6379:6379 redis

from fastapi import FastAPI
from celery import Celery
import openai, os, uuid

app = FastAPI()
celery_app = Celery(broker="redis://localhost:6379/0", backend="redis://localhost:6379/1")

@celery_app.task
def summarize(text: str, job_id: str) -> dict:
    # TODO: call openai chat completions with a summarize prompt
    # TODO: return {"job_id": job_id, "status": "complete", "summary": ...}
    pass

@app.post("/summarize")
def submit_summarize(payload: dict):
    job_id = str(uuid.uuid4())
    # TODO: enqueue the summarize task with payload["text"] and job_id
    # TODO: return {"job_id": job_id, "status": "pending"}
    pass

@app.get("/jobs/{job_id}")
def get_job(job_id: str):
    # TODO: look up task result using celery_app.AsyncResult(job_id)
    # TODO: return status ("pending", "complete", "failed") and result if ready
    pass

Quick check

  1. Why should Celery's task_acks_late be set to True for AI worker tasks?

  2. A client cannot host a public HTTP endpoint. Which result delivery mechanism should you use?

  3. You're passing a 10MB PDF payload through SQS to your AI worker. What is the right approach?

Self-check: Without looking at your notes, describe the full lifecycle of an async AI job: from the client's POST request to webhook delivery. Then explain what happens to the job if the worker crashes halfway through, and how you'd configure your queue to handle that correctly.