Saltar al contenido principal

SQLAlchemy para web scraping: una capa de almacenamiento de producción

HT

Hinata Tomoda

Ingeniero web y analista independiente

20 min de lectura

Respuesta corta: en un pipeline de scraping, SQLAlchemy no es la parte que descarga, sino la que hace innecesario volver a descargar. Modela el registro que realmente necesitas, acota el pool de conexiones por debajo del límite de la propia base de datos, escribe cada fila mediante un ON CONFLICT DO UPDATE para que un reintento no cueste nada, y reintenta la transacción completa en lugar de la sentencia. Todo lo que sigue apunta a SQLAlchemy 2.0 — 2.0.52, publicada el 11 de agosto de 2026 — y cita los valores por defecto oficiales en lugar de adivinarlos.

Claves del artículo

  • La capa de almacenamiento es antes que nada una palanca de coste: el tráfico de proxy se factura por gigabyte, así que la petición más barata es la que tu base de datos ya respondió.
  • En SQLAlchemy 2.0 la anotación de tipo ES el esquema: Mapped[str] genera NOT NULL y Mapped[str | None] genera NULL, de modo que el verificador de tipos y el DDL no pueden divergir.
  • Los valores por defecto documentados del pool son pool_size 5, max_overflow 10, pool_timeout 30, pool_recycle -1 y pool_pre_ping False; cada proceso worker adicional multiplica los dos primeros contra el límite de conexiones de tu base.
  • La idempotencia viene de on_conflict_do_update() sobre una clave natural, con una cláusula WHERE que salta las filas cuyo hash de contenido no ha cambiado: un rastreo sin novedades no escribe nada.
  • Bajo asyncio las reglas son una AsyncSession por tarea y expire_on_commit=False; una sesión compartida o una carga perezosa accidental es el error MissingGreenlet que de otro modo depurarás a las 2 de la madrugada.
  • Cuando cae una conexión, la transacción se pierde con ella: la recuperación documentada es reintentar la operación desde el principio, algo que solo es seguro porque las escrituras son idempotentes.

Dónde encaja la base de datos en un pipeline de scraping

Un pipeline de producción tiene cinco etapas — obtención de URL, descarga, parseo, validación y almacenamiento — y nuestra guía de web scraping recorre la cadena completa. Solo una de esas etapas cuesta dinero por byte. El tráfico de proxies residenciales se factura por gigabyte, así que la economía del pipeline la decide con qué frecuencia tienes que descargar una página por segunda vez; nuestra estimación de costes de web scraping pone cifras reales a eso.

Eso replantea la capa de almacenamiento. Su trabajo no es «guardar las filas en algún sitio». Su trabajo es responder a tres preguntas lo bastante barato como para que el descargador nunca tenga que hacerlo:

  1. ¿He visto ya este registro y ha cambiado? Si la respuesta es no, el rastreo se ahorra una petición facturada.
  2. ¿Qué sobrevive a una caída? Un worker que muere a mitad de lote no puede dejar media página confirmada y la otra media perdida.
  3. ¿Pueden leer los analistas mientras el crawler escribe? Una transacción larguísima que retiene bloqueos convierte al almacén y al crawler en enemigos.

SQLAlchemy responde a las tres, en dos capas que conviene no mezclar. Core es el lenguaje de expresiones SQL: Table, select(), insert(), el motor y el pool de conexiones. El ORM añade encima clases mapeadas, un identity map y la Session como unidad de trabajo. El reparto productivo para un scraper es definir el esquema con modelos declarativos del ORM —así los tipos de Python y el DDL tienen una única fuente de verdad— y escribir los lotes con sentencias de estilo Core, porque insertar diez mil filas es una operación de conjuntos, no diez mil mutaciones de objetos.

Modela el registro, no la página que scrapeaste

El error de modelado más caro en un scraper es guardar el artefacto en lugar del hecho. El HTML crudo es enorme, cambia por motivos que no te importan y no se puede consultar. Guarda el registro extraído más los metadatos justos para decidir si es nuevo.

Python
from __future__ import annotations

import datetime as dt
from typing import Annotated, Any

from sqlalchemy import (
    BigInteger,
    ForeignKey,
    Index,
    MetaData,
    String,
    TIMESTAMP,
    UniqueConstraint,
    func,
)
from sqlalchemy.dialects.postgresql import JSONB
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column, relationship

