Files
colibre/docs/superpowers/plans/2026-04-15-duckdb-migration.md
T
Colin Maudry 035d7f23fa Plan de migration DuckDB — découpage en 20 tâches
20 étapes bite-sized couvrant : dépendance duckdb + .gitignore,
should_rebuild (TDD), build_database avec transforms Polars et verrou
fcntl, startup guard + query_marches, migration page-par-page
(marche → acheteur → titulaire → arbre → tableau → observatoire →
figures), suppression des globaux Polars, mesure RSS avant/après.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-15 13:37:49 +02:00

47 KiB
Raw Blame History

DuckDB Migration Implementation Plan

For agentic workers: REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (- [ ]) syntax for tracking.

Goal: Replace the in-memory global Polars dataframes in src/utils.py (the full df plus five derived frames) with an on-disk DuckDB database, reducing steady-state RSS. Per-request SQL returns small Polars frames that existing downstream code continues to consume.

Architecture: A new module src/db.py owns the DuckDB lifecycle — startup build under a fcntl.flock lock, atomic os.replace swap, read-only runtime connection with per-call cursors, and a query_marches() helper that returns pl.DataFrame. Row-level transforms stay in Polars (via connection.register("frame", pl_frame)) so booleans_to_strings and the null-name replacement remain the single source of truth. df_acheteurs and df_titulaires stay as in-memory Polars frames populated from DuckDB at import time (they feed the autocomplete search).

Tech Stack: Python 3.10+, DuckDB (Python API), Polars, Dash, pytest (+ dash Selenium testing), pre-commit (prettier, ruff).

Spec: docs/superpowers/specs/2026-04-15-duckdb-migration-design.md


File Structure

Created files

Path Responsibility
src/db.py DuckDB lifecycle: should_rebuild, build_database, module-level conn, schema, get_cursor(), query_marches()
tests/test_db.py Unit tests for should_rebuild, build_database, query_marches, and concurrent-startup lock behavior

Modified files

Path What changes
pyproject.toml Add duckdb dependency
.gitignore Ignore **/decp.duckdb, **/decp.duckdb.tmp, **/decp.duckdb.lock
src/utils.py Remove df global, replace df_acheteurs / df_titulaires population with DuckDB-backed queries, remove df_*_marches and df_*_departement globals, expose schema / columns from src/db.py
src/pages/marche.py df.lazy().filter(...)query_marches("uid = ?", (uid,)).lazy()
src/pages/acheteur.py Same pattern + replace df.collect_schema() with module schema, df.columns with schema.names()
src/pages/titulaire.py Same pattern as acheteur
src/pages/tableau.py df.lazy()query_marches().lazy(), df.columnsschema.names(), df.widthlen(schema.names())
src/pages/observatoire.py 3 call sites: df.lazy()query_marches().lazy(), df.columnsschema.names()
src/pages/arbre/departement.py df_acheteurs_departement / df_titulaires_departementget_cursor().execute("SELECT ... FROM {table} WHERE ...", [...]).pl()
src/pages/arbre/liste_marches_org.py df_acheteurs_marches / df_titulaires_marchesget_cursor().execute(...).pl()
src/figures.py df.columnsschema.names() (via src.db)
tests/conftest.py Ensure the DuckDB test file is rebuilt from the freshly-written test.parquet each session

Task 1: Add DuckDB dependency and ignore patterns

Files:

  • Modify: pyproject.toml (dependencies list at top)

  • Modify: .gitignore

  • Step 1: Add duckdb to pyproject.toml

Open pyproject.toml, find the dependencies = [...] block, and insert "duckdb", on a new line before "flask-caching",. Result:

dependencies = [
  "dash==3.4.0",
  "dash[compress]",
  "polars",
  "gunicorn",
  "dash-bootstrap-components",
  "python-dotenv",
  "xlsxwriter",
  "plotly[express]",
  "httpx",
  "pandas",
  "unidecode",
  "dash-leaflet",
  "dash-extensions",
  "duckdb",
  "flask-caching",
]
  • Step 2: Install the dependency

Run: uv pip install -e ".[dev]"

Expected: installation completes, python -c "import duckdb; print(duckdb.__version__)" prints a version.

  • Step 3: Add DuckDB artifacts to .gitignore

Append to .gitignore:

# DuckDB runtime artifacts (regenerated from decp_prod.parquet at startup)
**/decp.duckdb
**/decp.duckdb.tmp
**/decp.duckdb.lock
  • Step 4: Commit
git add pyproject.toml .gitignore uv.lock
git commit -m "Ajout de la dépendance duckdb"

Task 2: Write failing tests for should_rebuild

Files:

  • Create: tests/test_db.py

  • Step 1: Write the failing tests

Create tests/test_db.py:

import os
import time
from pathlib import Path

import pytest


@pytest.fixture
def parquet_and_db(tmp_path, monkeypatch):
    parquet = tmp_path / "source.parquet"
    db = tmp_path / "decp.duckdb"
    parquet.write_bytes(b"fake parquet content")
    monkeypatch.delenv("REBUILD_DUCKDB", raising=False)
    monkeypatch.delenv("DEVELOPMENT", raising=False)
    return parquet, db


def test_should_rebuild_when_db_missing(parquet_and_db, monkeypatch):
    from src.db import should_rebuild
    parquet, db = parquet_and_db
    monkeypatch.setenv("DEVELOPMENT", "true")
    assert should_rebuild(db, parquet) is True


def test_should_rebuild_prod_when_parquet_newer(parquet_and_db, monkeypatch):
    from src.db import should_rebuild
    parquet, db = parquet_and_db
    db.write_bytes(b"x")
    time.sleep(0.01)
    parquet.touch()
    monkeypatch.setenv("DEVELOPMENT", "false")
    assert should_rebuild(db, parquet) is True


