惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

D
DataBreaches.Net
罗磊的独立博客
雷峰网
雷峰网
量子位
V
Visual Studio Blog
Vercel News
Vercel News
freeCodeCamp Programming Tutorials: Python, JavaScript, Git & More
The Cloudflare Blog
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
宝玉的分享
宝玉的分享
月光博客
月光博客
Martin Fowler
Martin Fowler
aimingoo的专栏
aimingoo的专栏
H
Hackread – Cybersecurity News, Data Breaches, AI and More
Microsoft Security Blog
Microsoft Security Blog
博客园 - 叶小钗
腾讯CDC
Engineering at Meta
Engineering at Meta
博客园 - Franky
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
Y
Y Combinator Blog
Recent Announcements
Recent Announcements
Jina AI
Jina AI
A
About on SuperTechFans

DEV Community

Authentication Security Deep Dive: From Brute Force to Salted Hashing (With Java Examples) Why AI Systems Don’t Fail — They Drift Spilling beans for how i learn for exam😁"Reinforcement Learning Cheat Sheet" I Replaced Chrome with Safari for AI Browser Automation. Here's What Broke (and What Finally Worked) How Python Borrows Other People's Work The $40 Architecture: Processing 1 Billion API Requests with 99.99% Uptime Vibe Coding: A Workflow Guide (From Zero to SaaS) Most webhook security guides protect the wrong side. The scary part is delivery. Headless CMS for TanStack Start: Build a Blog with Cosmic EU Age Verification App "Hacked in 2 Minutes" — What Actually Happened Comfy Cloud’s delete function does not actually remove files Running AI Models on GPU Cloud Servers: A Beginner Guide Event-driven media intelligence with AWS Step Functions and Bedrock I scored 500 AI prompts across 8 quality dimensions — here's what broke How to Call Google Gemini API from Next.js (Free Tier, No Backend Needed) The Portal Protocol: Reclaiming Human Connection in the Age of AI How to Fix Your Team's Scattered Knowledge Problem With a Self-Hosted Forum Intro to tc Cloud Functors: A Graph-First Mental Model for the Modern Cloud Designing Multi-Tenant Backends With Both Ownership and Team Access I Built a Neumorphic CSS Library with 77+ Components — Here's What I Learned PostgreSQL Performance Optimization: Why Connection Pooling Is Critical at Scale Cómo construí un SaaS multi-rubro para gestionar expensas en Argentina con FastAPI + Vue 3 🚀 I Built an Ethical Hacking Scanner Tool – Open Source Project I Replaced /usage and /context in Claude Code With a Single Statusline A Pythonic Way to Handle Emails (IMAP/SMTP) with Auto-Discovery and AI-Ready Design I Collected 8.9 Million Polymarket Price Points — Here's What I Found About How Markets Really Move EcoTrack AI — Carbon Footprint Tracker & Dashboard Everyone's Using AI. No One Agrees How. 5 self-hosted ebook managers worth trying in 2026 Building Your First AI Agent with LangChain: From Chatbot to Autonomous Assistant
Cuatro intentos para que una tarea programada se ejecute ...
Franchesco Romero · 2026-06-03 · via DEV Community

Un backend de FastAPI corriendo en tres réplicas de Fargate manda una
sola notificación de "quedaste en el top 3" + bono de Coins cada viernes a las 18:00 de México. Un miembro la recibió tres veces. La solución tomó cuatro intentos y un cambio arquitectónico chico antes de que el bug se quedara resuelto.

El codebase: FastAPI sobre AWS ECS Fargate, desiredCount=3, un
Postgres en RDS, y un ciclo de worker en segundo plano
(moderator_bot_worker) dentro de cada réplica que despierta cada
15 minutos y reparte trabajo según la hora del reloj en Ciudad de
México. Las tareas tipo cron (resumen semanal, cierre del leaderboard) viven como branches de wall-clock adentro del mismo
ciclo del worker.

TL;DR

Intento Qué arregló Qué se le escapó
1. pg_try_advisory_xact_lock por ciclo Réplicas concurrentes dentro del mismo instante Ciclos secuenciales sobre la misma ranura — si el trabajo fallaba en silencio, cada ciclo posterior lo volvía a correr
2. Centinela por título del hilo adentro de la tarea El chequeo común de "¿ya escribimos el hilo del resumen?" Dependía de que bot_service.post_thread tuviera éxito. Cuando regresaba None en silencio (sección faltante), el centinela nunca persistía
3. Fila de reclamo en worker_runs Atómica, con alcance de la transacción: persiste con el trabajo, se revierte con el trabajo Sigue metiendo una decisión recurrente de cuándo correr dentro de cada réplica
4. app/jobs/ + objetivo de EventBridge (este post) Saca el "cuándo" de las réplicas por completo. AWS garantiza el disparo; la tarea es un contenedor efímero de ECS Latencia de arranque en frío de la tarea (~15-30s); un poco más de CDK

