Werken met de streaming- en GPU-engines

Data-pipelines schalen en optimaliseren met Polars

Liam Brannigan

Data Scientist & Polars Contributor

Query-uitvoermotoren

Stroomschema dat een query toont van de API via de queryoptimizer naar de in-memory-engine.

Data-pipelines schalen en optimaliseren met Polars

Query-uitvoermotoren

Stroomschema dat een query toont van de API via de queryoptimizer naar de in-memory-engine.

Data-pipelines schalen en optimaliseren met Polars

Query-uitvoermotoren

Stroomschema dat een query toont van de API via de queryoptimizer naar de in-memory- en streaming-engines.

Data-pipelines schalen en optimaliseren met Polars

Query-uitvoermotoren

Stroomschema dat een query toont van de API via de queryoptimizer naar de in-memory-, streaming- en GPU-engines.

Data-pipelines schalen en optimaliseren met Polars

Een luie samenvatting van requests

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

query = (
    requests
    .filter(pl.col("STATUS") == "Completed")
    .group_by("DEPARTMENT")
    .len()
    .sort("len", descending=True)
)
Data-pipelines schalen en optimaliseren met Polars

De standaardengine gebruiken

query.collect()
shape: (5, 2)
| DEPARTMENT        | len     |
| ---               | ---     |
| str               | u32     |
|-------------------|---------|
| 311 City Services | 4859161 |
| Sanitation        | 3406631 |
| Aviation          | 2337842 |
| Transportation    | 1468240 |
| Water             | 416535  |
Data-pipelines schalen en optimaliseren met Polars

De streamingengine gebruiken

query.collect(
    engine="streaming"
)
shape: (5, 2)
| DEPARTMENT        | len     |
| ---               | ---     |
| str               | u32     |
|-------------------|---------|
| 311 City Services | 4859161 |
| Sanitation        | 3406631 |
| Aviation          | 2337842 |
| Transportation    | 1468240 |
| Water             | 416535  |
Data-pipelines schalen en optimaliseren met Polars

De streamingengine gebruiken

Afbeelding van een grote tabel die wordt gestreamd.

  • Niet-ondersteunde streamingstappen vallen automatisch terug op de standaard-engine
Data-pipelines schalen en optimaliseren met Polars

De streamingengine gebruiken

import polars as pl
pl.Config.set_engine_affinity("streaming")
Data-pipelines schalen en optimaliseren met Polars

De GPU-engine gebruiken

query.collect(
    engine="gpu"
)

$$

  • Vereist NVIDIA GPU-hardware
Data-pipelines schalen en optimaliseren met Polars

Een engine kiezen

$$

  • Eerst de standaard-engine; wissel alleen bij geheugenzorgen

$$

  • Streaming voor grote queries buiten het geheugen

$$

  • GPU voor snelheid, maar alleen met NVIDIA-hardware

Een engine kiezen

Data-pipelines schalen en optimaliseren met Polars

Werken in batches

completed_requests = (
    requests
    .filter(pl.col("STATUS") == "Completed")
    .select("SR_NUMBER", "TYPE", "DEPARTMENT", "CREATED_DATE")
)
Data-pipelines schalen en optimaliseren met Polars

Elke batch verwerken

for batch in completed_requests.collect_batches(

):


Data-pipelines schalen en optimaliseren met Polars

Elke batch verwerken

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


Data-pipelines schalen en optimaliseren met Polars

Elke batch verwerken

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)
Data-pipelines schalen en optimaliseren met Polars

Laten we oefenen!

Data-pipelines schalen en optimaliseren met Polars

Preparing Video For Download...