def test_should_not_rebuild_prod_when_parquet_older(parquet_and_db, monkeypatch):
    from src.db import should_rebuild
    parquet, db = parquet_and_db
    parquet.touch()
    time.sleep(0.01)
    db.write_bytes(b"x")
    monkeypatch.setenv("DEVELOPMENT", "false")
    assert should_rebuild(db, parquet) is False


def test_should_not_rebuild_dev_even_when_parquet_newer(parquet_and_db, monkeypatch):
    from src.db import should_rebuild
    parquet, db = parquet_and_db
    db.write_bytes(b"x")
    time.sleep(0.01)
    parquet.touch()
    monkeypatch.setenv("DEVELOPMENT", "true")
    monkeypatch.delenv("REBUILD_DUCKDB", raising=False)
    assert should_rebuild(db, parquet) is False


def test_should_rebuild_dev_when_rebuild_forced(parquet_and_db, monkeypatch):
    from src.db import should_rebuild
    parquet, db = parquet_and_db
    db.write_bytes(b"x")
    time.sleep(0.01)
    parquet.touch()
    monkeypatch.setenv("DEVELOPMENT", "true")
    monkeypatch.setenv("REBUILD_DUCKDB", "true")
    assert should_rebuild(db, parquet) is True
  • Step 2: Run tests and verify they fail

Run: uv run pytest tests/test_db.py -v

Expected: FAIL with ModuleNotFoundError: No module named 'src.db' (or ImportError).


Task 3: Implement should_rebuild in src/db.py

Files:

  • Create: src/db.py

  • Step 1: Create minimal src/db.py

import logging
import os
from pathlib import Path

logger = logging.getLogger("decp.info")


def should_rebuild(db_path: Path, parquet_path: Path) -> bool:
    """Decide whether to rebuild the DuckDB database from the source Parquet.

    Rules:
      - Rebuild if the DuckDB file does not exist.
      - Otherwise, rebuild only if the source Parquet is newer than the DB,
        EXCEPT in development mode without REBUILD_DUCKDB=true (dev keeps
        a stable DB across reloads unless explicitly opted in).
    """
    db_path = Path(db_path)
    parquet_path = Path(parquet_path)

    if not db_path.exists():
        return True

    dev = os.getenv("DEVELOPMENT", "False").lower() == "true"
    force = os.getenv("REBUILD_DUCKDB", "False").lower() == "true"
    if dev and not force:
        return False

    return parquet_path.stat().st_mtime > db_path.stat().st_mtime
  • Step 2: Run tests and verify they pass

Run: uv run pytest tests/test_db.py -v

Expected: 5 passed.

  • Step 3: Commit
git add src/db.py tests/test_db.py
git commit -m "Ajout de src.db.should_rebuild et ses tests"

Task 4: Write failing tests for build_database and query_marches

Files:

  • Modify: tests/test_db.py

  • Step 1: Append build/query tests to tests/test_db.py

import datetime
import polars as pl


@pytest.fixture
def built_db(tmp_path, monkeypatch):
    """Build a DuckDB from a small Polars frame written as parquet."""
    parquet_path = tmp_path / "source.parquet"
    db_path = tmp_path / "decp.duckdb"

    data = pl.DataFrame(
        [
            {
                "uid": "1",
                "id": "1",
                "objet": "Travaux",
                "acheteur_id": "A1",
                "acheteur_nom": "Mairie",
                "acheteur_departement_code": "75",
                "titulaire_id": "T1",
                "titulaire_nom": "Entreprise",
                "titulaire_departement_code": "35",
                "titulaire_typeIdentifiant": "SIRET",
                "montant": 1000.0,
                "dateNotification": datetime.date(2025, 1, 1),
                "donneesActuelles": True,
                "marcheInnovant": True,
            },
            {
                "uid": "2",
                "id": "2",
                "objet": "Études",
                "acheteur_id": "A1",
                "acheteur_nom": "Mairie",
                "acheteur_departement_code": "75",
                "titulaire_id": "T2",
                "titulaire_nom": None,
                "titulaire_departement_code": "75",
                "titulaire_typeIdentifiant": "SIRET",
                "montant": 500.0,
                "dateNotification": datetime.date(2024, 6, 1),
                "donneesActuelles": True,
                "marcheInnovant": False,
            },
            {
                "uid": "3",
                "id": "3",
                "objet": "Ancien",
                "acheteur_id": "A2",
                "acheteur_nom": None,
                "acheteur_departement_code": "13",
                "titulaire_id": "T3",
                "titulaire_nom": "Autre",
                "titulaire_departement_code": "13",
                "titulaire_typeIdentifiant": "SIRET",
                "montant": 100.0,
                "dateNotification": datetime.date(2023, 1, 1),
                "donneesActuelles": False,  # must be filtered out
                "marcheInnovant": False,
            },
        ]
    )
    data.write_parquet(parquet_path)
    monkeypatch.setenv("DATA_FILE_PARQUET_PATH", str(parquet_path))

    from src.db import build_database
    build_database(db_path, parquet_path)
    return db_path


def test_build_filters_donnees_actuelles(built_db):
    import duckdb
    with duckdb.connect(str(built_db), read_only=True) as c:
        rows = c.execute("SELECT uid FROM decp ORDER BY uid").fetchall()
    assert [r[0] for r in rows] == ["1", "2"]


def test_build_converts_booleans_to_oui_non(built_db):
    import duckdb
    with duckdb.connect(str(built_db), read_only=True) as c:
        values = c.execute(
            "SELECT marcheInnovant FROM decp ORDER BY uid"
        ).fetchall()
    assert [v[0] for v in values] == ["oui", "non"]