Cada intento tapó un agujero más estrecho que el anterior. El último
es estructural — en lugar de poner una capa más de bloqueos, elimina
la pregunta "¿qué réplica dispara el cron?" al dejar de disparar crons
dentro de las réplicas en primer lugar.

Acá puedes encontrar el código de acompañaniento:

Code companion — locking evolution post

Each folder maps to one stage of the post's narrative:

00-the-bug/           the original race scenario, reproducible in 30 lines
01-advisory-lock/     pg_try_advisory_xact_lock per tick (attempt 1)
02-thread-sentinel/   select-by-title dedupe (attempt 2, the silent-fail trap)
03-claim-row/         worker_runs table + claim_run() helper (attempt 3)
04-cli-entrypoint/    jobs/ package + __main__.py + worker delegates
05-eventbridge/       phase-2 CDK sketch (not deployed yet)

Suggested reading order

  1. 00-the-bug/ — the minimal reproducer. Run two replicas of the script against a shared Postgres and watch the duplicate notifications land.
  2. Each stage folder in order. Each one is a self-contained improvement; you can stop at any stage and still have something functional. The post argues that stopping earlier than 03-claim-row/ leaves you exposed.
  3. 05-eventbridge/ — the structural ending. Not strictly needed if you're happy with the worker model, but it removes the entire "which replica fires the cron?" question.

Notes

  • 00-the-bug/ is a self-contained…

El incidente

Viernes 18:00 de México, un miembro abre sus notificaciones y ve:

🥉 ¡Quedaste en el top 3! +10 Coins por tu actividad esta semana. (hace 17 min)
🥉 ¡Quedaste en el top 3! +10 Coins por tu actividad esta semana. (hace 18 min)
🥉 ¡Quedaste en el top 3! +10 Coins por tu actividad esta semana. (hace 19 min)

Tres notificaciones idénticas, separadas por un minuto. Su wallet
muestra +30 en lugar de +10. Lo mismo le pasó al #1 y al #2 de la
semana.

El intervalo del worker es de 15 minutos. Las notificaciones están a un minuto de distancia. Eso no son 3 ciclos. Son 3 réplicas
disparando la misma ranura una detrás de otra, el bloqueo de cada
réplica liberado por la confirmación de la anterior, y la siguiente
entrando antes de que cualquier centinela alcanzara a protegerla.

Intento 1: bloqueo por ciclo

El primer instinto fue correcto: serializar las réplicas. Los bloqueos consultivos de Postgres son baratos y no requieren cambios de esquema.

# myapp/workers/scheduled_worker.py
_TICK_LOCK_KEY = 0x4D424F54  # entero arbitrario de 32 bits, único para este worker

async def _try_tick_lock(db) -> bool:
    result = await db.execute(
        text("SELECT pg_try_advisory_xact_lock(:k)"),
        {"k": _TICK_LOCK_KEY},
    )
    return bool(result.scalar_one())

async def _tick():
    async with AsyncSessionLocal() as db:
        if not await _try_tick_lock(db):
            return  # otra réplica ya está adentro de este ciclo
        ...
        await db.commit()

pg_try_advisory_xact_lock no es bloqueante — regresa falso en
lugar de esperar si otra sesión lo tiene. Se libera de manera
automática cuando termina la transacción (commit o rollback). Dos
réplicas pegando contra _tick en el mismo instante: una se lleva
true y corre el trabajo; la otra se lleva false y sale limpia.

Esto se liberó a producción, y la manifestación obvia desapareció — ya no había disparos triples sincrónicos.

Lo que se me escapó. El bloqueo tiene alcance de transacción:
vive nada más adentro de una transacción. Si el ciclo del worker
es corto (10s) y el intervalo entre ciclos dentro de una sola réplica
es de 15 min, todo bien. Pero tres réplicas haciendo ciclos en
horarios escalonados — A en T=0, B en T+1min, C en T+2min — cada
confirmación libera el bloqueo para la siguiente. El bloqueo mantiene
honestos a los ciclos simultáneos, no a los secuenciales.

