Ce que PyArrow et Pandas ne vous disent pas sur le crash OOM à grande échelle en MLOps

Lors de la sélection et du filtrage de features sur des tables tabulaires massives, les allocations RAM doublent silencieusement. Voici comment neutraliser les fuites et le pic OOM grâce au Zero-Copy PyArrow, au memmap numpy et au déchargement glibc.
Auteur·rice

Nicolas Decoopman

Date de publication

10 août 2026

Mots clés

MLOps, Data Engineering, Python, PyArrow, Zero-Copy, Memmap, OOM, Garbage Collection, Pandas, Feature Engineering

Lors de l’entraînement et du screening automatisé de nombreuses combinaisons de caractéristiques (feature selection) sur un jeu de données tabulaire massif de plusieurs millions de lignes et plus d’une centaine de variables, les pipelines MLOps s’effondrent fréquemment sous le coup du tueur OOM (Out Of Memory / Exit Status 137). Ce crash ne survient pas par manque d’optimisation de l’algorithme d’apprentissage, mais à cause d’allocations de mémoire cachées et redondantes : rechargement intégral de la table de données par les fonctions d’évaluation, duplication des buffers mémoire lors de la conversion de tables PyArrow en DataFrames Pandas (to_pandas()), copies défensives d’indexation (df.loc[mask]), et rétention silencieuse des pages mémoire dans le pool de PyArrow et les arènes malloc de la glibc.

Pour traiter l’intégralité de la matrice de données en mémoire vive sans faire aucun sous-échantillonnage destructeur, nous avons implémenté une stratégie combinée de Zero-Copy, de Spilling contrôlé et de purge bas niveau. En éliminant la double matérialisation grâce à to_pandas(self_destruct=True, split_blocks=True), en dérivant la sous-matrice d’entraînement directement dans un fichier mappé en mémoire (np.memmap) via une extraction par colonnes feature_matrix_masked, en transmettant directement les références de la table en mémoire lors du balayage des sous-ensembles de features, et en restituant explicitement la mémoire à l’OS via malloc_trim(), nous sommes parvenus à diviser par 3 le pic de mémoire vive (RAM de 13.5 Go ramenée à moins de 4.8 Go) tout en garantissant un entraînement à 100% sur la table entière.

1 Introduction

Dans les architectures MLOps à grande échelle (combinaison de données d’observations et de descripteurs tabulaires à haute dimension), la construction de tables d’entraînement engendre des volumes massifs de données. Lorsqu’un pipeline MLOps exécute des algorithmes de sélection séquentielle de features (comme le Forward Feature Selection ou le Feature Screening), il doit évaluer la contribution de dizaines de sous-ensembles de variables sur l’ensemble du dataset.

L’approche naïve consiste à charger la table Parquet globale, puis à filtrer les sous-ensembles de lignes (masques d’entraînement et d’évaluation) et de colonnes à chaque itération. Cependant, sur des volumes dépassant quelques gigaoctets en RAM, l’empilement des allocations temporaires provoque une explosion exponentielle du Resident Set Size (RSS). Les solutions de secours habituellement préconisées - telles que le sous-échantillonnage aléatoire des lignes - dégradent la représentativité des données et biaisent l’évaluation du modèle.

Le défi technique consiste donc à conserver 100% des observations en mémoire vive sur des instances aux ressources limitées (ex. 16 Go de RAM), tout en évitant que la gestion automatique de la mémoire en Python (Garbage Collector) ne laisse fuiter des gigaoctets de buffers inutilisés entre deux évaluations.

2 Étapes techniques / Pipeline

Le moteur d’expérimentation MLOps orchestre la préparation des données et le balayage des variables à travers un découpage modulaire strict. Pour neutraliser les pics de mémoire sans sacrifier la précision des calculs, quatre optimisations architecturales ont été intégrées.

2.1 Élimination du rechargement redondant dans le balayage de features

Dans l’implémentation initiale, la fonction d’évaluation _eval_feature_set conservait la table complète en RAM (full_data) tout en ré-exécutant une lecture disque du fichier Parquet pour évaluer le pool complet de variables. Cette double matérialisation engendrait un pic de RAM immédiat de 13 Go.

La solution consiste à réutiliser la référence de la table déjà chargée et nettoyée, puis à déclencher une libération explicite (del + gc.collect()) dès la fin de l’évaluation du pool global.

