Módulo 6 · Procesamiento asíncrono y mensajería

Lección 31 — Tareas programadas y lotes

Cron moderno y procesamiento por lotes sin pegarte un tiro en el pie.

Publicada
En esta lección
  1. Ejercicio 2 — El lote
  2. Ejercicio 3 — Idempotencia de ventana
  3. Ejercicio 4 — Locks
  4. Ejercicio 5 — El cierre
  5. Resumen del profesor
python
# core/celery.py — el schedule vive en el repo (revisado en PR, desplegado con la app)
app.conf.beat_schedule = {
    # ventana: minuto de negocio; idempotencia: máquina de estados + lock por minuto (29)
    "expirar-reservas": {"task": "reservations.tasks.expirar_reservas", "schedule": crontab()},
    # ventana: escaneo de "eventos en 24h"; idempotencia: (reserva, kind, fecha_evento)
    "recordatorios": {"task": "tickets.tasks.enviar_recordatorios", "schedule": crontab(minute="*/15")},
    # ventana: fecha < hoy; idempotencia: transición CLOSED no re-aplica
    "cierre-eventos": {"task": "events.tasks.cerrar_eventos_pasados", "schedule": crontab(hour=4, minute=30)},
    # ventana: mes; idempotencia: UniqueConstraint (org, mes)
    "liquidacion": {"task": "billing.tasks.liquidar_mes", "schedule": crontab(day_of_month=2, hour=5)},
}
  1. Los hallazgos típicos: un CRON en un Dockerfile viejo que ejecuta backups DOBLES (el del cron + el del orquestador) y una crontab de expiración en el servidor heredada de la 00 que ya nadie reconoce — exactamente la config invisible que la 27 prohíbe. Migración: al repo, con PR.
  1. Las TZ-dependientes son recordatorios y liquidación. Recordatorio: TZ-inmune por diseño (escanea "eventos cuyo starts_at local cae en las próximas 24 h" cada 15 min — la hora de corrida da igual). Liquidación: ancla la ventana al MES (no a la hora) y corre "el día 2 a las 5" con la TZ del negocio documentada en el schedule; el mes como clave hace que una corrida repetida no duplique.

Ejercicio 2 — El lote

1-2. El lote con checkpoint:

python
def liquidar_mes(mes: str, *, lote: int = 500, desde_id: int = 0) -> int:
    total, ultimo = 0, desde_id
    while True:
        with transaction.atomic():
            ids = list(Organization.objects.filter(activa=True, id__gt=ultimo)
                       .order_by("id").values_list("id", flat=True)[:lote])
            if not ids:
                break
            for org in Organization.objects.select_for_update(skip_locked=True).filter(id__in=ids):
                LedgerEntry.objects.get_or_create(org=org, mes=mes, defaults={"amount": calcular_comision(org, mes)})
                ultimo = org.id
            total += len(ids)
        time.sleep(0.05)
    return total

La reanudación es gratis con el get_or_create + id > ultimo: muerto en la #843, relanza con desde_id=842 (o escanea desde 0 y la idempotencia hace el resto). El test de muerte a medias: excepción en la 3ª iteración → re-lanzada → las entradas de los lotes 1-2 ya existen (get_or_create no duplica) y el total final coincide con la corrida limpia.

  1. Queries por lote: 2 (una de ids + una de orgs id__in + N de get_or_create). Sin iterator/values_list con 2M filas: el queryset materializa 2M objetos en memoria (~GBs) y la primera query devuelve TODO en un golpe — el cursor del servidor (iterator(chunk_size)) mantiene el buffer del cliente en chunk_size. El CaptureQueriesContext del test documenta el número como contrato de rendimiento.

Ejercicio 3 — Idempotencia de ventana

  1. La migración expand-contract (11) sobre un ledger con datos:
python
# 1) expand: constraint añadida como NOT VALID (no escanea la tabla entera)
migrations.RunSQL("ALTER TABLE billing_ledgerentry ADD CONSTRAINT uniq_ledger_org_mes UNIQUE (org_id, mes) NOT VALID;")
# 2) limpiar duplicados si la validación falla:
#    DELETE ... WHERE id NOT IN (SELECT MIN(id) FROM ... GROUP BY org_id, mes)
migrations.RunSQL("ALTER TABLE billing_ledgerentry VALIDATE CONSTRAINT uniq_ledger_org_mes;")