Para el cierre del leaderboard en un slot de 15 min, eso significa
hasta tres entradas secuenciales, una por minuto, exactamente lo que
mostraba la captura de pantalla.

Intento 2: centinela adentro de la tarea

La solución anterior resolvió el problema de las réplicas.
La función original ya tenía un chequeo de idempotencia distinto — "si ya existe un hilo con el título del resumen de esta semana, sal".

# myapp/workers/scheduled_worker.py
async def _leaderboard_close(db, now_mx):
    title = f"🏆 Top 3 de la semana — {now_mx.strftime('%d %b %Y')}"
    existing = (
        await db.execute(
            select(Thread).where(Thread.title == title).limit(1)
        )
    ).scalar_one_or_none()
    if existing:
        return  # ya corrimos esta slot
    ...
    # otorgar Coins, escribir notificaciones, postear el hilo
    await bot_service.post_thread(
        db, section_slug="general", title=title, body_md=...
    )

La suposición: una vez que el hilo del resumen existe, cada ciclo
posterior lee el título y se sale por la corta. Funciona si
post_thread escribe el hilo. No dice nada si post_thread no lo
escribe.

bot_service.post_thread era de mejor esfuerzo:

# myapp/services/bot_service.py
async def post_thread(db, *, section_slug, title, body_md):
    bot = await get_bot(db)
    if bot is None:
        logger.warning("post_thread: bot user missing")
        return None  # <-- falla silenciosa
    section = (
        await db.execute(
            select(Section).where(Section.slug == section_slug)
        )
    ).scalar_one_or_none()
    if section is None:
        logger.warning("post_thread: section %s missing", section_slug)
        return None  # <-- falla silenciosa
    db.add(Thread(...))
    await db.flush()
    return thread

Intento 3: claim de worker_runs

La solución de verdad tiene que ser atómica con el trabajo. Si el
trabajo se confirma, el centinela se confirma. Si algo se revierte, el centinela se revierte. Misma transacción.

Una tabla de una sola fila con clave primaria compuesta:

CREATE TABLE worker_runs (
    name TEXT NOT NULL,
    key  TEXT NOT NULL,
    ran_at TIMESTAMPTZ NOT NULL DEFAULT now(),
    PRIMARY KEY (name, key)
);

El helper para hacer claim:

# myapp/jobs/_idempotency.py
async def claim_run(db, name: str, key: str) -> bool:
    """Reserva (name, key). True si es nuestro, False si ya existía."""
    result = await db.execute(
        text(
            """
            INSERT INTO worker_runs (name, key)
            VALUES (:name, :key)
            ON CONFLICT (name, key) DO NOTHING
            RETURNING 1
            """
        ),
        {"name": name, "key": key},
    )
    return result.scalar_one_or_none() is not None

La tarea queda así:

async def run(db, now_mx):
    iso_year, iso_week, _ = now_mx.isocalendar()
    if not await claim_run(db, "leaderboard_close", f"{iso_year}-W{iso_week:02d}"):
        return  # alguien más es dueño de esta ranura
    # ... escribir notificaciones, transacciones, hilo ...
    # quien llama confirma o revierte

El INSERT ... ON CONFLICT DO NOTHING RETURNING 1 es la pieza
central. El RETURNING 1 nada más emite una fila cuando el insert
realmente sucedió. La sentencia completa es atómica a nivel de fila:
dos inserts concurrentes para el mismo (name, key) ven exactamente
un ganador.

Lo crucial: la fila participa en la transacción de quien la
llama
. Quien la llama (el ciclo del worker) confirma todo el
paquete al final: notificaciones + transacciones + hilo + la fila de
reclamo, o ninguno de ellos. Si post_thread regresa None callado
y el hilo del resumen nunca aterriza, el reclamo igual se revierte
si quien lo llamó lo trata como un error
— o, si se confirma de
todos modos, el siguiente ciclo igual ve el reclamo porque el
reclamo se insertó antes de la llamada rota.

Esta es la capa que los dos intentos anteriores no tenían. El
bloqueo consultivo era a nivel de proceso; el centinela del título
del hilo era a nivel de tarea pero consecuencia de efecto secundario;
este es a nivel de transacción y primario, sentado entre la decisión
de la ranura y cualquier efecto secundario.

Cómo se ve la secuencia con las tres capas