def test_build_replaces_null_org_names(built_db):
    import duckdb
    with duckdb.connect(str(built_db), read_only=True) as c:
        titulaire_2 = c.execute(
            "SELECT titulaire_nom FROM decp WHERE uid = '2'"
        ).fetchone()
    assert titulaire_2[0] == "[Identifiant non reconnu dans la base INSEE]"


def test_build_creates_derived_tables(built_db):
    import duckdb
    with duckdb.connect(str(built_db), read_only=True) as c:
        tables = {r[0] for r in c.execute("SHOW TABLES").fetchall()}
    assert {"decp", "acheteurs_marches", "titulaires_marches",
            "acheteurs_departement", "titulaires_departement"} <= tables


def test_query_marches_returns_polars_frame(built_db, monkeypatch):
    monkeypatch.setenv("DATA_FILE_PARQUET_PATH", str(built_db.parent / "source.parquet"))
    # Force src.db to load pointing at this test DB.
    import importlib, src.db
    importlib.reload(src.db)
    from src.db import query_marches
    frame = query_marches("acheteur_id = ?", ("A1",))
    assert isinstance(frame, pl.DataFrame)
    assert frame.height == 2
    assert set(frame["uid"].to_list()) == {"1", "2"}
  • Step 2: Run new tests and verify they fail

Run: uv run pytest tests/test_db.py -v -k "build or query_marches"

Expected: FAIL with ImportError: cannot import name 'build_database' from 'src.db' (or similar on query_marches).


Task 5: Implement build_database in src/db.py

Files:

  • Modify: src/db.py

  • Step 1: Add imports and build_database

Replace the contents of src/db.py with:

import fcntl
import logging
import os
from pathlib import Path

import duckdb
import polars as pl
import polars.selectors as cs
from polars.exceptions import ComputeError
from time import sleep

logger = logging.getLogger("decp.info")


def should_rebuild(db_path: Path, parquet_path: Path) -> bool:
    db_path = Path(db_path)
    parquet_path = Path(parquet_path)
    if not db_path.exists():
        return True
    dev = os.getenv("DEVELOPMENT", "False").lower() == "true"
    force = os.getenv("REBUILD_DUCKDB", "False").lower() == "true"
    if dev and not force:
        return False
    return parquet_path.stat().st_mtime > db_path.stat().st_mtime


def _load_source_frame(parquet_path: Path) -> pl.DataFrame:
    """Read the source parquet and apply the row-level transforms.

    Kept here (not in utils.py) so src.db has no dependency on utils.
    Mirrors the behavior previously in utils.get_decp_data().
    """
    try:
        lff: pl.LazyFrame = pl.scan_parquet(str(parquet_path))
    except ComputeError:
        logger.info("Lecture du parquet échouée, nouvelle tentative dans 10s...")
        sleep(10)
        lff = pl.scan_parquet(str(parquet_path))

    lff = lff.sort(by=["dateNotification", "uid"], descending=True, nulls_last=True)
    lff = lff.filter(pl.col("donneesActuelles")).drop("donneesActuelles")

    # booleans_to_strings: true → "oui", false → "non"
    lff = lff.with_columns(
        pl.col(cs.Boolean).cast(pl.String).str.replace("true", "oui").str.replace("false", "non")
    )

    for col in ["acheteur_nom", "titulaire_nom"]:
        lff = lff.with_columns(
            pl.when(pl.col(col).is_null())
            .then(pl.lit("[Identifiant non reconnu dans la base INSEE]"))
            .otherwise(pl.col(col))
            .name.keep()
        )

    return lff.collect()


def build_database(db_path: Path, parquet_path: Path) -> None:
    """Build the DuckDB database atomically under an exclusive lock.

    Caller MUST hold the fcntl.flock on the .lock file.
    """
    db_path = Path(db_path)
    parquet_path = Path(parquet_path)
    tmp_path = db_path.with_suffix(".duckdb.tmp")
    if tmp_path.exists():
        tmp_path.unlink()

    logger.info(f"Construction de la base DuckDB à partir de {parquet_path}...")
    frame = _load_source_frame(parquet_path)

    with duckdb.connect(str(tmp_path)) as w:
        w.register("frame", frame)
        w.execute("CREATE TABLE decp AS SELECT * FROM frame")
        w.execute(
            "CREATE TABLE acheteurs_marches AS "
            "SELECT DISTINCT uid, objet, acheteur_id FROM decp "
            "ORDER BY acheteur_id"
        )
        w.execute(
            "CREATE TABLE titulaires_marches AS "
            "SELECT DISTINCT uid, objet, titulaire_id FROM decp "
            "ORDER BY titulaire_id"
        )
        w.execute(
            "CREATE TABLE acheteurs_departement AS "
            "SELECT DISTINCT acheteur_id, acheteur_nom, acheteur_departement_code "
            "FROM decp ORDER BY acheteur_nom"
        )
        w.execute(
            "CREATE TABLE titulaires_departement AS "
            "SELECT DISTINCT titulaire_id, titulaire_nom, titulaire_departement_code "
            "FROM decp ORDER BY titulaire_nom"
        )

    os.replace(tmp_path, db_path)
    logger.info(f"Base DuckDB construite : {db_path}")
  • Step 2: Run build tests and verify they pass

Run: uv run pytest tests/test_db.py -v -k "build"

Expected: 4 passed (test_build_filters_donnees_actuelles, test_build_converts_booleans_to_oui_non, test_build_replaces_null_org_names, test_build_creates_derived_tables).


Task 6: Implement startup guard, read-only connection, schema, get_cursor, query_marches

Files:

  • Modify: src/db.py

  • Step 1: Append startup logic and query helpers

Add at the bottom of src/db.py:

