Trabalhando com os mecanismos de streaming e GPU

Dimensionamento e Otimização de Pipelines de Dados com Polars

Liam Brannigan

Data Scientist & Polars Contributor

Mecanismos de execução de consultas

Diagrama de fluxo comparando uma consulta indo da API, passando pelo otimizador, até o mecanismo em memória.

Dimensionamento e Otimização de Pipelines de Dados com Polars

Mecanismos de execução de consultas

Diagrama de fluxo comparando uma consulta indo da API, passando pelo otimizador, até o mecanismo em memória padrão.

Dimensionamento e Otimização de Pipelines de Dados com Polars

Mecanismos de execução de consultas

Diagrama de fluxo comparando uma consulta indo da API, passando pelo otimizador, até os mecanismos em memória e streaming.

Dimensionamento e Otimização de Pipelines de Dados com Polars

Mecanismos de execução de consultas

Diagrama de fluxo comparando uma consulta indo da API, passando pelo otimizador, até os mecanismos em memória, streaming e GPU.

Dimensionamento e Otimização de Pipelines de Dados com Polars

Um resumo lazy das solicitações

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

query = (
    requests
    .filter(pl.col("STATUS") == "Completed")
    .group_by("DEPARTMENT")
    .len()
    .sort("len", descending=True)
)
Dimensionamento e Otimização de Pipelines de Dados com Polars

Usando o mecanismo padrão

query.collect()
shape: (5, 2)
| DEPARTMENT        | len     |
| ---               | ---     |
| str               | u32     |
|-------------------|---------|
| 311 City Services | 4859161 |
| Sanitation        | 3406631 |
| Aviation          | 2337842 |
| Transportation    | 1468240 |
| Water             | 416535  |
Dimensionamento e Otimização de Pipelines de Dados com Polars

Usando o mecanismo de streaming

query.collect(
    engine="streaming"
)
shape: (5, 2)
| DEPARTMENT        | len     |
| ---               | ---     |
| str               | u32     |
|-------------------|---------|
| 311 City Services | 4859161 |
| Sanitation        | 3406631 |
| Aviation          | 2337842 |
| Transportation    | 1468240 |
| Water             | 416535  |
Dimensionamento e Otimização de Pipelines de Dados com Polars

Usando o mecanismo de streaming

Imagem de uma grande tabela sendo transmitida.

  • Etapas sem suporte em streaming caem automaticamente no mecanismo padrão
Dimensionamento e Otimização de Pipelines de Dados com Polars

Usando o mecanismo de streaming

import polars as pl
pl.Config.set_engine_affinity("streaming")
Dimensionamento e Otimização de Pipelines de Dados com Polars

Usando o mecanismo de GPU

query.collect(
    engine="gpu"
)

$$

  • Requer hardware NVIDIA GPU
Dimensionamento e Otimização de Pipelines de Dados com Polars

Escolhendo um mecanismo

$$

  • Primeiro use o mecanismo padrão; troque só se memória for problema

$$

  • Streaming para consultas grandes que não cabem na memória

$$

  • GPU para mais velocidade, mas só com hardware NVIDIA disponível

Escolhendo um mecanismo

Dimensionamento e Otimização de Pipelines de Dados com Polars

Trabalhando em lotes

completed_requests = (
    requests
    .filter(pl.col("STATUS") == "Completed")
    .select("SR_NUMBER", "TYPE", "DEPARTMENT", "CREATED_DATE")
)
Dimensionamento e Otimização de Pipelines de Dados com Polars

Processando cada lote

for batch in completed_requests.collect_batches(

):


Dimensionamento e Otimização de Pipelines de Dados com Polars

Processando cada lote

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


Dimensionamento e Otimização de Pipelines de Dados com Polars

Processando cada lote

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)
Dimensionamento e Otimização de Pipelines de Dados com Polars

Vamos praticar!

Dimensionamento e Otimização de Pipelines de Dados com Polars

Preparing Video For Download...