T=0   se dispara el ciclo de la réplica A
      ├─ pg_try_advisory_xact_lock → true (A gana el chequeo de concurrencia)
      ├─ claim_run('leaderboard_close', '2026-W22') → true (A gana la ranura)
      ├─ inserta 3 Notifications, 3 TokenTransactions, postea el Thread del resumen
      └─ COMMIT (la fila del reclamo + el trabajo aterrizan juntos)

T=60s  se dispara el ciclo de la réplica B
      ├─ pg_try_advisory_xact_lock → true (el bloqueo de A se liberó al confirmar A)
      ├─ claim_run('leaderboard_close', '2026-W22') → false (la fila ya existe)
      └─ return  ← lo que los intentos 1 + 2 no podían atrapar

T=120s se dispara el ciclo de la réplica C
      ├─ igual que B
      └─ return

El bloqueo evita que A y B simultáneas hagan el trabajo. La fila de
reclamo evita que B y C posteriores lo vuelvan a hacer. Juntos van
apretados.

¿Qué pasa si el trabajo falla a la mitad?

La transacción te protege. claim_run inserta la fila adentro de la
transacción de quien llama. La tarea que lo usa no es dueña de
ningún commit:

# myapp/workers/scheduled_worker.py
async def _tick():
    async with AsyncSessionLocal() as db:
        if not await _try_tick_lock(db):
            return
        ...
        if friday_18_window(now_mx):
            await leaderboard_close.run(db, now_mx)
        await db.commit()  # todo o nada

Si leaderboard_close.run levanta una excepción después de insertar
el reclamo, el async with revierte. El reclamo desaparece. El
siguiente ciclo se puede reintentar limpio. No hay una cola de
muertos permanente que andar limpiando.

Una prueba de regresión amarró este comportamiento:

@pytest.mark.asyncio
async def test_claim_run_rollback_releases_slot(db):
    from myapp.jobs._idempotency import claim_run

    assert await claim_run(db, "leaderboard_close", "2026-W42") is True
    await db.rollback()

    exists = (
        await db.execute(
            text("SELECT 1 FROM worker_runs WHERE name = :n AND key = :k"),
            {"n": "leaderboard_close", "k": "2026-W42"},
        )
    ).scalar_one_or_none()
    assert exists is None, (
        "Revertir tiene que liberar la ranura. Si la fila sobreviviera, un "
        "ciclo fallido bloquearía permanentemente cada reintento."
    )

En este punto todo funciona. ¿Por qué seguir?

Porque cada solución de arriba contesta "¿cómo hacemos idempotente el
ciclo del worker que ya se disparó?" en lugar de "¿debería el
worker estar disparando el ciclo?".

El worker sigue:

  • corriendo dentro de cada réplica de Fargate
  • despertando cada 15 minutos sin importar si hay algo programado
  • usando branches de wall-clock (if weekday == 4 and hour == 18 and minute < 15) para decidir qué hacer
  • va a seguir necesitando el bloqueo consultivo y el claim por siempre para tapar el desajuste fundamental de "tres réplicas, una sola tarea"

La respuesta estructural es: dejar de meter lógica de cron dentro de
las réplicas. AWS ya corre un planificador que maneja esto — varias
veces, con reintentos, con semántica de a lo-más una vez del lado del
consumidor vía idempotencia. Resumen: usarlo.

Intento 4: separar el "cuándo" del "dónde"

La refactorización que no es una solución encima de soluciones. Tres
piezas:

  1. Cada tarea ahora es invocable. El cuerpo de _leaderboard_close se movió a myapp/jobs/leaderboard_close.py, con esta forma:
   # myapp/jobs/leaderboard_close.py
   async def run(db, now_mx):
       iso_year, iso_week, _ = now_mx.isocalendar()
       if not await claim_run(db, "leaderboard_close", f"{iso_year}-W{iso_week:02d}"):
           return
       # ...

La función nunca confirma. Escribe sobre la sesión que le
pasaron. Quien la llama (el ciclo del worker, o la CLI, o una prueba) es dueño del límite de la transacción.

  1. Un punto de entrada por CLI. python -m myapp.jobs <name> corre exactamente una tarea:
   # myapp/jobs/__main__.py
   JOBS = {
       "leaderboard_close": leaderboard_close.run,
       "squad_health": squad_health.run,
       "weekly_digest": weekly_digest.run,
   }

   async def _run(job_name, now_iso):
       fn = JOBS[job_name]
       now_mx = (
           datetime.fromisoformat(now_iso).astimezone(MX_TZ)
           if now_iso else datetime.now(MX_TZ)
       )
       async with AsyncSessionLocal() as db:
           try:
               await fn(db, now_mx)
               await db.commit()
           except Exception:
               await db.rollback()
               raise

