Brains Up AnalyticsBRAINSUPAnalytics
PolarsPythonStreamingData EngineeringDuckDB

Polars streaming : traiter des données plus grandes que la RAM — sans Spark

Comment la streaming engine de Polars traite des jeux de données qui ne tiennent pas en mémoire avec scan_parquet + sink_parquet, en gardant une utilisation de RAM constante et sans cluster.

Il y a un réflexe courant dans les équipes data : dès qu'un fichier dépasse quelques gigaoctets et que pandas commence à saturer la mémoire avec un MemoryError, la réponse par défaut est « passe sur Spark ». Or une bonne partie de ces cas n'est pas un vrai problème de big data — c'est un problème de donnée qui ne tient pas en RAM d'un seul coup. Et c'est précisément le scénario que Polars résout dans un seul processus Python, sans cluster, sans JVM, sans infrastructure supplémentaire.

La différence tient à la façon dont la donnée est traitée. pandas est eager : il charge tout en mémoire et exécute opération par opération. Polars possède une streaming engine qui travaille la donnée par petits lots, en gardant une empreinte mémoire constante — que le fichier fasse 1 Go ou 500 Go.

Les trois concepts qui font que ça marche

1. Lazy API : scan_parquet ne lit rien

Quand vous appelez pl.read_parquet(), Polars charge le fichier entier en mémoire — comportement eager, comme pandas. pl.scan_parquet(), en revanche, fait quelque chose de différent : il ne lit pas les données. Il crée un LazyFrame, qui n'est qu'un plan logique décrivant les opérations que vous comptez exécuter.

Rien ne se passe tant que vous n'appelez pas explicitement .collect() (matérialise en mémoire) ou .sink_parquet() (écrit sur disque). C'est ce qui ouvre la porte au traitement larger-than-RAM.

import polars as pl

# rien n'est lu ici — seul un plan lazy est créé
lf = pl.scan_parquet("ventes/*.parquet")

2. Query planner : optimisation avant l'exécution

Comme Polars connaît le plan entier avant de l'exécuter, il optimise la requête — d'une manière qu'une approche eager ne peut pas :

  • Projection pushdown : ne lit sur disque que les colonnes que vous utilisez vraiment.
  • Predicate pushdown : pousse les filtres au moment de la lecture, en écartant les lignes avant de les charger.
  • Réordonnancement et fusion des opérations pour minimiser les passages sur les données.

Vous écrivez le code de façon déclarative et chaînée ; le planificateur s'occupe du reste.

res = (
    lf.filter(pl.col("etat") == "SP")
      .group_by("magasin")
      .agg(pl.col("montant").sum())
)

3. Streaming sinks : écrire des résultats plus grands que la RAM

L'étape finale est là où la magie opère. Au lieu de tout matérialiser avec .collect(), vous utilisez un sink :

res.sink_parquet("resume.parquet")

sink_parquet() évalue la requête en mode streaming et écrit le résultat directement sur disque, lot par lot. Cela garantit une utilisation de mémoire constante quelle que soit la taille du jeu de données — et permet que le résultat final soit plus grand que la RAM disponible.

Sous le capot, des opérateurs comme group_by, sort et les côtés build et probe d'un equi-join sont spillable : lorsqu'ils détectent une pression mémoire, ils écrivent l'état accumulé dans des fichiers temporaires, libérant de la RAM. Les file sinks (sink_parquet, sink_csv, sink_ipc) sont spillable par nature — ils écrivent déjà sur disque par définition.

Le pipeline complet

En rassemblant le tout, un pipeline qui traite 100 Go de données sur une machine avec 8 Go de RAM ressemble à ceci :

import polars as pl

# 1. plan lazy — rien n'est chargé
lf = pl.scan_parquet("ventes/*.parquet")

# 2. transformations optimisées par le query planner
res = (
    lf.filter(pl.col("etat") == "SP")
      .group_by("magasin")
      .agg(pl.col("montant").sum())
)

# 3. exécution en streaming, écriture sur disque lot par lot
res.sink_parquet("sortie.parquet")

Sans cluster. Sans MemoryError. Sans réécrire quoi que ce soit en PySpark.

Pourquoi c'est important (au-delà de la vitesse)

Polars est construit en Rust au-dessus d'Apache Arrow et livre généralement des gains de 5 à 50× sur pandas dans des charges réelles. Mais l'argument ici n'est pas seulement la performance — c'est l'architecture :

  • Moins d'infrastructure : un script Python remplace un cluster qu'il faut provisionner, surveiller et payer.
  • Moins de coût : pas de nœuds inactifs, pas d'overhead d'orchestration distribuée pour un volume qui tourne tranquillement sur une seule machine.
  • Plus de portabilité : le même code tourne sur votre ordinateur portable, dans un petit conteneur ou sur une VM — sans dépendances d'écosystème.

Et l'écosystème s'emboîte bien : DuckDB peut interroger un DataFrame Polars directement en mémoire via l'interface Arrow, avec zéro overhead de sérialisation. Beaucoup d'équipes utilisent les deux dans le même pipeline — Polars pour la manipulation de DataFrame et le feature engineering, DuckDB pour l'analyse en SQL sur les mêmes données.

Quand Spark vaut encore le coup

Être juste avec le bon outil compte. Spark garde du sens quand :

  • vous avez besoin d'une vraie mise à l'échelle horizontale, avec des données réparties sur plusieurs nœuds ;
  • le volume dépasse ce qu'une seule machine (même avec streaming et spill sur disque) peut traiter dans un délai raisonnable ;
  • vous avez déjà l'écosystème en place (Databricks, gouvernance dans Unity Catalog, jobs en production) et la cohérence opérationnelle pèse plus que la simplicité.

Pour le traitement single-node larger-than-RAM qui apparaît dans la plupart des équipes — agrégations, joins, nettoyage et transformation de gros fichiers —, la streaming engine de Polars couvre le cas largement. Et la roadmap renforce cette direction : Polars 2.0 apporte une streaming engine repensée (parallélisme morsel-driven + machines à états asynchrones en Rust), et Polars Cloud étend le moteur open-source à l'exécution serverless, la mise à l'échelle horizontale sur des données partitionnées et la tolérance aux pannes.

Résumé en trois étapes

  1. scan_parquet crée un plan lazy — rien n'est lu tant que vous ne le demandez pas.
  2. Chaînez filter / group_by / agg ; le query planner optimise la requête entière.
  3. sink_parquet s'exécute en streaming et écrit sur disque lot par lot — des résultats plus grands que la RAM, avec une mémoire constante et sans cluster.

Avant de provisionner un cluster pour votre prochain pipeline, mesurez si un scan_parquet + sink_parquet ne suffit pas. La plupart du temps, il suffit.

Articles liés

Vous avez aimé ? Découvrez les e-books pour du contenu approfondi.

E-books