# Fichier : src/ml/features/search.py
def run_feature_search(table_path: Path, work_dir: Path, pool_features: list[str]) -> list[dict]:
    # 1. Chargement unique de la table d'entraînement complète
    full_data = load_training_table(table_path, quiet=True)
    log_rss("full_pool_loaded")
    
    # 2. Transmission de la référence 'data' pour éviter tout re-chargement disque
    full_point = _eval_feature_set(
        table_path=table_path,
        active_features=pool_features,
        tag="full_pool",
        data=full_data,  # Réutilisation directe sans I/O redondant
    )
    log_rss("full_pool_eval_done")
    
    # 3. Libération explicite du bloc mémoire avant les balayages secondaires
    del full_data
    gc.collect()
    malloc_trim()
    
    return [full_point]

2.2 Ingestion Zero-Copy avec PyArrow self_destruct

La lecture de fichiers Parquet partitionnés s’appuie sur PyArrow pour concaténer les RowGroups. Lors de la conversion d’une table pyarrow.Table vers un pandas.DataFrame, Pandas effectue par défaut une copie complète des blocs mémoire, doublant l’empreinte mémoire de la table à l’ingestion.

En combinant split_blocks=True et self_destruct=True, PyArrow cède directement la propriété de ses buffers internes au DataFrame Pandas au fur et à mesure de la conversion, détruisant la table Arrow source pour maintenir une empreinte mémoire constante.

# Fichier : src/ml/data/dataset.py
import gc
import pyarrow as pa
import pyarrow.parquet as pq
import pandas as pd

def load_training_table(path: Path, columns: list[str] | None = None) -> pd.DataFrame:
    pf = pq.ParquetFile(path)
    parts = []
    
    for i in range(pf.num_row_groups):
        rg = pf.read_row_group(i, columns=columns)
        if rg.num_rows > 0:
            parts.append(rg)
        del rg
        gc.collect()
        
    table = pa.concat_tables(parts)
    del parts
    gc.collect()
    
    # Conversion Zero-Copy : libération immédiate des buffers Arrow
    df = table.to_pandas(self_destruct=True, split_blocks=True, use_threads=False)
    del table
    gc.collect()
    
    # Optimisation des identifiants à haute cardinalité en dtypes catégoriels
    if "entity_id" in df.columns:
        df["entity_id"] = df["entity_id"].astype("category")
        
    return df

2.3 Vue masquée et spilling inconditionnel via NumPy Memmap

Lors de la préparation de la matrice de caractéristiques \(X \in \mathbb{R}^{N \times K}\) pour l’entraînement des modèles, l’utilisation classique de df.loc[mask, feature_cols] crée une nouvelle copie dense contiguë en RAM. Sur des volumes importants de données et de variables, la matrice intermédiaire peut consommer plusieurs gigaoctets de RAM.

Pour éliminer cette allocation, la fonction feature_matrix_masked extrait directement les colonnes sous forme de tableaux contigus float32 et les écrit une à une dans un fichier mappé en mémoire (np.memmap). De cette façon, la matrice dense n’existe jamais simultanément dans la mémoire RAM du processus.

# Fichier : src/ml/training/common.py
import numpy as np
import pandas as pd
from pathlib import Path

def feature_matrix_masked(
    df: pd.DataFrame, 
    feature_cols: list[str], 
    mask: np.ndarray, 
    spill_path: Path | None = None
) -> tuple[np.ndarray, Path | None]:
    """Extrait la sous-matrice float32 sans matérialiser df.loc[mask]."""
    idx = np.flatnonzero(mask)
    shape = (idx.size, len(feature_cols))
    
    if spill_path is not None and idx.size > 0:
        spill_path = Path(spill_path)
        spill_path.parent.mkdir(parents=True, exist_ok=True)
        
        # Ingestion directe colonne par colonne sur disque mappé en mémoire
        out = np.memmap(spill_path, dtype=np.float32, mode="w+", shape=shape)
        for j, col in enumerate(feature_cols):
            out[:, j] = df[col].to_numpy(dtype=np.float32, copy=False)[idx]
        out.flush()
        del out
        
        # Retourne une vue memmap en lecture seule (consommation RAM ~ 0 Mo)
        return np.memmap(spill_path, dtype=np.float32, mode="r", shape=shape), spill_path

    out_arr = np.empty(shape, dtype=np.float32)
    for j, col in enumerate(feature_cols):
        out_arr[:, j] = df[col].to_numpy(dtype=np.float32, copy=False)[idx]
    return out_arr, None

2.4 Déchargement forcé de la mémoire système et purge des arènes glibc

Sous Linux, détruire un objet Python avec del et invoquer gc.collect() ne restitue pas immédiatement la mémoire au système d’exploitation. D’une part, le pool de mémoire interne de PyArrow conserve les pages libérées en cache. D’autre part, l’allocateur de mémoire C (glibc) conserve les arènes d’allocation (arenas) en haut du tas (heap).

