Ce que PyArrow et Pandas ne vous disent pas sur le crash OOM à grande échelle en MLOps
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 df2.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, None2.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):
return3 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 :
- Remplacer les conversions naïves Arrow/Pandas : Activer les drapeaux
self_destruct=Trueetsplit_blocks=Truesur toutes les requêtes d’ingestion Parquet pour éviter la double matérialisation en RAM. - 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. - Mettre en place le Spilling sur disque (
memmap) : Rediriger la génération des matrices d’entraînement volumineuses vers des fichiers binaires.dattemporaires stockés sur un disque NVMe rapide. - 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.
- Insérer un verrou de nettoyage mémoire (
malloc_trim) : Exécutermalloc_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 sousglibc. Sur macOS ou Windows, l’appel échoue silencieusement et se rabat uniquement surgc.collect(). - Exigence de disques I/O rapides : Le spilling inconditionnel via
np.memmapdé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=Truerend la tablepyarrow.Tableinutilisable 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 PyArrowself_destruct.src/ml/training/common.py: Implémentation des utilitaires systèmefeature_matrix_masked(spillingnp.memmap) etmalloc_trim(restitution des arènesglibc).src/ml/training/trainer.py: Moteur d’entraînement de modèles MLOps exploitant les matrices masquées mappées en mémoire.