# Column recipes declared once and annotated everywhere. Changing the policy
# for every timestamp in the schema is then a one-line edit, not a sweep.
bigint_pk = Annotated[int, mapped_column(BigInteger, primary_key=True)]
utc_now = Annotated[dt.datetime, mapped_column(server_default=func.now())]

# Every constraint gets a deterministic name. Alembic cannot generate a
# migration for a constraint the database named anonymously, so this is what
# makes autogenerate usable later.
NAMING_CONVENTION = {
    "ix": "ix_%(column_0_label)s",
    "uq": "uq_%(table_name)s_%(column_0_name)s",
    "ck": "ck_%(table_name)s_%(constraint_name)s",
    "fk": "fk_%(table_name)s_%(column_0_name)s_%(referred_table_name)s",
    "pk": "pk_%(table_name)s",
}


class Base(DeclarativeBase):
    metadata = MetaData(naming_convention=NAMING_CONVENTION)

    # A naive timestamp is the classic silent data loss in a crawler that runs
    # in more than one region. Fix it once, for every model in the schema.
    type_annotation_map = {
        dt.datetime: TIMESTAMP(timezone=True),
        dict[str, Any]: JSONB,
    }


class ScrapedPage(Base):
    __tablename__ = "scraped_page"

    id: Mapped[bigint_pk]
    source: Mapped[str] = mapped_column(String(64))
    external_id: Mapped[str] = mapped_column(String(256))
    url: Mapped[str] = mapped_column(String(2048))
    # Mapped[str] is NOT NULL; Mapped[str | None] is NULL. The annotation IS
    # the schema, so mypy and the DDL cannot disagree about optionality.
    title: Mapped[str | None]
    payload: Mapped[dict[str, Any]]
    # Hash what you EXTRACTED, never the raw HTML: hash the page and every
    # rotating ad slot and CSRF token reads as a content change.
    content_hash: Mapped[str] = mapped_column(String(64))
    first_seen_at: Mapped[utc_now]
    last_changed_at: Mapped[utc_now]

    observations: Mapped[list[PriceObservation]] = relationship(
        back_populates="page", cascade="all, delete-orphan"
    )

    __table_args__ = (
        # The natural key: what "the same page, seen again" means. Without it
        # the upsert in the next section has nothing to conflict on.
        UniqueConstraint("source", "external_id", name="uq_scraped_page_natural_key"),
        Index("ix_scraped_page_changed", "source", "last_changed_at"),
    )


class PriceObservation(Base):
    __tablename__ = "price_observation"

    id: Mapped[bigint_pk]
    page_id: Mapped[int] = mapped_column(ForeignKey("scraped_page.id", ondelete="CASCADE"))
    observed_at: Mapped[utc_now]
    # Money as integer minor units. A float price silently loses cents, and
    # you will not notice until someone reconciles a year of history.
    amount_minor: Mapped[int]
    currency: Mapped[str] = mapped_column(String(3))

    page: Mapped[ScrapedPage] = relationship(back_populates="observations")

Tres detalles ahí dentro hacen casi todo el trabajo.

Mapped[] decide la nulabilidad. La regla documentada es que Mapped[str] produce NOT NULL y Mapped[Optional[str]] —o Mapped[str | None]— produce NULL, y que primary_key=True implica NOT NULL en cualquier caso. Puedes forzarlo con mapped_column(nullable=...), pero por defecto tu verificador de tipos lee la misma verdad que tu esquema.

type_annotation_map centraliza la correspondencia de Python a SQL. SQLAlchemy ya mapea int a Integer, str a String, datetime.datetime a DateTime, uuid.UUID a Uuid y así sucesivamente. Sobrescribir ese mapa en la clase base es como se aplica una política de todo el proyecto —marcas de tiempo con zona horaria, JSONB en vez del JSON genérico— sin repetirla en 40 columnas.

La clave natural es una restricción real, no una convención. source más external_id es lo que hace reconocible la misma ficha de producto entre rastreos. Un id sustituto por sí solo no puede expresarlo, y todo lo de la siguiente sección depende de que la base de datos pueda detectar la colisión por sí misma.

El motor y el pool son tu presupuesto de concurrencia

