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. Objetivos
  2. 1. Programación: quién tiene la verdad
  3. 2. El lote que no tumba la BD
  4. 3. Idempotencia por ventana: el lote puede correr dos veces
  5. 4. Solapamiento: locks de corrida
  6. 5. El catálogo de programadas de TicketFlow
  7. Autoevaluación

Stack: Django/DRF · Proyecto: TicketFlow Estado: Publicada — cron moderno y lotes sin pegarte un tiro en el pie Prerrequisito: Lección 30 — Brokers de mensajería


Objetivos

  1. Programar trabajos periódicos con beat/crontab sabiendo quién es dueño de la verdad (el código, no la crontab del servidor).
  2. Escribir procesamiento por lotes con cursor + update_fields + commit por lote: la liquidación mensual sin tumbar la BD.
  3. Manejar solapamientos (una corrida que tarda más que su intervalo) con locks y idempotencia por ventana.

1. Programación: quién tiene la verdad

Tres casas para el cron de TicketFlow: la crontab del sistema (0 3 * /app/venv/bin/manage.py expirar_reservas), celery beat (29), y un scheduler del orquestador (K8s CronJob, 43). La regla: la definición vive en el repo y se despliega con la app — la crontab manual del servidor es la config que no versionas (27): se pierde con la máquina, diverge entre nodos y nadie la revisa en un PR. Beat con beat_schedule en settings/código (27: config del entorno, schedule del código — la FRECUENCIA es diseño, no entorno) o CronJob en el manifiesto del repo (44).

El detalle que muerde: la zona horaria. crontab(hour=3) en beat usa la TZ de settings (TIME_ZONE = "Europe/Madrid"); la crontab del sistema usa la del server (¿UTC? ¿quién sabe?). Para tareas ligadas a negocio ("recordatorio el día del evento, hora local del evento") calcula el disparo en el dominio con el Clock (24): una tarea diaria que escanea "eventos en 24 h" es inmune a la TZ; "corre a las 3" depende de la TZ — elige la primera siempre que puedas.

2. El lote que no tumba la BD

La liquidación mensual de comisiones (org de la 08) toca 2M de filas. El bucle naïf (for org in Organization.objects.all(): calcular(...)) carga todo en memoria y mantiene una transacción de horas. El patrón por lotes:

python
def liquidar_mes(mes: str, *, lote: int = 500) -> int:
    total = 0
    qs = (Organization.objects
          .filter(activa=True)
          .order_by("id")                       # orden ESTABLE: el criterio de corte
          .values_list("id", flat=True))
    ids = list(qs.iterator(chunk_size=lote))    # cursor del servidor, sin cargar 2M
    for i in range(0, len(ids), lote):
        with transaction.atomic():
            orgs = Organization.objects.select_for_update(skip_locked=True).filter(id__in=ids[i:i+lote])
            for org in orgs:
                comision = calcular_comision(org, mes)       # lógica pura, unitaria
                LedgerEntry.objects.create(org=org, mes=mes, amount=comision)
                total += 1
        time.sleep(0.05)                        # deja respirar a la replicación (55)
    return total

Piezas: .iterator(chunk_size) usa cursor del servidor (no materializa el queryset); commit POR LOTE (una transacción de 500 filas es sana; una de 2M bloquea, hincha el WAL y si muere a la 1.9M, pierdes TODO el trabajo); update_fields en los saves ("evita el signal-aviso y el full save"), y el orden estable por id para poder REANUDAR: si murió a mitad, la corrida siguiente filtra id > ultimo_procesado (el checkpoint).

3. Idempotencia por ventana: el lote puede correr dos veces

El lote mensual se solapa (la corrida del día 1 tarda 3 h, el beat del día 2 dispara igual) o se repite (at-least-once de la 30). La protección: clave de idempotencia de la corrida — (org, mes) ÚNICO en LedgerEntry (constraint de la 08):

python
class Meta:
    constraints = [models.UniqueConstraint(fields=["org", "mes"], name="uniq_ledger_org_mes")]

Con get_or_create o ignore_conflicts, la repetición no duplica dinero: la segunda pasada encuentra las filas existentes y salta. Generalización a TODA tarea programada: la "ventana de negocio" (el mes, el día, el minuto) es la clave de idempotencia, no el intento de ejecución. La tarea de recordatorios usa (reserva, kind, fecha_local_del_evento); la de expiración es naturalmente idempotente (la máquina de estados no re-expira, 00b).

4. Solapamiento: locks de corrida

Si la corrida tarda más que su intervalo, dos instancias corren en paralelo y aunque el lote sea idempotente, se disputan los mismos locks (skip_locked hace que se repartan filas, pero el checkpoint/reanudación se corrompe). Lock de corrida con Redis (12) o Postgres advisory locks:

python
def with_run_lock(name: str, ttl: int):
    def deco(fn):
        @wraps(fn)
        def wrapper(*args, **kwargs):
            got = cache.add(f"run-lock:{name}", socket.gethostname(), ttl)
            if not got:
                logger.info("corrida %s ya en marcha: skip", name)
                return None
            try:
                return fn(*args, **kwargs)
            finally:
                cache.delete(f"run-lock:{name}")
        return wrapper
    return deco

@with_run_lock("liquidar_mes", ttl=4 * 3600)
def liquidar_mes(mes: str, **kw): ...

El TTL > duración máxima esperada (si el worker muere, el lock expira solo); el skip con log es explícito (no error: es el comportamiento normal de un intervalo corto). La alternativa Postgres: pg_try_advisory_lock(hashtext('liquidar_mes')) — sobrevive mejor a los relojes raros de contenedores (43) y muere con la conexión (sin TTL que gestionar). Elige uno y documenta.

5. El catálogo de programadas de TicketFlow

TareaFrecuenciaLoteIdempotenciaLock
expirar_reservas (29)1 min200 por pasadamáquina de estadoslock por minuto
recordatorio 24 h antes15 min500(reserva, kind, fecha evento)lock de corrida
cierre de eventos pasadosdiaria1000ventana fecha < hoyadvisory lock
liquidación mensualmensual500UniqueConstraint (org, mes)advisory lock + checkpoint
export RGPD diferido (23)on-demand + diaria de cola200por request_id—

El diseño de una programada sana en una frase: ventana de negocio como clave de idempotencia, lote con commit por trozo, lock de corrida, y la lógica en el servicio (28) con el Clock inyectado — el disparador (beat/cron/CronJob) es un borde intercambiable.


Autoevaluación

  1. ¿Por qué la crontab manual del servidor viola lo aprendido en la 27 y dónde vive la definición en este proyecto?
  2. Enumera las 5 piezas del lote que no tumba la BD y qué rompe cada una si la omites.
  3. ¿Cómo conviertes la ventana de negocio (mes/día/minuto) en clave de idempotencia? Da el constraint concreto del ledger.
  4. Lock de corrida: ¿por qué TTL en Redis y por qué advisory lock de Postgres es la alternativa preferible en contenedores?
  5. "Recordatorio el día del evento, hora local del evento": ¿por qué escanear "eventos en 24 h" cada 15 min vence a "correr a las 3"?

Continúa con los ejercicios. Las solutions.md solo tras intentarlo.