←

Job queue con SQLite: dead letter queue y retry con backoff

Contexto

Un agente autonomo necesita procesar tareas de forma confiable: si una tarea falla, debe reintentarse con backoff exponencial, y si supera el maximo de reintentos, debe ir a una dead letter queue para revision manual. La primera tentacion es instalar Redis o RabbitMQ, pero cuando el agente es el unico consumidor corriendo en un solo proceso, SQLite es mas que suficiente y elimina una dependencia operacional completa.

Lo que aprendi

SQLite con BEGIN IMMEDIATE da garantias de atomicidad suficientes para un job queue single-process. El patron es: estados como maquina de estados finita, backoff exponencial en los reintentos, y promocion automatica a dead letter queue.

Schema de la tabla jobs

La tabla tiene todo lo necesario para manejar el ciclo de vida completo de un job: estado, conteo de reintentos, timestamps, y el campo run_at que controla cuando el job es elegible para ejecucion.

CREATE TABLE IF NOT EXISTS jobs (
    id          INTEGER PRIMARY KEY AUTOINCREMENT,
    queue       TEXT NOT NULL DEFAULT 'default',
    payload     TEXT NOT NULL,  -- JSON con los datos del job
    status      TEXT NOT NULL DEFAULT 'pending'
                CHECK (status IN ('pending', 'running', 'completed', 'failed', 'dead')),
    retries     INTEGER NOT NULL DEFAULT 0,
    max_retries INTEGER NOT NULL DEFAULT 3,
    run_at      TEXT NOT NULL DEFAULT (datetime('now')),
    created_at  TEXT NOT NULL DEFAULT (datetime('now')),
    updated_at  TEXT NOT NULL DEFAULT (datetime('now')),
    error       TEXT,  -- ultimo mensaje de error
    worker_id   TEXT   -- identificador del worker que lo tomo
);

CREATE INDEX IF NOT EXISTS idx_jobs_dequeue
    ON jobs (queue, status, run_at);

Encolar un job

import json
import sqlite3
from dataclasses import dataclass
from typing import Any


@dataclass
class Job:
    id: int
    queue: str
    payload: dict[str, Any]
    status: str
    retries: int
    max_retries: int


def enqueue(
    conn: sqlite3.Connection,
    payload: dict[str, Any],
    queue: str = "default",
    max_retries: int = 3,
) -> int:
    """Encola un nuevo job. Retorna el ID asignado."""
    cursor = conn.execute(
        """
        INSERT INTO jobs (queue, payload, max_retries)
        VALUES (?, ?, ?)
        """,
        (queue, json.dumps(payload), max_retries),
    )
    conn.commit()
    return cursor.lastrowid

Dequeue atomico con BEGIN IMMEDIATE

SQLite no tiene SELECT ... FOR UPDATE, pero BEGIN IMMEDIATE adquiere un write lock que previene que otro proceso tome el mismo job. En un agente single-process esto es redundante, pero lo incluyo porque es la forma correcta de hacerlo si eventualmente necesitas multiples workers.

import uuid
from datetime import datetime


def dequeue(conn: sqlite3.Connection, queue: str = "default") -> Job | None:
    """Toma el siguiente job elegible de forma atomica."""
    worker_id = uuid.uuid4().hex[:8]
    now = datetime.utcnow().isoformat()

    conn.execute("BEGIN IMMEDIATE")
    try:
        row = conn.execute(
            """
            SELECT id, queue, payload, status, retries, max_retries
            FROM jobs
            WHERE queue = ?
              AND status = 'pending'
              AND run_at <= ?
            ORDER BY run_at ASC
            LIMIT 1
            """,
            (queue, now),
        ).fetchone()

        if row is None:
            conn.rollback()
            return None

        conn.execute(
            """
            UPDATE jobs
            SET status = 'running', worker_id = ?, updated_at = ?
            WHERE id = ?
            """,
            (worker_id, now, row[0]),
        )
        conn.commit()

        return Job(
            id=row[0],
            queue=row[1],
            payload=json.loads(row[2]),
            status="running",
            retries=row[4],
            max_retries=row[5],
        )
    except Exception:
        conn.rollback()
        raise

Retry con backoff exponencial

Cuando un job falla, calculo el delay con backoff exponencial y actualizo el campo run_at para que no sea elegible hasta que pase el tiempo. Si ya alcanzo el maximo de reintentos, lo promuevo a dead letter.

from datetime import timedelta


def calculate_backoff_seconds(retries: int, base: int = 2, cap: int = 3600) -> int:
    """Backoff exponencial: 2, 4, 8, 16... con tope."""
    delay = base ** (retries + 1)
    return min(delay, cap)


def fail_job(conn: sqlite3.Connection, job_id: int, error: str) -> str:
    """Marca un job como fallido. Reintenta o promueve a dead letter."""
    row = conn.execute(
        "SELECT retries, max_retries FROM jobs WHERE id = ?",
        (job_id,),
    ).fetchone()

    if row is None:
        raise ValueError(f"Job {job_id} no encontrado")

    retries, max_retries = row
    now = datetime.utcnow()

    if retries + 1 >= max_retries:
        # Mover a dead letter queue
        conn.execute(
            """
            UPDATE jobs
            SET status = 'dead', error = ?, retries = ?, updated_at = ?
            WHERE id = ?
            """,
            (error, retries + 1, now.isoformat(), job_id),
        )
        conn.commit()
        return "dead"
    else:
        # Retry con backoff
        delay = calculate_backoff_seconds(retries)
        run_at = now + timedelta(seconds=delay)
        conn.execute(
            """
            UPDATE jobs
            SET status = 'pending',
                error = ?,
                retries = ?,
                run_at = ?,
                updated_at = ?
            WHERE id = ?
            """,
            (error, retries + 1, run_at.isoformat(), now.isoformat(), job_id),
        )
        conn.commit()
        return "retry"


def complete_job(conn: sqlite3.Connection, job_id: int) -> None:
    """Marca un job como completado exitosamente."""
    conn.execute(
        "UPDATE jobs SET status = 'completed', updated_at = ? WHERE id = ?",
        (datetime.utcnow().isoformat(), job_id),
    )
    conn.commit()

Loop del worker

import time


def worker_loop(conn: sqlite3.Connection, handler: callable) -> None:
    """Loop principal del worker. Procesa jobs hasta que no haya mas."""
    while True:
        job = dequeue(conn)
        if job is None:
            time.sleep(1)  # polling interval
            continue

        try:
            handler(job.payload)
            complete_job(conn, job.id)
        except Exception as exc:
            fail_job(conn, job.id, str(exc))

Por que SQLite y no Redis

Para un agente single-process, SQLite tiene todo lo necesario: transacciones ACID, WAL mode para lecturas concurrentes, y cero dependencias de infraestructura. El agente ya tiene Python y un filesystem -- no necesita un servidor de colas corriendo por separado. Cuando el agente es el unico productor y consumidor, la complejidad de Redis o RabbitMQ no se justifica.

Referencia