create_engine() es una fábrica, no una conexión: mantiene un pool de conexiones que se crea de forma perezosa y se comparte. Salvo en SQLite en memoria, el pool por defecto es QueuePool, y sus valores por defecto documentados son los números que te morderán primero.

ParámetroPor defectoQué decide en realidad
pool_size5Conexiones mantenidas abiertas por motor y por proceso
max_overflow10Conexiones extra en picos, descartadas tras su uso
pool_timeout30Segundos que un worker espera una conexión libre antes de fallar
pool_recycle-1 (desactivado)Antigüedad en segundos tras la cual se reemplaza la conexión
pool_pre_pingFalseSi se comprueba la conexión con una consulta barata al tomarla
insertmanyvalues_page_size1000Filas por INSERT agrupado de varias filas

La aritmética que importa es workers × (pool_size + max_overflow) ≤ límite de conexiones − margen. Con los valores por defecto, ocho procesos de rastreo pueden exigir 120 conexiones a una base que quizá permita 100 en total, y el fallo aparece como pool_timeout agotándose en los workers mientras una migración ni siquiera consigue conectar. Decide el techo de forma deliberada:

Python
import os

from sqlalchemy import create_engine

engine = create_engine(
    # Never a literal DSN. Credentials belong in the environment or a secrets
    # manager, and the URL carries the driver: postgresql+psycopg for sync.
    os.environ["SCRAPER_DATABASE_URL"],
    pool_size=5,
    max_overflow=5,          # hard ceiling of 10 per process, not 15
    pool_timeout=10,         # fail fast; a queued worker is a stalled worker
    pool_recycle=1800,       # under a proxy or a managed DB, assume idle kills
    pool_pre_ping=True,      # one cheap round trip beats one lost batch
    insertmanyvalues_page_size=500,
)

pool_pre_ping=True es la estrategia pesimista documentada frente a desconexiones: SQLAlchemy envía un ping propio del dialecto (normalmente SELECT 1) al tomar la conexión y, si falla, la descarta e invalida todas las conexiones del pool más antiguas que ese instante. Cuesta un viaje de ida y vuelta por cada toma y elimina toda la clase de fallos «la primera consulta tras un periodo de inactividad falla». pool_recycle es su complemento para motores que cierran conexiones ociosas por temporizador: la documentación lo señala como remedio inmediato al server has gone away de MySQL, y el mismo razonamiento vale para cualquier PostgreSQL gestionado detrás de un proxy de conexiones.

El fork es donde se rompen los pools

Los crawlers multiproceso se topan con un peligro concreto y documentado: las conexiones del pool no se comparten con un proceso bifurcado. Dos procesos acaban escribiendo por el mismo socket, y el síntoma es un estado de protocolo corrupto en vez de un error limpio. El remedio oficial es desechar el pool heredado en el hijo, sin cerrar los sockets del padre:

Python
from multiprocessing import Pool


def init_worker() -> None:
    # close=False drops the inherited connection references WITHOUT closing
    # them — the parent still owns those sockets. The child then opens its own.
    engine.dispose(close=False)


with Pool(processes=8, initializer=init_worker) as pool:
    pool.map(scrape_one, urls)

Una sesión por unidad de trabajo

La Session está documentada como «a mutable, stateful object that represents a single database transaction», y «cannot be shared among concurrent threads or asyncio tasks without careful synchronization». Léelo como una regla de diseño y no como una advertencia: la vida de una sesión es la de una transacción, que es el tiempo que retienes bloqueos.

Crea el sessionmaker una sola vez a nivel de módulo y abre una sesión por unidad de trabajo:

Python
from sqlalchemy.orm import sessionmaker

# Once, at import time. The factory is cheap and shareable; sessions are not.
SessionFactory = sessionmaker(engine)


def store_batch(rows: list[dict[str, object]]) -> None:
    # One transaction per BATCH. Per row, commit overhead dominates the run;
    # per crawl, one bad page rolls back an hour of work while every row you
    # did write stays locked away from readers.
    with SessionFactory() as session, session.begin():
        session.execute(upsert_pages(rows))