NOT VALID + VALIDATE = el constraint se aplica a filas nuevas sin lock largo sobre 2M filas: la forma adulta de añadir unicidad en producción.

  1. El recordatorio:
python
def enviar_recordatorios(clock: Clock):
    mañana = [e for e in Event.objects.filter(starts_at__range=(clock.now(), clock.now() + timedelta(hours=24)))]
    for r in Reservation.objects.filter(event__in=mañana, status="CONFIRMED"):
        clave = f"reminder:{r.public_ref}:24h:{r.event.starts_at.date()}"
        if cache.add(clave, "1", 60 * 60 * 48):
            transaction.on_commit(lambda r=r: email_recordatorio.delay(r.public_ref))

4 corridas el mismo día → la clave ya existe en las 2-4 → UN email. El TTL 48 h cubre el día del evento y algo más.

Ejercicio 4 — Locks

  1. El lock con skip explícito:
python
def with_run_lock(name, ttl):
    def deco(fn):
        @wraps(fn)
        def wrapper(*args, **kwargs):
            if not cache.add(f"run-lock:{name}", socket.gethostname(), ttl):
                logger.info("run-lock %s ocupado: skip", name)
                return None
            try:
                return fn(*args, **kwargs)
            finally:
                cache.delete(f"run-lock:{name}")
        return wrapper
    return deco

El test: cache.add mockeado a False → retorno None, cero llamadas a la lógica, log con "skip" (no excepción: el solape es comportamiento esperado, no error). El ADR (extracto): "Locks de corrida: Postgres advisory. Razón: mueren con la conexión (sin TTL que ajustar), inmunes al drift de reloj de contenedores (43), auditables con pg_locks. Alternativa Redis descartada: TTL mal ajustado = doble corrida". En un despliegue K8s (43), advisory gana; si los schedulers fuesen multi-BD, Redis.

2-3. Con TTL 4 h y corrida de 5 h: en la hora 4 el lock EXPIRA y la corrida del intervalo siguiente entra en paralelo — dos procesos en el mismo lote. Con skip_locked se reparten filas sin excepciones (el peor tipo de bug: todo "funciona"), pero el checkpoint avanza dos veces y get_or_create salva el dinero — la idempotencia de ventana es la red de seguridad que convierte el desastre en no-evento. El fix: heartbeat — un hilo/threading.Timer que renueva el TTL cada 60 s mientras la corrida vive, y el finally cancela el timer y borra el lock.

Ejercicio 5 — El cierre

python
def cerrar_eventos_pasados(clock: Clock, *, lote: int = 1000) -> int:
    total = 0
    while True:
        with transaction.atomic():
            evs = list(Event.objects.select_for_update(skip_locked=True)
                       .filter(status=Event.Status.OPEN, ends_at__lt=clock.now())[:lote])
            if not evs:
                break
            for e in evs:
                e.close()                     # máquina de estados: OPEN → CLOSED (00b)
                outbox.create(event_type="EventClosed", payload={"event_id": e.uuid})
                for r in e.reservations.filter(status=Reservation.Status.ACTIVE):
                    r.expire()                # huérfanos: ACTIVE en evento cerrado
                    outbox.create(event_type="ReservationExpired", payload={"ref": r.public_ref})
            total += len(evs)
    return total

El test integral: 2ª corrida → la query status=OPEN AND ends_at < now devuelve vacío (los eventos ya están CLOSED): cero cambios, cero eventos nuevos. La idempotencia por máquina de estados + ventana es el patrón que hace segura la corrida repetida — cada tarea programada de TicketFlow termina siendo "query de ventana + lote + lock + outbox": cuatro piezas y varía el dominio.


Resumen del profesor

  • El schedule vive en el repo; la ventana de negocio (mes/día/minuto) es la clave de idempotencia, no el intento de ejecución.
  • Lote = cursor + commit por trozo + checkpoint + get_or_create: muere, reanuda y no duplica.
  • Lock de corrida con skip amable (no error); advisory lock cuando el reloj del entorno es sospechoso; heartbeat si la corrida puede superar el TTL.