Làm việc với engine streaming và GPU

Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Liam Brannigan

Data Scientist & Polars Contributor

Engine thực thi truy vấn

Sơ đồ luồng so sánh truy vấn đi từ API, qua bộ tối ưu truy vấn tới engine trong bộ nhớ.

Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Engine thực thi truy vấn

Sơ đồ luồng so sánh truy vấn đi từ API, qua bộ tối ưu truy vấn tới engine trong bộ nhớ.

Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Engine thực thi truy vấn

Sơ đồ luồng so sánh truy vấn đi từ API, qua bộ tối ưu truy vấn tới engine trong bộ nhớ và streaming.

Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Engine thực thi truy vấn

Sơ đồ luồng so sánh truy vấn đi từ API, qua bộ tối ưu truy vấn tới engine trong bộ nhớ, streaming và GPU.

Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Tóm tắt lazy cho yêu cầu

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

query = (
    requests
    .filter(pl.col("STATUS") == "Completed")
    .group_by("DEPARTMENT")
    .len()
    .sort("len", descending=True)
)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Dùng engine mặc định

query.collect()
shape: (5, 2)
| DEPARTMENT        | len     |
| ---               | ---     |
| str               | u32     |
|-------------------|---------|
| 311 City Services | 4859161 |
| Sanitation        | 3406631 |
| Aviation          | 2337842 |
| Transportation    | 1468240 |
| Water             | 416535  |
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Nhắm tới engine streaming

query.collect(
    engine="streaming"
)
shape: (5, 2)
| DEPARTMENT        | len     |
| ---               | ---     |
| str               | u32     |
|-------------------|---------|
| 311 City Services | 4859161 |
| Sanitation        | 3406631 |
| Aviation          | 2337842 |
| Transportation    | 1468240 |
| Water             | 416535  |
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Nhắm tới engine streaming

Hình ảnh bảng lớn được stream.

  • Bước không hỗ trợ streaming sẽ tự động quay về engine mặc định
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Nhắm tới engine streaming

import polars as pl
pl.Config.set_engine_affinity("streaming")
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Nhắm tới engine GPU

query.collect(
    engine="gpu"
)

$$

  • Cần phần cứng GPU NVIDIA
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Chọn engine

$$

  • Dùng engine mặc định trước; chỉ đổi khi lo ngại bộ nhớ

$$

  • Streaming cho truy vấn lớn, vượt bộ nhớ

$$

  • GPU để tăng tốc, chỉ khi có phần cứng NVIDIA

Chọn engine

Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Làm việc theo lô

completed_requests = (
    requests
    .filter(pl.col("STATUS") == "Completed")
    .select("SR_NUMBER", "TYPE", "DEPARTMENT", "CREATED_DATE")
)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Xử lý từng lô

for batch in completed_requests.collect_batches(

):


Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Xử lý từng lô

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


Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Xử lý từng lô

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)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Passons à la pratique !

Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Preparing Video For Download...