session.begin() como gestor de contexto confirma al terminar bien y revierte ante cualquier excepción, así que no existe camino por el que escape un lote a medio escribir. Los lotes de unas 200 a 1.000 filas son el punto dulce habitual: lo bastante grandes para que los viajes de ida y vuelta dejen de dominar, lo bastante pequeños para que revertir sea barato y la duración de los bloqueos se quede en decenas de milisegundos.

Haz idempotente cada escritura

Esta es la sección que justifica el artículo. Un crawler es un sistema que se volverá a ejecutar: reintentos, rellenos históricos, un operador que relanza el trabajo de ayer. Si escribir dos veces la misma página produce duplicados o pisa datos buenos, cada uno de esos eventos se convierte en incidente. INSERT ... ON CONFLICT DO UPDATE traslada la deduplicación a la base de datos, donde es atómica.

Python
from sqlalchemy import func
from sqlalchemy.dialects.postgresql import insert


def upsert_pages(rows: list[dict[str, object]]):
    stmt = insert(ScrapedPage).values(rows)
    return stmt.on_conflict_do_update(
        # Infers the unique index behind the natural key.
        index_elements=[ScrapedPage.source, ScrapedPage.external_id],
        # `excluded` is the row PostgreSQL tried to insert — the fresh scrape.
        set_={
            "url": stmt.excluded.url,
            "title": stmt.excluded.title,
            "payload": stmt.excluded.payload,
            "content_hash": stmt.excluded.content_hash,
            "last_changed_at": func.now(),
        },
        # The line that pays for itself: a re-crawl that found nothing new
        # writes no row at all. No dead tuple, no WAL, no index churn.
        where=ScrapedPage.content_hash != stmt.excluded.content_hash,
    )

Fíjate en la ruta de importación: el upsert vive en el insert específico del dialecto, insert de sqlalchemy.dialects.postgresql, no en el genérico. SQLite ofrece la misma grafía on_conflict_do_update(), MySQL y MariaDB ofrecen on_duplicate_key_update(), y no hay construcción portable que cubra las tres. Es una razón sólida para desarrollar contra el motor que despliegas.

La cláusula where merece un momento de honestidad sobre su compromiso. Saltarse la actualización significa que last_changed_at registra cuándo cambió el contenido por última vez, que es justo lo que quiere un pipeline de detección de cambios, pero ya no registras cuándo comprobaste por última vez. Si necesitas ambas cosas, conserva la actualización protegida para el registro y escribe la marca «visto en» barata en una tabla estrecha aparte contra la que nadie haga joins.

La recompensa es un feed de cambios gratis. Como las filas sin cambios se saltan, RETURNING devuelve exactamente las que se movieron:

Python
changed = session.scalars(
    upsert_pages(rows).returning(ScrapedPage),
    # Refreshes any of these objects already living in the session's identity
    # map — unlike a plain INSERT, some of these rows already existed.
    execution_options={"populate_existing": True},
).all()

for page in changed:
    enqueue_downstream(page.id)

Esa es la diferencia entre «reindexamos todo cada noche» y «reindexamos el 0,4 % que cambió». En un corpus de un millón de páginas, es la diferencia entre un trabajo de minutos y otro que ocupa la noche entera.

Escrituras masivas: deja que insertmanyvalues agrupe

SQLAlchemy 2.0 acepta una lista de diccionarios como conjunto de parámetros de un insert(), y el ORM interpreta las claves como nombres de atributo en vez de nombres de columna: un cambio deliberado de la 2.0 que importa en cuanto un atributo mapeado y su columna se escriben distinto.

Python
from sqlalchemy import insert

session.execute(
    # render_nulls keeps every row in one batch. Scraped records are ragged by
    # nature — an optional field that is None would otherwise split the run
    # into several statements at exactly the rows you have most of.
    insert(ScrapedPage).execution_options(render_nulls=True),
    rows,
)

Por debajo, la función insertmanyvalues reescribe eso en sentencias INSERT ... VALUES agrupadas de varias filas. Está activa por defecto en PostgreSQL, MySQL, SQLite, SQL Server y Oracle, y el tamaño de lote sigue a insertmanyvalues_page_size, que «defaults to 1000, but may also be subject to dialect-specific limiting factors». Bájalo cuando las filas lleven grandes cargas JSON: la ganancia viene de reducir viajes, y una sentencia que supere los límites de parámetros del servidor devuelve esa ganancia de inmediato.

Cuando necesites las claves primarias generadas, pídelas en el orden de los parámetros:

