Travailler avec les moteurs streaming et GPU

Mise à l'échelle et optimisation des pipelines de données avec Polars

Liam Brannigan

Data Scientist & Polars Contributor

Moteurs d'exécution des requêtes

Schéma de flux comparant une requête passant par l'API, l'optimiseur de requêtes, puis le moteur en mémoire.

Mise à l'échelle et optimisation des pipelines de données avec Polars

Moteurs d'exécution des requêtes

Schéma de flux comparant une requête passant par l'API, l'optimiseur de requêtes, puis le moteur en mémoire.

Mise à l'échelle et optimisation des pipelines de données avec Polars

Moteurs d'exécution des requêtes

Schéma de flux comparant une requête passant par l'API, l'optimiseur de requêtes, puis les moteurs en mémoire et en streaming.

Mise à l'échelle et optimisation des pipelines de données avec Polars

Moteurs d'exécution des requêtes

Schéma de flux comparant une requête passant par l'API, l'optimiseur de requêtes, puis les moteurs en mémoire, en streaming et GPU.

Mise à l'échelle et optimisation des pipelines de données avec Polars

Résumé paresseux d'une requête

requests = pl.scan_parquet("311_Service_Requests.parquet")

query = (
    requests
    .filter(pl.col("STATUS") == "Completed")
    .group_by("DEPARTMENT")
    .len()
    .sort("len", descending=True)
)
Mise à l'échelle et optimisation des pipelines de données avec Polars

Utiliser le moteur par défaut

query.collect()
shape: (5, 2)
| DEPARTMENT        | len     |
| ---               | ---     |
| str               | u32     |
|-------------------|---------|
| 311 City Services | 4859161 |
| Sanitation        | 3406631 |
| Aviation          | 2337842 |
| Transportation    | 1468240 |
| Water             | 416535  |
Mise à l'échelle et optimisation des pipelines de données avec Polars

Cibler le moteur streaming

query.collect(
    engine="streaming"
)
shape: (5, 2)
| DEPARTMENT        | len     |
| ---               | ---     |
| str               | u32     |
|-------------------|---------|
| 311 City Services | 4859161 |
| Sanitation        | 3406631 |
| Aviation          | 2337842 |
| Transportation    | 1468240 |
| Water             | 416535  |
Mise à l'échelle et optimisation des pipelines de données avec Polars

Cibler le moteur streaming

Image d'une grande table diffusée en flux.

  • Les étapes non prises en charge en streaming basculent automatiquement sur le moteur par défaut
Mise à l'échelle et optimisation des pipelines de données avec Polars

Cibler le moteur streaming

import polars as pl
pl.Config.set_engine_affinity("streaming")
Mise à l'échelle et optimisation des pipelines de données avec Polars

Cibler le moteur GPU

query.collect(
    engine="gpu"
)

$$

  • Nécessite un GPU NVIDIA
Mise à l'échelle et optimisation des pipelines de données avec Polars

Choisir un moteur

$$

  • D'abord le moteur par défaut ; changer seulement si la mémoire pose problème

$$

  • Streaming pour les requêtes volumineuses hors mémoire

$$

  • GPU pour la vitesse, seulement si un matériel NVIDIA est dispo

Choisir un moteur

Mise à l'échelle et optimisation des pipelines de données avec Polars

Travailler par lots

completed_requests = (
    requests
    .filter(pl.col("STATUS") == "Completed")
    .select("SR_NUMBER", "TYPE", "DEPARTMENT", "CREATED_DATE")
)
Mise à l'échelle et optimisation des pipelines de données avec Polars

Traiter chaque lot

for batch in completed_requests.collect_batches(

):


Mise à l'échelle et optimisation des pipelines de données avec Polars

Traiter chaque lot

for batch in completed_requests.collect_batches(
    chunk_size=50_000,
):


Mise à l'échelle et optimisation des pipelines de données avec Polars

Traiter chaque lot

for batch in completed_requests.collect_batches(
    chunk_size=50_000,
):
    print(batch.shape)
    batch_json = batch.write_json()
(50000, 4)
(50000, 4)
...
(50000, 4)
(12409, 4)
Mise à l'échelle et optimisation des pipelines de données avec Polars

Passons à la pratique !

Mise à l'échelle et optimisation des pipelines de données avec Polars

Preparing Video For Download...