Aller au contenu principal

SQLAlchemy pour le web scraping : une couche de stockage de production

HT

Hinata Tomoda

Ingénieur web & testeur indépendant

20 min de lecture

Réponse courte : dans un pipeline de scraping, SQLAlchemy n'est pas la partie qui récupère — c'est la partie qui rend inutile une nouvelle récupération. Modélisez l'enregistrement dont vous avez réellement besoin, bornez le pool de connexions sous la limite de la base elle-même, écrivez chaque ligne via un ON CONFLICT DO UPDATE pour qu'une reprise ne coûte rien, et rejouez la transaction entière plutôt que l'instruction. Tout ce qui suit vise SQLAlchemy 2.0 — 2.0.52, publiée le 11 août 2026 — et cite les valeurs par défaut officielles au lieu de les deviner.

À retenir

  • La couche de stockage est d'abord un levier de coût : le trafic proxy se facture au gigaoctet, donc la requête la moins chère est celle à laquelle votre base a déjà répondu.
  • Dans SQLAlchemy 2.0 l'annotation de type EST le schéma — Mapped[str] donne NOT NULL et Mapped[str | None] donne NULL — le vérificateur de types et le DDL ne peuvent donc pas diverger.
  • Les valeurs par défaut documentées du pool sont pool_size 5, max_overflow 10, pool_timeout 30, pool_recycle -1 et pool_pre_ping False ; chaque processus worker supplémentaire multiplie les deux premières face à la limite de connexions de votre base.
  • L'idempotence vient de on_conflict_do_update() sur une clé naturelle, avec une clause WHERE qui saute les lignes dont l'empreinte de contenu n'a pas bougé — un nouveau crawl sans changement n'écrit alors rien du tout.
  • En asyncio, deux règles : une AsyncSession par tâche et expire_on_commit=False ; une session partagée ou un chargement paresseux accidentel est l'erreur MissingGreenlet que vous déboguerez sinon à 2 heures du matin.
  • Quand une connexion tombe, la transaction disparaît avec elle — la reprise documentée consiste à rejouer l'opération depuis le début, ce qui n'est sûr que parce que les écritures sont idempotentes.

Où se place la base de données dans un pipeline de scraping

Un pipeline de production comporte cinq étapes — collecte d'URL, récupération, parsing, validation et stockage — et notre guide du web scraping parcourt toute la chaîne. Une seule de ces étapes coûte de l'argent à l'octet. Le trafic des proxys résidentiels se facture au gigaoctet : l'économie du pipeline se décide donc sur la fréquence à laquelle vous devez récupérer une page une seconde fois, et notre estimation des coûts du web scraping met des chiffres concrets là-dessus.

Cela recadre la couche de stockage. Son rôle n'est pas de « ranger les lignes quelque part ». Son rôle est de répondre à trois questions assez bon marché pour que le récupérateur n'ait jamais à le faire :

  1. Ai-je déjà vu cet enregistrement, et a-t-il changé ? Si la réponse est non, le crawl économise une requête facturée.
  2. Que survit à un crash ? Un worker qui meurt au milieu d'un lot ne doit pas laisser la moitié d'une page validée et l'autre perdue.
  3. Les analystes peuvent-ils lire pendant que le crawler écrit ? Une transaction interminable qui tient des verrous transforme l'entrepôt et le crawler en adversaires.

SQLAlchemy répond aux trois, en deux couches qu'il vaut mieux ne pas confondre. Core est le langage d'expression SQL : Table, select(), insert(), le moteur et le pool de connexions. L'ORM ajoute par-dessus les classes mappées, une identity map et la Session façon unité de travail. Le partage productif pour un scraper : définir le schéma avec des modèles déclaratifs ORM — les types Python et le DDL ont alors une seule source de vérité — et écrire les lots avec des instructions de style Core, car insérer dix mille lignes est une opération ensembliste, pas dix mille mutations d'objets.

Modélisez l'enregistrement, pas la page que vous avez scrapée

