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:
SELECTluegoINSERT/UPDATE: race condition entre check y write.- Try
INSERT, catch error,UPDATE: feo, dos round-trips, posibles deadlocks. - 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 columnaid(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-tablaEXCLUDEDcon 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:
ON CONFLICTresuelve idempotency a nivel SQL, atómicamente.- Combinado con bulk insert, da imports re-runnable.
EXCLUDEDvstabla.colmatters mucho — distinguir bien.- 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_insertdesqlalchemy.dialects.postgresql, no elinsertgené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_inserten SQLAlchemy 2.0. - Distinguir cuándo usar
EXCLUDEDvstabla.col. - Aplicar
WHEREen 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
- PostgreSQL Docs —
INSERT ... ON CONFLICT— referencia oficial. - SQLAlchemy 2.0 —
pg_insert.on_conflict_do_update— referencia. - PostgreSQL Wiki — UPSERT — historia y casos.
- Brandur Leach — Postgres queries patterns — patterns con ON CONFLICT.
- Citus Data — UPSERT performance — casos a escala.
- Heap — Idempotent ETL — ON CONFLICT en pipeline.
Cápsula 05 de 08 — Módulo 7 — SQL Patterns for Production APIs Guide