Arbeiten mit Streaming- und GPU-Engines

Skalieren und Optimieren von Data-Pipelines mit Polars

Liam Brannigan

Data Scientist & Polars Contributor

Abfrage-Engines

Ablaufdiagramm: Eine Abfrage geht von der API über den Query-Optimizer zur In-Memory-Engine.

Skalieren und Optimieren von Data-Pipelines mit Polars

Abfrage-Engines

Ablaufdiagramm: Eine Abfrage geht von der API über den Query-Optimizer zur In-Memory-Engine.

Skalieren und Optimieren von Data-Pipelines mit Polars

Abfrage-Engines

Ablaufdiagramm: Eine Abfrage geht von der API über den Query-Optimizer zu In-Memory- und Streaming-Engines.

Skalieren und Optimieren von Data-Pipelines mit Polars

Abfrage-Engines

Ablaufdiagramm: Eine Abfrage geht von der API über den Query-Optimizer zu In-Memory-, Streaming- und GPU-Engines.

Skalieren und Optimieren von Data-Pipelines mit Polars

Eine faule Request-Zusammenfassung

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

query = (
    requests
    .filter(pl.col("STATUS") == "Completed")
    .group_by("DEPARTMENT")
    .len()
    .sort("len", descending=True)
)
Skalieren und Optimieren von Data-Pipelines mit Polars

Standard-Engine verwenden

query.collect()
shape: (5, 2)
| DEPARTMENT        | len     |
| ---               | ---     |
| str               | u32     |
|-------------------|---------|
| 311 City Services | 4859161 |
| Sanitation        | 3406631 |
| Aviation          | 2337842 |
| Transportation    | 1468240 |
| Water             | 416535  |
Skalieren und Optimieren von Data-Pipelines mit Polars

Streaming-Engine gezielt nutzen

query.collect(
    engine="streaming"
)
shape: (5, 2)
| DEPARTMENT        | len     |
| ---               | ---     |
| str               | u32     |
|-------------------|---------|
| 311 City Services | 4859161 |
| Sanitation        | 3406631 |
| Aviation          | 2337842 |
| Transportation    | 1468240 |
| Water             | 416535  |
Skalieren und Optimieren von Data-Pipelines mit Polars

Streaming-Engine gezielt nutzen

Bild einer großen, gestreamten Tabelle.

  • Nicht unterstützte Streaming-Schritte fallen automatisch auf die Standard-Engine zurück
Skalieren und Optimieren von Data-Pipelines mit Polars

Streaming-Engine gezielt nutzen

import polars as pl
pl.Config.set_engine_affinity("streaming")
Skalieren und Optimieren von Data-Pipelines mit Polars

GPU-Engine gezielt nutzen

query.collect(
    engine="gpu"
)

$$

  • Benötigt NVIDIA-GPU-Hardware
Skalieren und Optimieren von Data-Pipelines mit Polars

Engine auswählen

$$

  • Zuerst Standard-Engine; nur bei Speicherbedarf wechseln

$$

  • Streaming für große Abfragen, die nicht in den Speicher passen

$$

  • GPU für Tempo, aber nur mit NVIDIA-Hardware

Engine auswählen

Skalieren und Optimieren von Data-Pipelines mit Polars

In Batches arbeiten

completed_requests = (
    requests
    .filter(pl.col("STATUS") == "Completed")
    .select("SR_NUMBER", "TYPE", "DEPARTMENT", "CREATED_DATE")
)
Skalieren und Optimieren von Data-Pipelines mit Polars

Jeden Batch verarbeiten

for batch in completed_requests.collect_batches(

):


Skalieren und Optimieren von Data-Pipelines mit Polars

Jeden Batch verarbeiten

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


Skalieren und Optimieren von Data-Pipelines mit Polars

Jeden Batch verarbeiten

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)
Skalieren und Optimieren von Data-Pipelines mit Polars

Lass uns üben!

Skalieren und Optimieren von Data-Pipelines mit Polars

Preparing Video For Download...