L'erreur de modélisation la plus coûteuse dans un scraper est de stocker l'artefact au lieu du fait. Le HTML brut est énorme, change pour des raisons qui ne vous intéressent pas, et ne s'interroge pas. Stockez l'enregistrement extrait, plus juste assez de métadonnées pour décider s'il est nouveau.

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")

Trois détails y font l'essentiel du travail.

Mapped[] décide de la nullabilité. La règle documentée est que Mapped[str] produit NOT NULL et que Mapped[Optional[str]] — ou Mapped[str | None] — produit NULL, primary_key=True impliquant NOT NULL quoi qu'il arrive. Vous pouvez toujours forcer avec mapped_column(nullable=...), mais par défaut votre vérificateur de types lit la même vérité que votre schéma.

type_annotation_map centralise la correspondance Python vers SQL. SQLAlchemy mappe déjà int sur Integer, str sur String, datetime.datetime sur DateTime, uuid.UUID sur Uuid, et ainsi de suite. Surcharger cette table sur la classe de base est la façon d'appliquer une politique valable pour tout le projet — horodatages avec fuseau, JSONB plutôt que JSON générique — sans la répéter sur 40 colonnes.

La clé naturelle est une vraie contrainte, pas une convention. source plus external_id, c'est ce qui rend la même fiche produit reconnaissable d'un crawl à l'autre. Un id de substitution seul ne peut pas l'exprimer, et tout le chapitre suivant repose sur la capacité de la base à détecter elle-même la collision.

Le moteur et le pool sont votre budget de concurrence

create_engine() est une fabrique, pas une connexion : elle détient un pool de connexions créé paresseusement et partagé. Hors SQLite en mémoire, le pool par défaut est QueuePool, et ses valeurs par défaut documentées sont les chiffres qui vous mordront en premier.

ParamètreDéfautCe qu'il décide réellement
pool_size5Connexions maintenues ouvertes par moteur et par processus
max_overflow10Connexions supplémentaires en pointe, jetées après usage
pool_timeout30Secondes d'attente d'une connexion libre avant erreur
pool_recycle-1 (désactivé)Âge en secondes au-delà duquel la connexion est remplacée
pool_pre_pingFalseTester la vivacité par une requête légère au retrait
insertmanyvalues_page_size1000Lignes par INSERT multi-lignes groupé

Le calcul qui compte est workers × (pool_size + max_overflow) ≤ limite de connexions − marge. Avec les valeurs par défaut, huit processus de crawl peuvent exiger 120 connexions d'une base qui n'en autorise peut-être que 100, et la panne se manifeste par des pool_timeout qui expirent côté workers pendant qu'une migration ne se connecte même plus. Décidez le plafond délibérément :

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 est la stratégie pessimiste documentée face aux déconnexions : SQLAlchemy envoie au retrait un ping propre au dialecte (typiquement SELECT 1) et, en cas d'échec, jette cette connexion et invalide toutes les connexions du pool plus anciennes que l'instant courant. Cela coûte un aller-retour par retrait et supprime toute la classe « la première requête après une période d'inactivité échoue ». pool_recycle est le complément pour les moteurs qui ferment les connexions inactives au chronomètre — la documentation le cite comme remède immédiat au server has gone away de MySQL, et le raisonnement vaut pour tout PostgreSQL managé derrière un proxy de connexions.

Le fork est l'endroit où les pools cassent

Les crawlers multiprocessus rencontrent un danger précis et documenté : les connexions du pool ne sont pas partagées avec un processus forké. Deux processus finissent par écrire dans la même socket, et le symptôme est un état de protocole corrompu plutôt qu'une erreur propre. Le remède officiel est de disposer du pool hérité dans l'enfant, sans fermer les sockets du parent :

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)

Une session par unité de travail

