Módulo 7: Bulk Operations

Upserts atómicos con `ON CONFLICT`

"Upsert" = INSERT or UPDATE. El caso típico: importás datos donde algunos registros ya existen. Sin upsert, tus opciones son:

  1. SELECT luego INSERT/UPDATE: race condition entre check y write.
  2. Try INSERT, catch error, UPDATE: feo, dos round-trips, posibles deadlocks.
  3. Truncate + re-insert: pierde datos, no idempotente.

PostgreSQL tiene la respuesta canónica desde 9.5: INSERT ... ON CONFLICT. Es atómico (sin race), idempotente (re-run el mismo INSERT no causa errores), y bulk-friendly (combina con executemany y COPY-to-temp).

En esta cápsula aprendés la sintaxis exacta, las tres variantes (DO NOTHING, DO UPDATE, partial), cómo usar EXCLUDED correctamente, y los anti-patterns más comunes.


La sintaxis básica

INSERT INTO tasks (id, title, status)
VALUES (1, 'Task 1', 'pending')
ON CONFLICT (id) DO UPDATE SET
    title = EXCLUDED.title,
    status = EXCLUDED.status,
    updated_at = NOW();

Lectura paso a paso:

  • INSERT INTO ... VALUES: intentá insertar.
  • ON CONFLICT (id): si hay conflict en la columna id (porque ya existe).
  • DO UPDATE SET ...: en lugar de fallar, hacé este UPDATE.
  • EXCLUDED.title: refiere al valor que iba a insertarse. PostgreSQL crea una pseudo-tabla EXCLUDED con la fila propuesta.

EXCLUDED es la palabra clave clave. Permite referirte al valor "nuevo" en el UPDATE, distinguiéndolo de la fila existente.


Tres variantes de ON CONFLICT

Variante 1: DO NOTHING (skip duplicates)

INSERT INTO tasks (id, title) VALUES (1, 'Task 1')
ON CONFLICT (id) DO NOTHING;

Si id=1 ya existe, no hace nada. Útil para imports idempotentes donde no querés sobrescribir.

Caso de uso típico: log de eventos donde duplicates son posibles pero el primero gana.

Variante 2: DO UPDATE (overwrite)

INSERT INTO tasks (id, title, status) VALUES (1, 'Updated', 'completed')
ON CONFLICT (id) DO UPDATE SET
    title = EXCLUDED.title,
    status = EXCLUDED.status;

Si id=1 existe, actualiza title y status con los valores nuevos.

Caso de uso típico: sync desde external source. El último valor siempre gana.

Variante 3: DO UPDATE parcial (con WHERE)

INSERT INTO tasks (id, title, status, version) VALUES (1, 'Updated', 'completed', 5)
ON CONFLICT (id) DO UPDATE SET
    title = EXCLUDED.title,
    status = EXCLUDED.status,
    version = EXCLUDED.version
WHERE tasks.version < EXCLUDED.version;  -- Solo update si la nueva version es más reciente

Update conditional. La fila se actualiza solo si la condición se cumple. Útil para "last-write-wins" con timestamps o version numbers.


Conflict target: qué columna(s)

ON CONFLICT (cols) requiere que las columnas tengan constraint UNIQUE o PRIMARY KEY:

-- ✅ id es PK
INSERT INTO tasks (id, ...) VALUES (...) ON CONFLICT (id) DO UPDATE ...;

-- ✅ email es UNIQUE
INSERT INTO users (email, ...) VALUES (...) ON CONFLICT (email) DO UPDATE ...;

-- ❌ status no es UNIQUE — error
INSERT INTO tasks (...) VALUES (...) ON CONFLICT (status) DO UPDATE ...;
-- ERROR: there is no unique or exclusion constraint matching the ON CONFLICT specification

Conflict en composite key

-- Constraint UNIQUE composite
ALTER TABLE order_items ADD CONSTRAINT uq_order_book UNIQUE (order_id, book_id);

INSERT INTO order_items (order_id, book_id, quantity) VALUES (1, 42, 3)
ON CONFLICT (order_id, book_id) DO UPDATE SET
    quantity = order_items.quantity + EXCLUDED.quantity;

Útil para "agregá si no existe, sumá quantity si ya existe".

Conflict en constraint name

-- Si tenés varios UNIQUE constraints, podés especificar cual
INSERT INTO ...
ON CONFLICT ON CONSTRAINT my_unique_constraint DO UPDATE ...;