Python
page_ids = session.scalars(
    # sort_by_parameter_order guarantees the returned ids line up with the
    # input rows, which is what lets you attach child records without a
    # second SELECT. Added in SQLAlchemy 2.0.10.
    insert(ScrapedPage).returning(ScrapedPage.id, sort_by_parameter_order=True),
    rows,
).all()

Dos cosas que no conviene hacer. No recorras session.add() sobre diez mil objetos para confirmar una vez: pagas toda la contabilidad de la unidad de trabajo por filas que nadie va a mutar. Y no construyas el SQL a mano para «evitar la sobrecarga del ORM»; la vía parametrizada es precisamente lo que mantiene el texto scrapeado —entrada no confiable por definición— fuera de tu gramática SQL.

Pipelines asíncronos: una AsyncSession por tarea

Si tu descargador ya es asyncio —la mayoría de los modernos lo son, y nuestras notas sobre proxies para agentes de IA explican por qué se ha extendido esa forma—, ejecutar la base de datos en el mismo bucle de eventos evita un salto de hilo. Las reglas son estrechas y no negociables.

Python
import asyncio
import os

from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine

# The URL carries the async driver: postgresql+asyncpg://...
engine = create_async_engine(os.environ["SCRAPER_DATABASE_URL"], pool_size=10, max_overflow=0)

# expire_on_commit=False so attributes stay readable after commit. The default
# would expire them, and refreshing an expired attribute is implicit IO —
# exactly what asyncio cannot do behind your back.
AsyncSessionFactory = async_sessionmaker(engine, expire_on_commit=False)


async def store(rows: list[dict[str, object]], gate: asyncio.Semaphore) -> None:
    # One AsyncSession per task: a single instance is explicitly documented as
    # unsafe in concurrent tasks. The semaphore keeps the task count from
    # outrunning the pool, which is the other half of the same budget.
    async with gate, AsyncSessionFactory() as session, session.begin():
        await session.execute(upsert_pages(rows))


async def main(batches: list[list[dict[str, object]]]) -> None:
    gate = asyncio.Semaphore(10)
    try:
        await asyncio.gather(*(store(batch, gate) for batch in batches))
    finally:
        # Closes the pool. Skip it and the loop shuts down under open sockets.
        await engine.dispose()

El error que te encontrarás si lo haces mal se llama MissingGreenlet, y casi siempre significa una carga perezosa disparada fuera del contexto asíncrono. La documentación de asyncio ofrece tres soluciones, en el orden en que yo las probaría: cargar por adelantado con selectinload() en la consulta, añadir el mixin AsyncAttrs a tu base y escribir await obj.awaitable_attrs.things, o bajar a session.run_sync() para un bloque de código ORM síncrono corriente.

Una advertencia sobre el dimensionado: lo asíncrono no eleva el límite de conexiones de tu base, solo hace mucho más fácil alcanzarlo. Mil tareas concurrentes contra pool_size=10 no son un interbloqueo, pero todas las tareas a partir de la décima quedan en cola, y pool_timeout decide si esa cola falla en voz alta o se convierte en silencio en tu latencia.

Releer sin N+1

La ruta de lectura es donde se exportan los datos scrapeados, y donde vive el clásico fallo de rendimiento del ORM: iterar sobre padres tocando una relación en cada vuelta emite una consulta por padre. La respuesta de SQLAlchemy es declarar la estrategia de carga en la consulta.

Python
from sqlalchemy import select
from sqlalchemy.orm import raiseload, selectinload

stmt = (
    select(ScrapedPage)
    .options(
        # selectinload is the documented default for COLLECTIONS: a second
        # SELECT with an IN clause, so no join multiplies the parent rows.
        selectinload(ScrapedPage.observations),
        # Everything not eager-loaded above now RAISES instead of quietly
        # emitting a query. An N+1 becomes a test failure, not a slow export.
        raiseload("*"),
    )
    .where(ScrapedPage.source == "example-retailer")
    .order_by(ScrapedPage.last_changed_at.desc())
    .limit(500)
)
pages = session.scalars(stmt).all()

Para una relación escalar de muchos a uno —una observación hacia su página— joinedload() es la mejor forma, porque el join añade un conjunto de columnas en vez de un segundo viaje. La única trampa que la documentación señala expresamente: joinedload() sobre una colección multiplica filas, así que hay que deduplicar los resultados con .unique() antes de contar nada.