def _resolve_db_path() -> Path:
    parquet = os.getenv("DATA_FILE_PARQUET_PATH")
    if not parquet:
        raise RuntimeError("DATA_FILE_PARQUET_PATH is not set")
    return Path(parquet).parent / "decp.duckdb"


def _ensure_database() -> Path:
    db_path = _resolve_db_path()
    parquet_path = Path(os.getenv("DATA_FILE_PARQUET_PATH"))
    lock_path = db_path.with_suffix(".duckdb.lock")

    with open(lock_path, "w") as lock_fd:
        fcntl.flock(lock_fd, fcntl.LOCK_EX)
        if should_rebuild(db_path, parquet_path):
            build_database(db_path, parquet_path)
    return db_path


DB_PATH = _ensure_database()
conn: duckdb.DuckDBPyConnection = duckdb.connect(str(DB_PATH), read_only=True)
schema: pl.Schema = conn.execute("SELECT * FROM decp LIMIT 0").pl().schema


def get_cursor() -> duckdb.DuckDBPyConnection:
    """Return a per-request cursor that shares the process-wide connection."""
    return conn.cursor()


def query_marches(
    where_sql: str = "TRUE",
    params: tuple = (),
    columns: list[str] | None = None,
    order_by: str | None = None,
    limit: int | None = None,
) -> pl.DataFrame:
    """Run a parameterized SELECT against the decp table and return Polars.

    `where_sql` and `order_by` are trusted SQL fragments (callers are internal
    code, never user input). `params` values are passed through DuckDB's
    parameter binding.
    """
    cols = ", ".join(columns) if columns else "*"
    sql = f"SELECT {cols} FROM decp WHERE {where_sql}"
    if order_by:
        sql += f" ORDER BY {order_by}"
    if limit is not None:
        sql += f" LIMIT {int(limit)}"
    return get_cursor().execute(sql, list(params)).pl()
  • Step 2: Run the full test_db.py suite

Run: uv run pytest tests/test_db.py -v

Expected: all tests pass (9 total).

  • Step 3: Commit
git add src/db.py tests/test_db.py
git commit -m "Implémentation de build_database, query_marches et connexion DuckDB"

Task 7: Write failing test for concurrent build locking

Files:

  • Modify: tests/test_db.py

  • Step 1: Append a lock-serialization test

import threading


def test_concurrent_build_serialized(tmp_path, monkeypatch):
    """Two workers starting at once must not stomp on each other's tmp file."""
    parquet_path = tmp_path / "source.parquet"
    db_path = tmp_path / "decp.duckdb"

    pl.DataFrame(
        [
            {
                "uid": "1", "id": "1", "objet": "x",
                "acheteur_id": "A", "acheteur_nom": "n",
                "titulaire_id": "T", "titulaire_nom": "n",
                "acheteur_departement_code": "75",
                "titulaire_departement_code": "75",
                "titulaire_typeIdentifiant": "SIRET",
                "montant": 1.0,
                "dateNotification": datetime.date(2025, 1, 1),
                "donneesActuelles": True,
                "marcheInnovant": False,
            }
        ]
    ).write_parquet(parquet_path)

    monkeypatch.setenv("DATA_FILE_PARQUET_PATH", str(parquet_path))

    from src.db import build_database, should_rebuild
    import fcntl

    errors: list[Exception] = []

    def worker():
        try:
            lock_path = db_path.with_suffix(".duckdb.lock")
            with open(lock_path, "w") as lock_fd:
                fcntl.flock(lock_fd, fcntl.LOCK_EX)
                if should_rebuild(db_path, parquet_path):
                    build_database(db_path, parquet_path)
        except Exception as exc:
            errors.append(exc)

    threads = [threading.Thread(target=worker) for _ in range(3)]
    for t in threads:
        t.start()
    for t in threads:
        t.join()

    assert errors == []
    assert db_path.exists()
    assert not db_path.with_suffix(".duckdb.tmp").exists()
  • Step 2: Run it

Run: uv run pytest tests/test_db.py::test_concurrent_build_serialized -v

Expected: PASS (the lock and should_rebuild short-circuit inside the worker already implements correct serialization — this test verifies the contract, no code changes needed).

If it fails, fix the issue in src/db.py (most likely: build_database does not clean up tmp_path when re-entered). Then re-run.

  • Step 3: Commit
git add tests/test_db.py
git commit -m "Test de sérialisation de la construction concurrente de la base"

Task 8: Make tests/conftest.py rebuild the DuckDB with the test parquet

Files:

  • Modify: tests/conftest.py

  • Modify: pyproject.toml (pytest env block)

  • Step 1: Force a fresh DuckDB rebuild in the test session

The existing test_data fixture writes tests/test.parquet at session start. We must delete any leftover tests/decp.duckdb before src.db is imported by the app, because src.db is imported at module load by src/utils.py, which is in turn imported by pages.

Update tests/conftest.py:

import datetime
import os
from pathlib import Path

import polars as pl
import pytest
from selenium.webdriver.chrome.options import Options