Implementación en SQLAlchemy 2.0

Single row upsert

from sqlalchemy.dialects.postgresql import insert as pg_insert


async def upsert_task(session, task_data: dict):
    stmt = pg_insert(Task).values(**task_data)
    stmt = stmt.on_conflict_do_update(
        index_elements=["id"],  # ON CONFLICT (id)
        set_={
            "title": stmt.excluded.title,
            "status": stmt.excluded.status,
        }
    )

    await session.execute(stmt)
    await session.commit()

stmt.excluded es la traducción de EXCLUDED.col en SQL.

Bulk upsert con ON CONFLICT

async def bulk_upsert_tasks(session, records: list[dict]):
    stmt = pg_insert(Task).values(records)
    stmt = stmt.on_conflict_do_update(
        index_elements=["id"],
        set_={
            "title": stmt.excluded.title,
            "status": stmt.excluded.status,
        }
    )

    await session.execute(stmt)
    await session.commit()

pg_insert(Task).values(list_of_dicts) ejecuta INSERT múltiple con ON CONFLICT en cada fila. Atomic.

DO NOTHING

stmt = pg_insert(Task).values(records)
stmt = stmt.on_conflict_do_nothing(index_elements=["id"])

await session.execute(stmt)

Conditional UPDATE

from sqlalchemy import text

stmt = pg_insert(Task).values(records)
stmt = stmt.on_conflict_do_update(
    index_elements=["id"],
    set_={
        "title": stmt.excluded.title,
        "version": stmt.excluded.version,
    },
    where=Task.version < stmt.excluded.version,  # Conditional
)

Patrón completo en endpoint

from datetime import datetime, timezone
from fastapi import APIRouter, Depends
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.dialects.postgresql import insert as pg_insert


router = APIRouter()


@router.post("/tasks/sync")
async def sync_tasks(
    tasks: list[TaskSyncData],
    db: AsyncSession = Depends(get_db),
):
    """Sync tasks from external source (idempotente).

    - Si la task no existe (por external_id), insertarla.
    - Si existe, actualizarla con los nuevos valores.
    """
    records = [
        {
            "external_id": t.external_id,
            "title": t.title,
            "status": t.status,
            "synced_at": datetime.now(timezone.utc),
        }
        for t in tasks
    ]

    stmt = pg_insert(Task).values(records)
    stmt = stmt.on_conflict_do_update(
        index_elements=["external_id"],
        set_={
            "title": stmt.excluded.title,
            "status": stmt.excluded.status,
            "synced_at": stmt.excluded.synced_at,
        }
    )

    await db.execute(stmt)
    await db.commit()

    return {"synced": len(records)}

Cliente puede llamar este endpoint múltiples veces con la misma data — siempre llega al mismo estado final. Idempotencia perfecta.


Casos de uso comunes

1. Sync desde external source

# Cron job que sync con sistema externo
async def sync_from_external():
    external_data = fetch_from_external_api()
    await bulk_upsert_tasks(external_data)
    # Re-correrlo es safe

2. Counter incremental

-- Tracking de page views: incrementar contador o crear si no existe
INSERT INTO page_views (page_id, view_count, last_viewed)
VALUES ($1, 1, NOW())
ON CONFLICT (page_id) DO UPDATE SET
    view_count = page_views.view_count + 1,
    last_viewed = NOW();

page_views.view_count + 1 (referencia tabla original) es importante. EXCLUDED.view_count + 1 sería incorrecto (siempre 1+1=2).

3. Idempotency para mensajes deduplicados

-- Dedupe mensajes por message_id
INSERT INTO messages (message_id, content, created_at) VALUES ($1, $2, $3)
ON CONFLICT (message_id) DO NOTHING
RETURNING id;

Si el mensaje ya existe, no devuelve fila. App sabe que ya lo procesó.

4. Last-write-wins con version

INSERT INTO documents (id, content, version) VALUES ($1, $2, $3)
ON CONFLICT (id) DO UPDATE SET
    content = EXCLUDED.content,
    version = EXCLUDED.version
WHERE documents.version < EXCLUDED.version;

Solo updatea si la versión nueva es mayor. Race-safe.


Referenciar tabla actual vs EXCLUDED

Confusión común: cuándo usar tasks.col vs EXCLUDED.col:

INSERT INTO tasks (id, view_count) VALUES (1, 1)
ON CONFLICT (id) DO UPDATE SET
    view_count = tasks.view_count + 1;     -- valor actual + 1
    -- vs
    view_count = EXCLUDED.view_count + 1;  -- valor que iba a insertar (1) + 1 = 2 siempre
  • tasks.col = valor actual de la fila existente.
  • EXCLUDED.col = valor que ibas a insertar.

Para counters: tasks.view_count + 1. Para overwrite: EXCLUDED.col. Para suma incremental: tasks.col + EXCLUDED.col.


Limitaciones

ON CONFLICT no funciona con COPY directo

-- ❌ No existe COPY ... ON CONFLICT
COPY tasks FROM STDIN ON CONFLICT (id) DO UPDATE SET ...;  -- Error sintáctico

Para bulk upserts, el patrón es COPY a tabla temporal + INSERT...ON CONFLICT desde la temp (cápsula 06).

Solo una columna conflict target a la vez

No podés tener múltiples ON CONFLICT con diferentes columnas en una sola query. Solo un target.

RETURNING puede ser confuso

INSERT ... ON CONFLICT DO UPDATE ... RETURNING id, (xmax = 0) AS inserted;

xmax = 0 es truco para distinguir si fue INSERT (true) o UPDATE (false). Útil para reportar stats.


Trampas y errores comunes

1. ON CONFLICT (col) en columna sin UNIQUE constraint.

Error: "there is no unique or exclusion constraint matching the ON CONFLICT specification". Verificar que tenés UNIQUE o PRIMARY KEY en esa columna.

2. Confundir EXCLUDED con NEW.

NEW solo existe en triggers PL/pgSQL. En ON CONFLICT, es EXCLUDED.

3. EXCLUDED para counters.

-- ❌ Wrong: siempre suma 1+1=2 (o el valor inicial)
ON CONFLICT (id) DO UPDATE SET count = EXCLUDED.count + 1;

-- ✅ Right: incrementa el valor actual
ON CONFLICT (id) DO UPDATE SET count = tasks.count + 1;

4. Olvidar columnas en SET.

Si tu UPDATE solo setea title, otras columnas mantienen su valor existente. Para "overwrite completo", listar todas las columnas.

5. WHERE en ON CONFLICT con condición incorrecta.

-- Si la condición no se cumple, NO falla — simplemente no actualiza.
INSERT ... ON CONFLICT (id) DO UPDATE SET ... WHERE FALSE;

Esto es como DO NOTHING. Tener cuidado con la lógica del WHERE.

6. ON CONFLICT DO UPDATE con triggers PL/pgSQL.

Triggers se disparan en el INSERT subyacente y en el UPDATE resultante. Si tu trigger tiene side effects (logging, audit), se dispara dos veces. Verificar.

7. Performance esperada vs real.

ON CONFLICT agrega overhead per-row (check del unique constraint). En bulk muy grande, puede ser ~10-30% más lento que INSERT puro. Para 100k filas: 4s sin ON CONFLICT, 5-6s con.

8. bulk_insert (SQLAlchemy ORM bulk) NO soporta ON CONFLICT.

session.execute(insert(Model), [...]) no soporta on_conflict_do_update fácilmente. Usar pg_insert from sqlalchemy.dialects.postgresql que sí.


Ejercicio: implementar sync endpoint

Setup:

class Task(Base):
    __tablename__ = "tasks"

    id: Mapped[int] = mapped_column(primary_key=True)
    external_id: Mapped[str] = mapped_column(String(100), unique=True)
    title: Mapped[str] = mapped_column(String(200))
    status: Mapped[str] = mapped_column(String(50))
    synced_at: Mapped[datetime] = mapped_column(server_default=func.now())

external_id es UNIQUE — el conflict target.

Paso 1: implementar POST /tasks/sync con upsert.

from sqlalchemy.dialects.postgresql import insert as pg_insert


@router.post("/tasks/sync")
async def sync_tasks(tasks: list[TaskSync], db: AsyncSession = Depends(get_db)):
    records = [t.model_dump() for t in tasks]
    stmt = pg_insert(Task).values(records)
    stmt = stmt.on_conflict_do_update(
        index_elements=["external_id"],
        set_={
            "title": stmt.excluded.title,
            "status": stmt.excluded.status,
            "synced_at": stmt.excluded.synced_at,
        }
    )
    await db.execute(stmt)
    await db.commit()
    return {"synced": len(records)}