raiseload("*") está infravalorado. Actívalo en las pruebas y en cualquier exportación por lotes, y toda carga perezosa accidental se convierte en un fallo ruidoso y localizado durante el desarrollo, en lugar de un misterioso trabajo de una hora en producción.

Observabilidad: mide el SQL, no la sensación

echo=True es un interruptor de desarrollo. En producción registra todas las sentencias, incluidos los parámetros ligados, que en un scraper significan contenido scrapeado y, en la propia conexión, credenciales. Usa eventos en su lugar y registra la sentencia sin los parámetros.

Python
import logging
import time

from sqlalchemy import event

log = logging.getLogger("scraper.sql")
SLOW_QUERY_SECONDS = 0.5


@event.listens_for(engine, "before_cursor_execute")
def _start_timer(conn, cursor, statement, parameters, context, executemany):
    conn.info.setdefault("query_start", []).append(time.perf_counter())


@event.listens_for(engine, "after_cursor_execute")
def _record_duration(conn, cursor, statement, parameters, context, executemany):
    elapsed = time.perf_counter() - conn.info["query_start"].pop()
    if elapsed >= SLOW_QUERY_SECONDS:
        # The statement, never `parameters`: that tuple is the payload.
        log.warning("slow sql %.3fs executemany=%s %s", elapsed, executemany, statement)

La firma del evento es fija —(conn, cursor, statement, parameters, context, executemany) en ambos hooks— y la pila sobre conn.info mantiene bien emparejadas las ejecuciones anidadas. Vale la pena exportar tres métricas: duración de la sentencia por operación, proporción de executemany (sube cuando tu agrupado funciona) y tiempo de espera al tomar del pool, el número que te dice si el crawler va lento por el sitio objetivo o por tu propio pool_size.

Reintenta la transacción, nunca la sentencia

La documentación es tajante: «when a connection is lost, the entire transaction is lost. There is no useful way that the database can reconnect and retry and continue where it left off». La forma recomendada es «retry the entire operation from the start of the transaction».

Python
import random
import time
from collections.abc import Callable
from typing import TypeVar

from sqlalchemy.exc import DBAPIError

T = TypeVar("T")
# Serialization failure and deadlock detected: transient by definition.
RETRYABLE_SQLSTATES = frozenset({"40001", "40P01"})


def _is_retryable(error: DBAPIError) -> bool:
    if error.connection_invalidated:
        return True
    orig = error.orig
    return getattr(orig, "sqlstate", None) in RETRYABLE_SQLSTATES


def with_retry(work: Callable[[], T], attempts: int = 5) -> T:
    """Re-run a whole unit of work, from BEGIN, on a transient failure."""
    for attempt in range(1, attempts + 1):
        try:
            return work()
        except DBAPIError as error:
            if attempt == attempts or not _is_retryable(error):
                raise
            # Full jitter. A fleet that backs off on one schedule reconnects
            # as a thundering herd and knocks the database over a second time.
            time.sleep(random.uniform(0, min(2**attempt * 0.1, 5.0)))
    raise AssertionError("unreachable")


with_retry(lambda: store_batch(rows))

Esto solo es seguro porque la escritura es un upsert. Reintentar un INSERT simple tras un fallo ambiguo es como se generan duplicados; reintentar un ON CONFLICT DO UPDATE converge en la misma fila por muchas veces que se ejecute. Aquí la idempotencia no es un lujo: es la condición previa que hace posible la recuperación automática.

pool_pre_ping y este decorador resuelven mitades distintas del problema. El pre-ping atrapa la conexión que murió mientras estaba ociosa en el pool, antes de empezar el trabajo. El reintento atrapa la conexión que muere a mitad de transacción, cuando ya no hay nada que salvar.

Cambios de esquema: autogenerate es un borrador, no una migración

Alembic es la herramienta de migraciones del mismo proyecto, y su propia documentación dice que autogenerate «is not intended to be perfect» y que «it is always necessary to manually review and correct the candidate migrations that autogenerate produces». Tómatelo al pie de la letra.

Shell
alembic revision --autogenerate -m "add price_observation"
alembic upgrade head