@pytest.fixture(scope="session", autouse=True)
def test_data():
    data = [
        {
            "uid": "1",
            "id": "1",
            "acheteur_nom": "ACHETEUR 1",
            "acheteur_id": "123",
            "titulaire_nom": "TITULAIRE 1",
            "titulaire_id": "345",
            "montant": 10,
            "dateNotification": datetime.date(2025, 1, 1),
            "codeCPV": "71600000",
            "donneesActuelles": True,
            "acheteur_departement_code": "75",
            "acheteur_departement_nom": "Paris",
            "acheteur_commune_nom": "Paris",
            "titulaire_departement_code": "35",
            "titulaire_departement_nom": "Ille-et-Vilaine",
            "titulaire_commune_nom": "Rennes",
            "titulaire_distance": 10,
            "titulaire_typeIdentifiant": "SIRET",
            "objet": "Objet test",
            "dureeRestanteMois": 12,
            "lieuExecution_code": "75001",
            "sourceFile": "test.xml",
            "sourceDataset": "test_dataset",
            "datePublicationDonnees": datetime.date(2025, 1, 1),
            "considerationsSociales": "",
            "considerationsEnvironnementales": "",
            "type": "Marché",
            "acheteur_categorie": "Collectivité",
            "titulaire_categorie": "PME",
        }
    ]
    parquet_path = Path(os.path.abspath("tests/test.parquet"))
    db_path = parquet_path.parent / "decp.duckdb"
    print(f"Writing test data to: {parquet_path}")

    pl.DataFrame(data).write_parquet(parquet_path)

    # Remove any stale DuckDB from a previous run so src.db rebuilds from
    # the freshly-written parquet at import time.
    for artifact in (db_path, db_path.with_suffix(".duckdb.tmp")):
        if artifact.exists():
            artifact.unlink()

    yield str(parquet_path)


def pytest_setup_options():
    options = Options()
    options.add_argument("--window-size=1200,1200 ")
    options.add_experimental_option(
        "prefs",
        {
            "download.default_directory": "/home/colin/git/decp.info",
            "download.prompt_for_download": False,
            "download.directory_upgrade": True,
            "safebrowsing.enabled": True,
        },
    )
    return options
  • Step 2: Verify existing tests still import cleanly

Run: uv run pytest tests/test_db.py -v

Expected: all 10 DB tests still pass (the conftest changes only affect integration tests).

  • Step 3: Commit
git add tests/conftest.py
git commit -m "Suppression de la base DuckDB de test obsolète avant chaque session"

Task 9: Migrate src/utils.py — introduce src.db imports, keep old globals

Files:

  • Modify: src/utils.py

This task does not remove the old globals yet. It wires src.db in and re-points schema / columns so pages can start using them without a flag day.

  • Step 1: Import from src.db at the top of utils.py

After the from unidecode import unidecode line, add:

from src.db import conn as duckdb_conn  # noqa: F401  (exposed for convenience)
from src.db import get_cursor, query_marches, schema  # noqa: F401
  • Step 2: Replace the module-level assignments at the bottom

Find lines 891913 (the block starting with df: pl.DataFrame = get_decp_data()).

Replace the block with:

df: pl.DataFrame = get_decp_data()
# schema and columns now come from src.db; overwrite in case any local code
# still reads them directly from utils.
schema = schema  # re-exported from src.db
columns = schema.names()

df_acheteurs = get_org_data(df, "acheteur")
df_titulaires = get_org_data(df, "titulaire")
df_acheteurs_departement: pl.DataFrame = (
    df_acheteurs.select(["acheteur_id", "acheteur_nom", "acheteur_departement_code"])
    .unique()
    .sort("acheteur_nom")
)
df_titulaires_departement: pl.DataFrame = (
    df_titulaires.select(
        ["titulaire_id", "titulaire_nom", "titulaire_departement_code"]
    )
    .unique()
    .sort("titulaire_nom")
)
df_acheteurs_marches: pl.DataFrame = (
    df.select("uid", "objet", "acheteur_id").unique().sort("acheteur_id")
)
df_titulaires_marches: pl.DataFrame = (
    df.select("uid", "objet", "titulaire_id").unique().sort("titulaire_id")
)

No removals — we keep the old globals alive. Pages that already use schema / columns now transparently get them from src.db.

  • Step 3: Run the Selenium test suite

Run: uv run pytest tests/test_main.py -v

Expected: all tests pass. The app loads; nothing has changed behaviorally yet.

  • Step 4: Commit
git add src/utils.py
git commit -m "Intégration de src.db dans utils (coexistence avec les globaux)"

Task 10: Migrate src/pages/marche.py (1 call site)

Files:

  • Modify: src/pages/marche.py

  • Step 1: Replace df with query_marches

Edit the import block:

from src.utils import (
    data_schema,
    format_values,
    make_org_jsonld,
    meta_content,
    unformat_montant,
)
from src.db import query_marches

Replace the body of get_marche_data:

def get_marche_data(url) -> tuple[dict, list]:
    marche_uid = url.split("/")[-1]

    # Filtre SQL côté DuckDB, puis Polars pour le post-traitement
    dff_marche = query_marches("uid = ?", (marche_uid,))
    if dff_marche.height == 0:
        return {}, []

    lff = dff_marche.lazy()
    dff_titulaires = lff.select(cs.starts_with("titulaire")).collect(engine="streaming")
    dff_marche_unique = lff.unique("uid").collect(engine="streaming")
    dff_marche_unique = format_values(dff_marche_unique)

    return dff_marche_unique.to_dicts()[0], dff_titulaires.to_dicts()
  • Step 2: Smoke-test

Run: uv run pytest tests/test_main.py -v

Expected: all tests pass.

Manually: uv run run.py, navigate to /marches/<any_uid_from_test_data> (e.g. /marches/1), confirm page renders with the buyer, amount, and titulaire list.

  • Step 3: Commit
git add src/pages/marche.py
git commit -m "marche.py : utilisation de query_marches au lieu du df global"

Task 11: Migrate src/pages/acheteur.py

Files:

  • Modify: src/pages/acheteur.py

  • Step 1: Replace df imports

Update the import block:

from src.utils import (
    columns,
    df_acheteurs,
    filter_table_data,
    format_number,
    get_annuaire_data,
    get_button_properties,
    get_default_hidden_columns,
    # ... any other existing non-df imports
)
from src.db import query_marches, schema

(Remove df from the src.utils import list.)

  • Step 2: Replace df.columns at line 73
columns=[{"id": col, "name": col} for col in schema.names()],
  • Step 3: Replace df.collect_schema() at line 303