La Session est documentée comme « a mutable, stateful object that represents a single database transaction », et elle « cannot be shared among concurrent threads or asyncio tasks without careful synchronization ». Lisez cela comme une règle de conception plutôt qu'un avertissement : la durée de vie d'une session est celle d'une transaction, qui est la durée pendant laquelle vous tenez des verrous.

Créez le sessionmaker une fois au niveau module et ouvrez une session par unité de travail :

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() en gestionnaire de contexte valide en cas de succès et annule à la moindre exception : aucun chemin ne laisse s'échapper un lot à moitié écrit. Des lots de 200 à 1 000 lignes environ sont le point d'équilibre habituel : assez gros pour que les allers-retours cessent de dominer, assez petits pour qu'une annulation reste bon marché et que la durée de verrouillage se compte en dizaines de millisecondes.

Rendez chaque écriture idempotente

C'est le chapitre qui justifie l'article. Un crawler est un système qui rejouera : reprises, rattrapages, un opérateur qui relance le job d'hier. Si écrire deux fois la même page produit des doublons ou écrase de bonnes données, chacun de ces événements devient un incident. INSERT ... ON CONFLICT DO UPDATE déplace la déduplication dans la base, où elle est atomique.

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,
    )

Notez le chemin d'import : l'upsert vit sur l'insert propre au dialecte, insert de sqlalchemy.dialects.postgresql, pas sur le générique. SQLite propose la même écriture on_conflict_do_update(), MySQL et MariaDB proposent on_duplicate_key_update(), et aucune construction portable ne couvre les trois. C'est une raison sérieuse de développer sur le moteur que vous déployez.

La clause where mérite une minute d'honnêteté sur son compromis. Sauter la mise à jour signifie que last_changed_at enregistre le moment où le contenu a changé pour la dernière fois — exactement ce que veut un pipeline de détection de changement — mais vous n'enregistrez plus quand vous avez vérifié pour la dernière fois. S'il vous faut les deux, gardez la mise à jour protégée pour l'enregistrement lui-même et écrivez l'horodatage « vu à » bon marché dans une table étroite séparée, que personne ne joint.

Le bénéfice est un flux de changements offert. Comme les lignes inchangées sont sautées, RETURNING renvoie exactement les lignes qui ont bougé :

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)

C'est la différence entre « on réindexe tout chaque nuit » et « on réindexe les 0,4 % qui ont changé ». Sur un corpus d'un million de pages, c'est la différence entre un job de quelques minutes et un job qui tourne toute la nuit.

Écritures en masse : laissez insertmanyvalues grouper

SQLAlchemy 2.0 accepte une liste de dictionnaires comme jeu de paramètres d'un insert(), et l'ORM interprète les clés comme des noms d'attributs plutôt que des noms de colonnes — un changement délibéré de la 2.0, qui compte dès qu'un attribut mappé et sa colonne s'écrivent différemment.

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,
)

Dessous, la fonctionnalité insertmanyvalues réécrit cela en instructions INSERT ... VALUES multi-lignes groupées. Elle est active par défaut pour PostgreSQL, MySQL, SQLite, SQL Server et Oracle, et la taille de lot suit insertmanyvalues_page_size, qui « defaults to 1000, but may also be subject to dialect-specific limiting factors ». Baissez-la quand les lignes portent de gros contenus JSON : le gain vient de la réduction des allers-retours, et une instruction dépassant les limites de paramètres du serveur vous le reprend aussitôt.

Quand vous avez besoin des clés primaires générées, demandez-les dans l'ordre des paramètres :

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()

Deux choses à ne pas faire. Ne bouclez pas session.add() sur dix mille objets pour valider une fois : vous payez toute la comptabilité de l'unité de travail pour des lignes que personne ne mutera. Et ne construisez pas le SQL à la main pour « éviter le surcoût de l'ORM » ; c'est le chemin paramétré qui tient le texte scrapé — entrée non fiable par définition — hors de votre grammaire SQL.

Pipelines asynchrones : une AsyncSession par tâche