Detecta de forma fiable altas y bajas de tablas y columnas, cambios de nulabilidad, cambios básicos en índices y restricciones únicas con nombre, cambios básicos en claves foráneas y restricciones CHECK con nombre. No detecta renombrados de tabla ni de columna —ambos aparecen como un borrado más un alta, lo que en un corpus scrapeado significa borrar la columna y reconstruirla vacía— y no ve en absoluto las restricciones con nombre anónimo. Por eso la naming_convention está en el MetaData del primer bloque de código: sin ella, autogenerate no tiene con qué emparejar las restricciones.

El flujo que aguanta: generar, leer el fichero línea a línea, reescribir los pares borrado-alta como op.alter_column(..., new_column_name=...) y ejecutar la migración contra una copia restaurada de producción antes de que se acerque a producción.

Pruebas: revertir en lugar de limpiar

Las pruebas de un scraper necesitan una base que se comporte como la real y una fixture que no deje residuos. SQLAlchemy 2.0 documenta el patrón: abrir una transacción externa sobre una conexión, ligar la sesión a esa conexión en modo savepoint y revertir la transacción externa al terminar la prueba.

Python
import pytest
from sqlalchemy.orm import Session


@pytest.fixture()
def session(engine):
    connection = engine.connect()
    transaction = connection.begin()
    # join_transaction_mode="create_savepoint" lets the CODE UNDER TEST call
    # session.commit() for real — the commit lands on a SAVEPOINT inside the
    # outer transaction, which the fixture then throws away wholesale.
    db = Session(bind=connection, join_transaction_mode="create_savepoint")
    try:
        yield db
    finally:
        db.close()
        transaction.rollback()
        connection.close()

Ejecútalas contra el mismo motor que despliegas. SQLite es un objetivo válido para las partes de Python puro, pero la semántica de ON CONFLICT, los operadores JSONB y el orden de los NULL difieren, y una suite que pasa en SQLite mientras producción corre Postgres está probando otro programa.

Cinco fallos que encuentro una y otra vez en bases de scrapers

SíntomaCausa raízSolución
Filas duplicadas tras cada reintentoSin restricción única en la clave natural, ON CONFLICT no tiene nada que detectarAñadir primero la restricción y luego el upsert; deduplicar el histórico una vez
La tabla crece pero el número de filas noUn upsert sin guarda reescribe cada fila en cada rastreoAñadir la guarda de content_hash a la cláusula where
Los workers se atascan y la base parece ociosapool_size × workers supera el límite del servidor; todos esperan en pool_timeoutFijar el techo de forma deliberada; exportar la espera del pool como métrica
La primera consulta tras un rato de calma fallaEl servidor o un proxy de conexiones cierra las conexiones ociosaspool_pre_ping=True más un pool_recycle por debajo del tiempo de inactividad
La exportación nocturna tarda horasCarga perezosa de una relación dentro de un bucleselectinload() en la consulta y raiseload("*") para que siga así

Ninguno es exótico. Son los mismos cinco cada vez, y cuatro de los cinco son configuración, no código.

Qué gana con esto la parte de descarga

Una capa de almacenamiento construida así cambia lo que se le permite hacer al crawler. Como las escrituras son idempotentes, un worker puede caerse y reiniciarse sin conciliación. Como el upsert informa de qué cambió, el planificador puede rastrear más a menudo las páginas que cambian rápido y menos las tranquilas, que es la mayor palanca sobre el gasto en proxies y encaja directamente con las prácticas de cómo scrapear sin que te bloqueen. Y como el feed de cambios existe, los trabajos posteriores como la monitorización de precios consumen un flujo en lugar de releer una tabla.

La etapa de descarga se lleva la atención porque ahí están los bloqueos y las facturas. La etapa de almacenamiento es donde decides cuánta descarga tendrás que pagar.

Preguntas frecuentes

Ambos, en rutas distintas. Define el esquema con modelos declarativos del ORM para que las anotaciones de tipo de Python y el DDL no se separen, y escribe los lotes con la construcción insert() de estilo Core, a la que SQLAlchemy 2.0 permite pasar una lista de diccionarios. Así los modelos siguen tipados y lo que ejecuta la base de datos sigue siendo trabajo por conjuntos.
Volver a la guía completa

Artículos relacionados