dff = pl.DataFrame(schema=schema)
  • Step 4: Replace df.lazy().filter(...) in get_acheteur_marches_data
def get_acheteur_marches_data(url, ach_year: str) -> tuple:
    acheteur_siret = url.split("/")[-1]
    lff = query_marches("acheteur_id = ?", (acheteur_siret,)).lazy()
    if ach_year and ach_year != "Toutes les années":
        ach_year = int(ach_year)
        lff = lff.filter(pl.col("dateNotification").dt.year() == ach_year)
    lff = lff.sort(["dateNotification", "uid"], descending=True, nulls_last=True)
    dff: pl.DataFrame = lff.collect(engine="streaming")
    download_disabled, download_text, download_title = get_button_properties(dff.height)
    data = dff.to_dicts()
    return data, download_disabled, download_text, download_title
  • Step 5: Run tests and smoke-test

Run: uv run pytest tests/test_main.py -v

Manually: uv run run.py, navigate to /acheteurs/123 (test data SIRET), confirm KPIs, charts, and table render.

  • Step 6: Commit
git add src/pages/acheteur.py
git commit -m "acheteur.py : migration vers query_marches et schema depuis src.db"

Task 12: Migrate src/pages/titulaire.py

Files:

  • Modify: src/pages/titulaire.py

  • Step 1: Replace imports and call sites

Remove df from the src.utils import list. Add:

from src.db import query_marches, schema
  • Step 2: Replace df.columns at line 72
columns=[{"id": col, "name": col} for col in schema.names()],
  • Step 3: Replace df.lazy().filter(...) in get_titulaire_marches_data
def get_titulaire_marches_data(url, titulaire_year: str) -> tuple:
    titulaire_siret = url.split("/")[-1]
    lff = query_marches(
        "titulaire_id = ? AND titulaire_typeIdentifiant = 'SIRET'",
        (titulaire_siret,),
    ).lazy()
    if titulaire_year and titulaire_year != "Toutes les années":
        lff = lff.filter(
            pl.col("dateNotification").cast(pl.String).str.starts_with(titulaire_year)
        )
    lff = lff.sort(["dateNotification", "uid"], descending=True, nulls_last=True)
    lff = lff.fill_null("")
    dff: pl.DataFrame = lff.collect(engine="streaming")
    # ... rest of function unchanged
  • Step 4: Run tests and smoke-test

Run: uv run pytest tests/test_main.py -v

Manually: uv run run.py, navigate to /titulaires/345 (test SIRET), confirm it renders.

  • Step 5: Commit
git add src/pages/titulaire.py
git commit -m "titulaire.py : migration vers query_marches et schema depuis src.db"

Task 13: Migrate src/pages/arbre/departement.py

Files:

  • Modify: src/pages/arbre/departement.py

  • Step 1: Replace imports and queries

import polars as pl
from dash import Input, Output, callback, dcc, html, register_page

from src.utils import departements
from src.db import get_cursor

# ... (register_page and layout unchanged) ...


@callback(
    Output(component_id="departement_marches", component_property="children"),
    Input(component_id="departement_url", component_property="pathname"),
)
def departement_marches(url):
    departement = url.split("/")[-1]

    def make_link_list(org_type) -> list:
        table = (
            "acheteurs_departement" if org_type == "acheteur"
            else "titulaires_departement" if org_type == "titulaire"
            else None
        )
        if table is None:
            raise ValueError
        col_prefix = org_type
        rows = get_cursor().execute(
            f"SELECT {col_prefix}_id, {col_prefix}_nom "
            f"FROM {table} "
            f"WHERE {col_prefix}_departement_code = ? "
            f"ORDER BY {col_prefix}_nom",
            [departement],
        ).fetchall()

        link_list = []
        for org_id, org_nom in rows:
            li = html.Li(
                [
                    dcc.Link(
                        org_nom,
                        href=url + f"/{org_type}/{org_id}",
                        title=f"Marchés publics de {org_nom}",
                    ),
                    " ",
                    dcc.Link(
                        "(page dédiée)",
                        href=f"/{org_type}s/{org_id}",
                        title=f"Page dédiée aux marchés publics de {org_nom}",
                    ),
                ]
            )
            link_list.append(li)
        return link_list

    content = [
        html.H3("Acheteurs publics du département"),
        html.Ul(make_link_list("acheteur")),
        html.H3("Titulaires du département"),
        html.Ul(make_link_list("titulaire")),
    ]
    return content
  • Step 2: Run tests and smoke-test

Run: uv run pytest tests/test_main.py -v

Manually: uv run run.py, navigate to /departements/75, confirm list of acheteurs and titulaires renders.

  • Step 3: Commit
git add src/pages/arbre/departement.py
git commit -m "arbre/departement.py : requête DuckDB directe pour les listes départementales"

Task 14: Migrate src/pages/arbre/liste_marches_org.py

Files:

  • Modify: src/pages/arbre/liste_marches_org.py

  • Step 1: Replace imports and queries

import polars as pl
from dash import Input, Output, callback, dcc, html, register_page

from src.utils import df_acheteurs, df_titulaires
from src.db import get_cursor

name = "Liste des marchés publics"


def make_org_nom_verbe(org_type, org_id) -> tuple:
    if org_type == "titulaire":
        source = df_titulaires
        verbe = "remportés"
    elif org_type == "acheteur":
        source = df_acheteurs
        verbe = "attribués"
    else:
        raise ValueError

    org_nom = (
        source.filter(pl.col(f"{org_type}_id") == org_id)
        .select(f"{org_type}_nom")
        .item(0, 0)
    )
    return org_nom, verbe


# ... (register_page / layout unchanged) ...