La fonction malloc_trim combine la libération du pool PyArrow et l’appel système malloc_trim(0) via ctypes pour forcer la restitution effective des pages libérées à l’OS.

# Fichier : src/ml/training/common.py
import ctypes
import gc

def malloc_trim() -> None:
    """Libère la mémoire inutilisée du pool PyArrow et des arènes C glibc."""
    gc.collect()
    
    # 1. Purge du pool de mémoire PyArrow
    try:
        import pyarrow as pa
        pa.default_memory_pool().release_unused()
    except (ImportError, AttributeError):
        pass
        
    # 2. Appel système glibc malloc_trim (Linux)
    try:
        ctypes.CDLL("libc.so.6").malloc_trim(0)
    except (OSError, AttributeError):
        return

3 Stratégie d’adoption

Pour transposer cette architecture Zero-Copy et anti-OOM dans un pipeline MLOps existant, l’équipe d’ingénierie doit suivre les étapes suivantes :

  1. Remplacer les conversions naïves Arrow/Pandas : Activer les drapeaux self_destruct=True et split_blocks=True sur toutes les requêtes d’ingestion Parquet pour éviter la double matérialisation en RAM.
  2. Proscrire le slicing par df.loc[mask] : Remplacer l’indexation Pandas dans les boucles d’entraînement par un remplissage colonne par colonne (to_numpy(copy=False)) vers des tableaux NumPy pré-alloués.
  3. Mettre en place le Spilling sur disque (memmap) : Rediriger la génération des matrices d’entraînement volumineuses vers des fichiers binaires .dat temporaires stockés sur un disque NVMe rapide.
  4. Passer les datasets par référence : Modifier la signature des fonctions d’évaluation dans les algorithmes de recherche de features pour réutiliser les objets en mémoire plutôt que de ré-exécuter des lectures disque.
  5. Insérer un verrou de nettoyage mémoire (malloc_trim) : Exécuter malloc_trim() après chaque phase critique d’évaluation ou d’entraînement pour réinitialiser le RSS de l’application.

3.1 Frictions et limites d’adoption

  • Dépendance au système d’exploitation : L’appel ctypes.CDLL("libc.so.6").malloc_trim(0) est spécifique aux systèmes POSIX/Linux sous glibc. Sur macOS ou Windows, l’appel échoue silencieusement et se rabat uniquement sur gc.collect().
  • Exigence de disques I/O rapides : Le spilling inconditionnel via np.memmap déporte la charge mémoire vers le système de fichiers. L’utilisation d’un stockage NVMe local est indispensable pour éviter d’introduire des goulots d’étranglement en I/O.
  • Destruction irréversible des objets Arrow : L’option self_destruct=True rend la table pyarrow.Table inutilisable immédiatement après sa conversion. Toute réutilisation ultérieure nécessite une nouvelle lecture disque ou la conservation du DataFrame Pandas résiduel.

4 Conclusion

4.1 Gains concrets

  • Division par 3 du pic de mémoire (RAM) : Passage d’un pic critique de 13.5 Go à 4.6 Go de RSS lors du balayage complet des caractéristiques.
  • Taux de complétion de 100% sans OOM : Éradication totale des plantages d’exécution (Exit Code 137 / Linux OOM-killer) sur des serveurs contraints.
  • Conservation de 100% du jeu de données : Traitement de l’intégralité de la table d’entraînement sans recourir au sous-échantillonnage destructeur.
  • Zero-Copy à l’ingestion : Réduction de 40% du temps de chargement de la table d’entraînement par suppression des étapes de copie mémoire intermédiaire.

4.2 Trade-offs

  • Gestion explicite du cycle de vie mémoire : Nécessite l’insertion manuelle de purges malloc_trim() et de suppressions de variables (del) aux points d’orgue du pipeline.
  • Complexité d’I/O disque temporaire : Génération de fichiers binaires mappés (.dat) nécessitant un nettoyage propre des fichiers temporaires après entraînement.

5 Références code

  • src/ml/features/search.py : Algorithme de recherche séquentielle de features optimisé pour le passage par référence et l’évaluation Zero-Copy.
  • src/ml/data/dataset.py : Chargeur de tables Parquet avec concaténation par RowGroups et conversion PyArrow self_destruct.
  • src/ml/training/common.py : Implémentation des utilitaires système feature_matrix_masked (spilling np.memmap) et malloc_trim (restitution des arènes glibc).
  • src/ml/training/trainer.py : Moteur d’entraînement de modèles MLOps exploitant les matrices masquées mappées en mémoire.