Si votre récupérateur est déjà en asyncio — la plupart des modernes le sont, et nos notes sur les proxys pour agents IA expliquent pourquoi cette forme s'est répandue — faire tourner la base sur la même boucle d'événements évite un saut de thread. Les règles sont étroites et non négociables.

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()

L'erreur que vous rencontrerez en cas d'écart s'appelle MissingGreenlet, et elle signale presque toujours un chargement paresseux déclenché hors du contexte asynchrone. La documentation asyncio propose trois remèdes, dans l'ordre où je les essaierais : charger en amont avec selectinload() dans la requête, ajouter le mixin AsyncAttrs à votre base et écrire await obj.awaitable_attrs.things, ou basculer dans session.run_sync() pour un bloc de code ORM synchrone ordinaire.

Une mise en garde sur le dimensionnement : l'asynchrone n'augmente pas la limite de connexions de votre base — il rend simplement bien plus facile de l'atteindre. Mille tâches concurrentes face à pool_size=10 ne créent pas d'interblocage, mais toutes les tâches au-delà de la dixième font la queue, et pool_timeout décide si cette file échoue bruyamment ou devient silencieusement votre latence.

Relire sans N+1

Le chemin de lecture est celui par lequel les données scrapées sont exportées, et celui où loge le bug de performance classique de l'ORM : itérer sur des parents en touchant une relation à chaque tour émet une requête par parent. La réponse de SQLAlchemy est d'énoncer la stratégie de chargement dans la requête.

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()

Pour une relation scalaire plusieurs-vers-un — une observation vers sa page — joinedload() est la meilleure forme, car la jointure ajoute un jeu de colonnes plutôt qu'un second aller-retour. Le seul piège que la documentation signale explicitement : joinedload() sur une collection multiplie les lignes, il faut donc dédupliquer avec .unique() avant de compter quoi que ce soit.

raiseload("*") est sous-estimé. Activez-le dans les tests et dans tout export par lot, et chaque chargement paresseux accidentel devient un échec bruyant et localisé pendant le développement, au lieu d'un job mystérieux d'une heure en production.

Observabilité : mesurez le SQL, pas l'impression

echo=True est un interrupteur de développement. En production, il journalise chaque instruction — y compris les paramètres liés, ce qui dans un scraper signifie du contenu scrapé et, sur la connexion elle-même, des identifiants. Utilisez plutôt les événements, et journalisez l'instruction sans les paramètres.

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 signature d'événement est figée — (conn, cursor, statement, parameters, context, executemany) pour les deux hooks — et la pile posée sur conn.info garde les exécutions imbriquées correctement appariées. Trois métriques méritent d'être exportées : la durée par type d'opération, la part d'executemany (elle monte quand votre groupage fonctionne), et le temps d'attente au retrait du pool, le chiffre qui vous dit si le crawler est lent à cause du site cible ou de votre propre pool_size.

Rejouez la transaction, jamais l'instruction

La documentation est catégorique : « 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 forme recommandée est de « 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))

Ce n'est sûr que parce que l'écriture est un upsert. Rejouer un INSERT nu après un échec ambigu, c'est fabriquer des doublons ; rejouer un ON CONFLICT DO UPDATE converge vers la même ligne quel que soit le nombre d'exécutions. L'idempotence n'est pas un confort ici — c'est la condition qui rend la reprise automatique possible.

pool_pre_ping et ce décorateur résolvent deux moitiés différentes du problème. Le pre-ping attrape la connexion morte pendant qu'elle dormait dans le pool, avant que vous ayez commencé le travail. La reprise attrape la connexion qui meurt en pleine transaction, quand il n'y a plus rien à sauver.

Changements de schéma : autogenerate est un brouillon, pas une migration

Alembic est l'outil de migration du même projet, et sa propre documentation dit qu'autogenerate « is not intended to be perfect » et qu'« it is always necessary to manually review and correct the candidate migrations that autogenerate produces ». Prenez-le au pied de la lettre.

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