@callback(
    Output(component_id="liste_marches", component_property="children"),
    Input(component_id="liste_marches_url", component_property="pathname"),
)
def liste_marches(url):
    org_type = url.split("/")[-2]
    org_id = url.split("/")[-1]

    def make_link_list() -> list:
        table = (
            "acheteurs_marches" if org_type == "acheteur"
            else "titulaires_marches" if org_type == "titulaire"
            else None
        )
        if table is None:
            raise ValueError
        rows = get_cursor().execute(
            f"SELECT uid, objet FROM {table} WHERE {org_type}_id = ?",
            [org_id],
        ).fetchall()

        return [
            html.Li(
                dcc.Link(
                    objet,
                    href=f"/marches/{uid}",
                    title=f"Marchés public attribué : {objet}",
                )
            )
            for uid, objet in rows
        ]

    nom, verbe = make_org_nom_verbe(org_type, org_id)
    return [
        html.H3(f"Marchés publics {verbe} par {nom}"),
        html.Ul(make_link_list()),
    ]
  • Step 2: Run tests and smoke-test

Run: uv run pytest tests/test_main.py -v

Manually: uv run run.py, navigate to /departements/75/acheteur/123, confirm list of marchés renders.

  • Step 3: Commit
git add src/pages/arbre/liste_marches_org.py
git commit -m "arbre/liste_marches_org.py : requêtes DuckDB pour les listes de marchés"

Task 15: Migrate src/pages/tableau.py

Files:

  • Modify: src/pages/tableau.py

  • Step 1: Replace imports

Remove df from src.utils imports; remove schema from src.utils imports if present, re-import it from src.db:

from src.utils import (
    columns,
    filter_table_data,
    get_default_hidden_columns,
    invert_columns,
    logger,
    meta_content,
    prepare_table_data,
    sort_table_data,
    update_date_iso,
)
from src.db import query_marches, schema
  • Step 2: Replace df.columns at line 64
columns=[{"id": col, "name": col} for col in schema.names()],
  • Step 3: Replace df.width at line 131

In the dcc.Markdown block, replace {str(df.width)} with {len(schema.names())}.

  • Step 4: Replace df.lazy() at line 319 in download_data
def download_data(n_clicks, filter_query, sort_by, hidden_columns: list = None):
    lff: pl.LazyFrame = query_marches().lazy()

    if hidden_columns:
        lff = lff.drop(hidden_columns)

    if filter_query:
        lff = filter_table_data(lff, filter_query, "tab download")

    if sort_by and len(sort_by) > 0:
        lff = sort_table_data(lff, sort_by)

    def to_bytes(buffer):
        lff.collect(engine="streaming").write_excel(buffer, worksheet="DECP")

    date = datetime.now().strftime("%Y-%m-%d_%H:%M:%S")
    return dcc.send_bytes(to_bytes, filename=f"decp_{date}.xlsx")

Note: query_marches() with no arguments returns the full table into a Polars frame, which is then lazily filtered. This is equivalent to today's df.lazy() — memory profile is identical to the current path during download. We accept this: the download endpoint is infrequent and downloads are bounded by filter_query in practice. A future optimization could push the filter into SQL, but it requires reimplementing filter_table_data in SQL. Out of scope.

  • Step 5: Run tests and smoke-test

Run: uv run pytest tests/test_main.py -v

Manually: uv run run.py, navigate to /tableau, apply a filter, try the download button.

  • Step 6: Commit
git add src/pages/tableau.py
git commit -m "tableau.py : migration vers query_marches et schema depuis src.db"

Task 16: Migrate src/pages/observatoire.py

Files:

  • Modify: src/pages/observatoire.py

  • Step 1: Replace imports

Remove df from src.utils imports; add from src.db import query_marches, schema. The columns import from src.utils stays (it already resolves to schema.names() via Task 9).

  • Step 2: Replace df.columns at line 75

Wherever df.columns is used, replace with schema.names().

  • Step 3: Replace df.lazy() at lines 668, 793, 884

Each call site currently reads df.lazy(). Replace with query_marches().lazy():

# Line ~668
lff: pl.LazyFrame = query_marches().lazy()
lff = prepare_dashboard_data(lff=lff, **filter_params)

# Line ~793
lff = prepare_dashboard_data(lff=query_marches().lazy(), **(filter_params or {}))

# Line ~884
lff = prepare_dashboard_data(lff=query_marches().lazy(), **(filter_params or {}))

Same caveat as tableau.py download: this materializes the full table before filtering in Polars. Since observatoire results are already cached via @cache.memoize(timeout=3600), the rebuild cost is amortized. Future optimization (pushing prepare_dashboard_data filters into SQL) is tracked separately and out of scope.

  • Step 4: Run tests and smoke-test

Run: uv run pytest tests/test_main.py -v

Manually: uv run run.py, navigate to /observatoire, apply and remove filters, confirm charts render.

  • Step 5: Commit
git add src/pages/observatoire.py
git commit -m "observatoire.py : migration vers query_marches et schema depuis src.db"

Task 17: Migrate src/figures.py

Files:

  • Modify: src/figures.py

  • Step 1: Remove df from the src.utils import block

from src.utils import (
    add_links,
    data_schema,
    departements_geojson,
    format_number,
    setup_table_columns,
)
from src.db import schema
  • Step 2: Replace df.columns at line 777
for col in schema.names()
  • Step 3: Run tests and smoke-test

Run: uv run pytest tests/test_main.py -v

Manually: uv run run.py, exercise pages that render charts (/acheteurs/123, /observatoire).

  • Step 4: Commit
git add src/figures.py
git commit -m "figures.py : suppression de l'import df global, usage de schema depuis src.db"

Task 18: Remove the in-memory globals from src/utils.py

Files:

  • Modify: src/utils.py

