Exercise 1 — The catalog
# core/celery.py — the schedule lives in the repo (reviewed in PRs, deployed with the app)
app.conf.beat_schedule = {
# window: business minute; idempotency: state machine + per-minute lock (29)
"expirar-reservas": {"task": "reservations.tasks.expirar_reservas", "schedule": crontab()},
# window: scan of "events within 24h"; idempotency: (reservation, kind, event_date)
"recordatorios": {"task": "tickets.tasks.enviar_recordatorios", "schedule": crontab(minute="*/15")},
# window: date < today; idempotency: the CLOSED transition doesn't re-apply
"cierre-eventos": {"task": "events.tasks.cerrar_eventos_pasados", "schedule": crontab(hour=4, minute=30)},
# window: month; idempotency: UniqueConstraint (org, month)
"liquidacion": {"task": "billing.tasks.liquidar_mes", "schedule": crontab(day_of_month=2, hour=5)},
}- The typical findings: a
CRONin an old Dockerfile running DOUBLE backups (the cron's one + the orchestrator's) and an expiration crontab on the server inherited from 00 that nobody recognizes anymore — exactly the invisible config 27 forbids. Migration: into the repo, with a PR.
- The TZ-dependent ones are reminders and settlement. Reminder: TZ-immune by design (scans "events whose local
starts_atfalls within the next 24 h" every 15 min — the run's hour doesn't matter). Settlement: anchor the window to the MONTH (not the hour) and run "on day 2 at 5" with the business's TZ documented in the schedule; the month as key makes a repeated run not duplicate.
Exercise 2 — The batch
1-2. The batch with checkpoint:
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 totalResuming is free with get_or_create + id > ultimo: dead at #843, relaunch with desde_id=842 (or scan from 0 and idempotency does the rest). The die-halfway test: exception on the 3rd iteration → relaunched → the batches 1-2 entries already exist (get_or_create doesn't duplicate) and the final total matches the clean run.
- Queries per batch: 2 (one for ids + one for the
id__inorgs + N get_or_creates). Withoutiterator/values_listwith 2M rows: the queryset materializes 2M objects in memory (~GBs) and the first query returns EVERYTHING in one blow — the server-side cursor (iterator(chunk_size)) keeps the client's buffer atchunk_size. The test'sCaptureQueriesContextdocuments the number as a performance contract.
Exercise 3 — Window idempotency
- The expand-contract migration (11) over a ledger with data:
# 1) expand: constraint added as NOT VALID (doesn't scan the whole table)
migrations.RunSQL("ALTER TABLE billing_ledgerentry ADD CONSTRAINT uniq_ledger_org_mes UNIQUE (org_id, mes) NOT VALID;")
# 2) clean duplicates if validation fails:
# 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 = the constraint applies to new rows without a long lock over 2M rows: the grown-up way of adding uniqueness in production.
- The reminder:
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 runs the same day → the key already exists on runs 2-4 → ONE email. The 48 h TTL covers the event's day and a bit more.
Exercise 4 — Locks
- The lock with the polite skip:
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 busy: skip", name)
return None
try:
return fn(*args, **kwargs)
finally:
cache.delete(f"run-lock:{name}")
return wrapper
return decoThe test: cache.add mocked to False → returns None, zero calls to the logic, a log with "skip" (no exception: the overlap is expected behavior, not an error). The ADR (excerpt): "Run locks: Postgres advisory. Reason: they die with the connection (no TTL to tune), immune to containers' clock drift (43), auditable with pg_locks. Discarded Redis alternative: a badly tuned TTL = double run". In a K8s deployment (43), advisory wins; if the schedulers were multi-DB, Redis.
2-3. With a 4 h TTL and a 5 h run: at hour 4 the lock EXPIRES and the next interval's run enters in parallel — two processes in the same batch. With skip_locked they split rows without exceptions (the worst kind of bug: everything "works"), but the checkpoint advances twice and get_or_create saves the money — window idempotency is the safety net that turns the disaster into a non-event. The fix: a heartbeat — a threading.Timer renewing the TTL every 60 s while the run lives, and the finally cancels the timer and deletes the lock.
Exercise 5 — The closing
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() # state machine: 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() # orphans: ACTIVE in a closed event
outbox.create(event_type="ReservationExpired", payload={"ref": r.public_ref})
total += len(evs)
return totalThe integral test: 2nd run → the status=OPEN AND ends_at < now query returns empty (the events are already CLOSED): zero changes, zero new events. Idempotency by state machine + window is the pattern that makes the repeated run safe — every TicketFlow scheduled task ends up being "window query + batch + lock + outbox": four pieces, and only the domain varies.
Professor's summary
- The schedule lives in the repo; the business window (month/day/minute) is the idempotency key, not the execution attempt.
- Batch = cursor + commit per chunk + checkpoint +
get_or_create: it dies, resumes and doesn't duplicate. - Run lock with a polite skip (not an error); advisory lock when the environment's clock is suspect; a heartbeat if the run can outlast the TTL.