Il détecte de façon fiable l'ajout et la suppression de tables et de colonnes, les changements de nullabilité, les changements simples d'index et de contraintes d'unicité nommées, les changements simples de clés étrangères et les contraintes CHECK nommées. Il ne détecte pas les renommages de table ni de colonne — les deux apparaissent comme une suppression suivie d'un ajout, ce qui sur un corpus scrapé revient à supprimer la colonne et la reconstruire vide — et il ne voit pas du tout les contraintes nommées anonymement. C'est pourquoi la naming_convention figure sur le MetaData dans le premier bloc de code : sans elle, autogenerate n'a rien pour rapprocher les contraintes.

Le déroulé qui tient : générer, lire le fichier ligne à ligne, réécrire les paires suppression-ajout en op.alter_column(..., new_column_name=...), et exécuter la migration sur une copie restaurée de la production avant qu'elle n'approche la production.

Tests : annuler plutôt que nettoyer

Les tests d'un scraper ont besoin d'une base qui se comporte comme la vraie et d'une fixture qui ne laisse pas de résidu. SQLAlchemy 2.0 documente le motif : ouvrir une transaction externe sur une connexion, lier la session à cette connexion en mode savepoint, et annuler la transaction externe à la fin du test.

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()

Exécutez-les sur le même moteur que celui que vous déployez. SQLite est une cible correcte pour les parties purement Python, mais la sémantique d'ON CONFLICT, les opérateurs JSONB et l'ordre des NULL diffèrent tous, et une suite qui passe sur SQLite alors que la production tourne sur Postgres teste un autre programme.

Cinq pannes que je retrouve dans les bases de scrapers

SymptômeCause racineCorrection
Lignes en double après chaque reprisePas de contrainte d'unicité sur la clé naturelle, ON CONFLICT n'a rien à détecterAjouter la contrainte d'abord, puis l'upsert ; dédupliquer l'historique une fois
La table gonfle mais le nombre de lignes stagneUn upsert sans garde réécrit chaque ligne à chaque crawlAjouter la garde content_hash à la clause where
Les workers bloquent, la base semble oisivepool_size × workers dépasse la limite serveur ; tout le monde attend dans pool_timeoutFixer le plafond délibérément ; exporter l'attente au retrait comme métrique
La première requête après une accalmie échoueConnexions inactives fermées par le serveur ou un proxy de connexionspool_pre_ping=True plus un pool_recycle sous le délai d'inactivité
L'export nocturne dure des heuresChargement paresseux d'une relation dans une boucleselectinload() dans la requête, raiseload("*") pour que ça le reste

Rien d'exotique là-dedans. Ce sont les mêmes cinq à chaque fois, et quatre sur cinq relèvent de la configuration, pas du code.

Ce que cela apporte au côté récupération

Une couche de stockage bâtie ainsi change ce que le crawler a le droit de faire. Comme les écritures sont idempotentes, un worker peut planter et redémarrer sans réconciliation. Comme l'upsert signale ce qui a changé, l'ordonnanceur peut crawler plus souvent les pages qui bougent et plus rarement les autres — c'est le plus grand levier sur la dépense proxy, et il s'articule directement avec les pratiques de scraper sans se faire bloquer. Et parce que le flux de changements existe, les jobs en aval comme la surveillance des prix consomment un flux au lieu de relire une table.

C'est l'étape de récupération qui attire l'attention, parce que c'est là que sont les blocages et les factures. L'étape de stockage est celle où vous décidez combien de récupération vous aurez à payer.

Questions fréquentes

Les deux, sur des chemins distincts. Définissez le schéma avec les modèles déclaratifs de l'ORM pour que les annotations de type Python et le DDL ne divergent jamais, puis écrivez les lots avec la construction insert() de style Core, à laquelle SQLAlchemy 2.0 permet de passer une liste de dictionnaires. Les modèles restent typés et ce que la base exécute reste ensembliste.
Retour au guide complet

Articles associés