No page now references df, df_acheteurs_departement, df_titulaires_departement, df_acheteurs_marches, df_titulaires_marches (verify with grep below). df_acheteurs and df_titulaires remain (homepage search).

  • Step 1: Verify no consumers remain

Run:

rg '\bdf\b|df_acheteurs_departement|df_titulaires_departement|df_acheteurs_marches|df_titulaires_marches' src/

Expected: only matches inside src/utils.py (the globals themselves and get_decp_data/get_org_data) and no matches in src/pages/ or src/figures.py.

  • Step 2: Repoint df_acheteurs / df_titulaires to DuckDB

In src/utils.py, replace the bottom-of-file block (lines 891913 in the current file) with:

# df_acheteurs / df_titulaires sont conservés en mémoire pour alimenter
# la recherche sur la page d'accueil (autocomplétion, filtrage par sous-chaîne
# à chaque frappe). Les colonnes reproduisent la sortie historique de
# get_org_data(df, org_type).
def _build_org_frame(org_type: str) -> pl.DataFrame:
    org_cols = [
        c for c in schema.names()
        if c.startswith(f"{org_type}_")
        and c not in (f"{org_type}_latitude", f"{org_type}_longitude")
    ]
    select_list = ", ".join(org_cols)
    group_list = ", ".join(org_cols)
    sql = (
        f"SELECT {select_list}, COUNT(*) AS \"Marchés\" "
        f"FROM decp GROUP BY {group_list}"
    )
    return get_cursor().execute(sql).pl()


df_acheteurs = _build_org_frame("acheteur")
df_titulaires = _build_org_frame("titulaire")

columns = schema.names()
  • Step 3: Delete now-unused functions

Remove get_decp_data (lines 234272) — it has been inlined into src.db._load_source_frame. Remove get_org_data (lines 275284) — replaced by _build_org_frame.

  • Step 4: Remove the from polars.exceptions import ComputeError import if no longer used

Run:

rg 'ComputeError' src/utils.py

If empty, remove the import line near the top of utils.py.

Same check for from time import localtime, sleep:

rg '\bsleep\b|\blocaltime\b' src/utils.py

Remove anything unused.

  • Step 5: Run the full test suite

Run: uv run pytest -v

Expected: all tests pass.

  • Step 6: Manual smoke test of every page
uv run run.py

Visit in a browser:

  • / (recherche) — type in the acheteur and titulaire search fields, confirm autocomplete works.

  • /acheteurs/123 — confirm renders.

  • /titulaires/345 — confirm renders.

  • /marches/1 — confirm renders.

  • /tableau — filter, sort, download.

  • /observatoire — apply filters, confirm charts.

  • /departements and /departements/75 — confirm lists.

  • /departements/75/acheteur/123 — confirm marchés list.

  • Step 7: Commit

git add src/utils.py
git commit -m "Suppression des dataframes globaux remplacés par DuckDB"

Task 19: Measure memory impact

Files: none (documentation commit only)

  • Step 1: Measure RSS before (from a reference point on main)

Check out main in a separate worktree, start the app pointing at decp_prod.parquet, wait for import to finish, then:

ps -o rss= -p $(pgrep -f "gunicorn app:server" | head -1)

Record the value (expect ~several GB with 1.5M rows in memory).

  • Step 2: Measure RSS after (on this branch)

Back on the feature branch, same command. Record the value.

  • Step 3: Document the result

Append a ## Outcome section to the spec file docs/superpowers/specs/2026-04-15-duckdb-migration-design.md with the two RSS measurements and the ratio.

  • Step 4: Commit
git add docs/superpowers/specs/2026-04-15-duckdb-migration-design.md
git commit -m "Mesure de l'impact mémoire post-migration DuckDB"

Task 20: Final verification and branch readiness

  • Step 1: Run full test suite clean
uv run pytest -v

Expected: all green.

  • Step 2: Run pre-commit hooks on all files
pre-commit run --all-files

Expected: all hooks pass.

  • Step 3: Confirm no stray references remain
rg '\bdf\b' src/ | rg -v 'dff|df_acheteurs|df_titulaires|data_schema|self\.df'

Expected: no results (or only matches inside comments / docstrings).

  • Step 4: Verify the DB file is ignored by git
git status --ignored | rg duckdb

Expected: decp.duckdb (and .tmp, .lock if present) listed under ignored files.

Ready for PR / merge via the finishing-a-development-branch skill.


Self-Review Notes

  • Spec coverage: every section of docs/superpowers/specs/2026-04-15-duckdb-migration-design.md maps to a task:

    • Goal / Approach summary → Task 6 (query_marches), Tasks 1017 (migration), Task 18 (in-memory helpers).
    • Cache invalidation rule → Tasks 23 (should_rebuild + tests).
    • Concurrency → Tasks 57 (build_database + lock test).
    • Build logic → Task 5 (reuses _load_source_frame with Polars transforms).
    • Module layout — src/db.py public surface → Task 6.
    • Configuration (DATA_FILE_PARQUET_PATH reused; no new env var except REBUILD_DUCKDB) → Task 6 (_resolve_db_path).
    • Testing → Tasks 2, 4, 7, 8.
    • Migration order → Tasks 10 (marche), 11 (acheteur), 12 (titulaire), 1314 (arbre), 15 (tableau), 16 (observatoire), 17 (figures), 18 (remove globals).
    • Out-of-scope items (cache.py, observatoire-localstorage work, parquet schema, SQL views) → not touched.
  • Type consistency: schema is pl.Schema everywhere; conn is duckdb.DuckDBPyConnection; query_marches signature matches the spec. _load_source_frame is internal to src.db (single underscore, not exported).

  • Placeholder scan: no "TBD", no "similar to task N"; each step shows complete code or exact commands.