Module 6 · Async processing and messaging

Lesson 31 — Scheduled tasks and batches

Modern cron and batch processing without shooting yourself in the foot.

Published
In this lesson
  1. Objectives
  2. 1. Scheduling: who holds the truth
  3. 2. The batch that doesn't knock the DB over
  4. 3. Idempotency per window: the batch may run twice
  5. 4. Overlap: run locks
  6. 5. TicketFlow's scheduled catalog
  7. Self-assessment

Stack: Django/DRF · Project: TicketFlow Status: Published — modern cron and batches without shooting yourself in the foot Prerequisite: Lesson 30 — Message brokers


Objectives

  1. Schedule periodic jobs with beat/crontab knowing who owns the truth (the code, not the server's crontab).
  2. Write batch processing with cursor + update_fields + commit per batch: the monthly settlement without knocking the DB over.
  3. Handle overlaps (a run taking longer than its interval) with locks and per-window idempotency.

1. Scheduling: who holds the truth

Three homes for TicketFlow's cron: the system crontab (0 3 * /app/venv/bin/manage.py expirar_reservas), celery beat (29), and an orchestrator scheduler (K8s CronJob, 43). The rule: the definition lives in the repo and deploys with the app — the server's manual crontab is the config you don't version (27): it is lost with the machine, diverges between nodes and nobody reviews it in a PR. Beat with beat_schedule in settings/code (27: config from the environment, schedule from code — the FREQUENCY is design, not environment) or a CronJob in the repo's manifest (44).

The detail that bites: the timezone. crontab(hour=3) in beat uses settings' TZ (TIME_ZONE = "Europe/Madrid"); the system crontab uses the server's (UTC? who knows?). For business-bound tasks ("reminder on the event's day, the event's local time") compute the firing in the domain with the Clock (24): a daily task that scans "events within 24 h" is TZ-immune; "runs at 3" depends on the TZ — pick the former whenever you can.

2. The batch that doesn't knock the DB over

The monthly commission settlement (08's org) touches 2M rows. The naïve loop (for org in Organization.objects.all(): calcular(...)) loads everything into memory and holds a transaction for hours. The batch pattern:

python
def liquidar_mes(mes: str, *, lote: int = 500) -> int:
    total = 0
    qs = (Organization.objects
          .filter(activa=True)
          .order_by("id")                       # STABLE order: the cutoff criterion
          .values_list("id", flat=True))
    ids = list(qs.iterator(chunk_size=lote))    # server-side cursor, no 2M loaded
    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)       # pure logic, unit-testable
                LedgerEntry.objects.create(org=org, mes=mes, amount=comision)
                total += 1
        time.sleep(0.05)                        # let replication breathe (55)
    return total

Pieces: .iterator(chunk_size) uses a server-side cursor (doesn't materialize the queryset); commit PER BATCH (a 500-row transaction is healthy; a 2M one blocks, bloats the WAL and if it dies at 1.9M you lose ALL the work); update_fields on saves ("avoids the signal-notice and the full save"), and the stable order by id so you can RESUME: if it died halfway, the next run filters id > ultimo_procesado (the checkpoint).

3. Idempotency per window: the batch may run twice

The monthly batch overlaps (the day-1 run takes 3 h, the day-2 beat fires anyway) or repeats (30's at-least-once). The protection: the run's idempotency key — (org, mes) UNIQUE in LedgerEntry (08's constraint):

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

With get_or_create or ignore_conflicts, the repetition doesn't duplicate money: the second pass finds the existing rows and skips. Generalization to EVERY scheduled task: the "business window" (the month, the day, the minute) is the idempotency key, not the execution attempt. The reminders task uses (reserva, kind, fecha_local_del_evento); the expiration one is naturally idempotent (the state machine doesn't re-expire, 00b).

4. Overlap: run locks

If the run takes longer than its interval, two instances run in parallel and even if the batch is idempotent, they fight over the same locks (skip_locked makes them split rows, but the checkpoint/resume gets corrupted). Run lock with Redis (12) or 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("run %s already going: 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): ...

The TTL > maximum expected duration (if the worker dies, the lock expires on its own); the skip with a log is explicit (not an error: it is the normal behavior of a short interval). The Postgres alternative: pg_try_advisory_lock(hashtext('liquidar_mes')) — it survives containers' weird clocks better (43) and dies with the connection (no TTL to manage). Pick one and document it.

5. TicketFlow's scheduled catalog

TaskFrequencyBatchIdempotencyLock
expirar_reservas (29)1 min200 per passstate machineper-minute lock
24h-before reminder15 min500(reservation, kind, event date)run lock
closing past eventsdaily1000window date < todayadvisory lock
monthly settlementmonthly500UniqueConstraint (org, month)advisory lock + checkpoint
deferred GDPR export (23)on-demand + daily queue drain200per request_id—

The healthy scheduled task's design in one sentence: business window as the idempotency key, batch with commit per chunk, run lock, and the logic in the service (28) with the Clock injected — the trigger (beat/cron/CronJob) is an interchangeable edge.


Self-assessment

  1. Why does the server's manual crontab violate what you learned in 27, and where does the definition live in this project?
  2. List the 5 pieces of the batch that doesn't knock the DB over and what each one's omission breaks.
  3. How do you turn the business window (month/day/minute) into the idempotency key? Give the ledger's concrete constraint.
  4. Run lock: why a TTL in Redis, and why is Postgres's advisory lock the preferable alternative in containers?
  5. "Reminder on the event's day, the event's local time": why does scanning "events within 24 h" beat "run at 3"?

Continue with the exercises. The solutions only after trying it yourself.