Paso 2: test idempotencia.

# Ejecutar dos veces — segundo no debe fallar
data = [{"external_id": "ext-1", "title": "T1", "status": "pending"}]
await sync_tasks(data, db)
await sync_tasks(data, db)  # No error, mismo estado final

Paso 3: test que valores se actualizan.

# Primera ejecución: insert
await sync_tasks([{"external_id": "ext-1", "title": "Original", "status": "pending"}], db)

# Segunda: update
await sync_tasks([{"external_id": "ext-1", "title": "Updated", "status": "completed"}], db)

# Verificar
task = await db.scalar(select(Task).where(Task.external_id == "ext-1"))
assert task.title == "Updated"
assert task.status == "completed"

Paso 4: agregar reporting (cuántos insert vs update).

stmt = pg_insert(Task).values(records).returning(
    Task.id,
    text("(xmax = 0) AS inserted")  # true si insert, false si update
)
stmt = stmt.on_conflict_do_update(...)

result = await db.execute(stmt)
rows = result.all()

inserted = sum(1 for r in rows if r.inserted)
updated = len(rows) - inserted

return {"inserted": inserted, "updated": updated}

Paso 5: agregar conditional update (solo si version es mayor).

stmt = stmt.on_conflict_do_update(
    index_elements=["external_id"],
    set_={...},
    where=Task.version < stmt.excluded.version,
)
Ver discusión

Paso 1 — implementación: funciona out-of-the-box.

Paso 2 — idempotencia:

Ejecutar 100 veces el mismo payload. Estado final es el mismo. No errores.

Paso 3 — actualización:

Tasks existentes se actualizan. synced_at también, lo que da auditabilidad.

Paso 4 — reporting:

xmax = 0 es PostgreSQL trick para identificar si fue INSERT vs UPDATE. Útil para devolver stats al cliente.

Paso 5 — conditional:

Si una task tiene version=5 localmente y recibís data con version=3 (out-of-order), no se actualiza. Protección contra race conditions de delivery.

Lecciones clave:

  1. ON CONFLICT resuelve idempotency a nivel SQL, atómicamente.
  2. Combinado con bulk insert, da imports re-runnable.
  3. EXCLUDED vs tabla.col matters mucho — distinguir bien.
  4. Conditional WHERE protege contra updates fuera de orden.

Resumen y siguiente paso

Lo que aprendiste:

  • INSERT ... ON CONFLICT (col): upsert atómico nativo de PostgreSQL.
  • 3 variantes: DO NOTHING (skip), DO UPDATE (overwrite), DO UPDATE WHERE (conditional).
  • EXCLUDED = valor que iba a insertarse. tabla.col = valor actual.
  • Conflict target: requiere UNIQUE/PK constraint en la(s) columna(s).
  • SQLAlchemy 2.0: usar pg_insert de sqlalchemy.dialects.postgresql, no el insert genérico.
  • Bulk upsert funciona con lista de dicts.
  • Limitación: COPY directo no soporta ON CONFLICT — necesita patrón con tabla temporal (cápsula 06).

Antes de avanzar, deberías poder:

  • Escribir SQL INSERT ... ON CONFLICT (col) DO UPDATE SET col = EXCLUDED.col.
  • Implementar upsert con pg_insert en SQLAlchemy 2.0.
  • Distinguir cuándo usar EXCLUDED vs tabla.col.
  • Aplicar WHERE en ON CONFLICT para conditional updates.

En la siguiente cápsula combinamos COPY con ON CONFLICT — el patrón canónico para bulk upserts grandes. Como COPY no soporta ON CONFLICT directo, el patrón es: COPY a tabla temporal → INSERT...ON CONFLICT desde la temp. Vas a aprender el setup completo y por qué este patrón es el que usa cualquier ETL de PostgreSQL serio.


Recursos

  1. PostgreSQL Docs — INSERT ... ON CONFLICT — referencia oficial.
  2. SQLAlchemy 2.0 — pg_insert.on_conflict_do_update — referencia.
  3. PostgreSQL Wiki — UPSERT — historia y casos.
  4. Brandur Leach — Postgres queries patterns — patterns con ON CONFLICT.
  5. Citus Data — UPSERT performance — casos a escala.
  6. Heap — Idempotent ETL — ON CONFLICT en pipeline.

Cápsula 05 de 08 — Módulo 7 — SQL Patterns for Production APIs Guide