Es la misma forma que el worker usa internamente — abrir una
sesión, correr la tarea, confirmar-o-revertir. La fase 1 de la
refactorización hace que el worker y la CLI pasen por aquí,
idénticamente.

  1. EventBridge Scheduler → ECS RunTask. En CDK:
   // infra/lib/jobs-stack.ts (borrador de fase 2)
   new scheduler.CfnSchedule(this, 'LeaderboardClose', {
     scheduleExpression: 'cron(0 18 ? * FRI *)',
     scheduleExpressionTimezone: 'America/Mexico_City',
     target: {
       arn: 'arn:aws:scheduler:::aws-sdk:ecs:runTask',
       roleArn: schedulerRole.roleArn,
       input: JSON.stringify({
         Cluster: cluster.clusterArn,
         TaskDefinition: backendTaskDef.taskDefinitionArn,
         LaunchType: 'FARGATE',
         Overrides: {
           ContainerOverrides: [{
             Name: 'backend',
             Command: ['python', '-m', 'myapp.jobs', 'leaderboard_close'],
           }],
         },
       }),
     },
   });

EventBridge es dueño del "cuándo". ECS RunTask es dueño del
"dónde" — exactamente un contenedor efímero, sin réplicas
compitiendo. La fila de claim_run se queda como tercera capa:
EventBridge tiene entrega de al-menos-una-vez del lado del
planificador, así que si un parpadeo de red reintenta la llamada
a RunTask, el segundo contenedor es un no-op limpio.

Qué se simplifica

  • _tick pierde las tres branches de wall-clock. Nada más corre _anniversary + _zombie_threads (que sí necesitan dispararse cada 15 minutos, no a una hora fija del reloj).
  • _INTERVAL_SECONDS ya no es un concepto de cron; es el intervalo de sondeo para el trabajo que de verdad es de flujo continuo.
  • Las tareas nuevas requieren: un archivo en myapp/jobs/, una entrada en JOBS, un recurso de calendario en CDK. Se acaba el "¿calculamos bien la frontera de hour == 18 and minute < 15?".
  • La función _tick de 99 líneas con cinco asuntos (bloqueo consultivo, aniversario, hilos zombi, tres branches de wall-clock, commit) se vuelve un despachador de 30 líneas.

Lo que NO ayudó

  • Agregar un bloqueo distribuido basado en Redis. Misma forma que el bloqueo consultivo (exclusión mutua a nivel de réplica), peor historia ante particiones (que Redis pierda la llave durante un failover significa que una ranura se dispara dos veces). No metas otra dependencia para resolver un problema que la base de datos ya maneja de manera atómica.
  • Poner el trabajo detrás de un solo líder electo. "Una réplica es el planificador" vía elección de líder (consul, etcd) necesita renovación de leases, traspaso, y maquinaria de quórum que no vale la pena para cuatro tareas cron.
  • Mover la idempotencia a la capa de notificaciones. "Si el usuario ya recibió LEADERBOARD_TOP3 esta semana, salta." Suena atractivo pero acopla cada consumidor a la preocupación del planificador. La fila de worker_runs está aguas arriba de cada consumidor y mantiene el alcance contenido.

Lecciones

  • Cada solución hizo el bug menos probable hasta que un cambio estructural hizo el bug imposible en esta capa. No tomes el cambio estructural como tu primer movimiento al primer reporte. Sí tómalo cuando el tercer intento es "otro bloqueo más".
  • Atómico a nivel de fila, atómico con el trabajo, atómico ante fallas. INSERT ... ON CONFLICT DO NOTHING RETURNING 1 trae esas tres propiedades de regalo. Échale mano antes de echar mano de un servicio de bloqueo distribuido.
  • El mejor esfuerzo y la falla silenciosa no se llevan con la idempotencia. Si una función puede regresar None para decir "no hice mi trabajo", cada quien que llame y dependa de su efecto secundario para deduplicar es un duplicado a futuro. O la función levanta excepción, o quien la llama trata el None como falla, o quien la llama usa su propio centinela atómico.
  • La branch de wall-clock dentro de un loop huele mal. if weekday == 4 and hour == 18 and minute < 15 es reinventar cron con peor semántica, contra un sistema que ya tiene un planificador.