스트리밍 및 GPU 엔진 사용

Polars로 데이터 파이프라인 확장 및 최적화하기

Liam Brannigan

Data Scientist & Polars Contributor

쿼리 실행 엔진

API에서 쿼리 최적화기를 거쳐 인메모리 엔진으로 전달되는 쿼리 흐름 다이어그램.

Polars로 데이터 파이프라인 확장 및 최적화하기

쿼리 실행 엔진

API에서 쿼리 최적화기를 거쳐 인메모리 엔진으로 전달되는 쿼리 흐름 다이어그램.

Polars로 데이터 파이프라인 확장 및 최적화하기

쿼리 실행 엔진

API에서 쿼리 최적화기를 거쳐 인메모리 및 스트리밍 엔진으로 전달되는 쿼리 흐름 다이어그램.

Polars로 데이터 파이프라인 확장 및 최적화하기

쿼리 실행 엔진

API에서 쿼리 최적화기를 거쳐 인메모리, 스트리밍, GPU 엔진으로 전달되는 쿼리 흐름 다이어그램.

Polars로 데이터 파이프라인 확장 및 최적화하기

지연 요청 요약

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

query = (
    requests
    .filter(pl.col("STATUS") == "Completed")
    .group_by("DEPARTMENT")
    .len()
    .sort("len", descending=True)
)
Polars로 데이터 파이프라인 확장 및 최적화하기

기본 엔진 사용

query.collect()
shape: (5, 2)
| DEPARTMENT        | len     |
| ---               | ---     |
| str               | u32     |
|-------------------|---------|
| 311 City Services | 4859161 |
| Sanitation        | 3406631 |
| Aviation          | 2337842 |
| Transportation    | 1468240 |
| Water             | 416535  |
Polars로 데이터 파이프라인 확장 및 최적화하기

스트리밍 엔진 지정

query.collect(
    engine="streaming"
)
shape: (5, 2)
| DEPARTMENT        | len     |
| ---               | ---     |
| str               | u32     |
|-------------------|---------|
| 311 City Services | 4859161 |
| Sanitation        | 3406631 |
| Aviation          | 2337842 |
| Transportation    | 1468240 |
| Water             | 416535  |
Polars로 데이터 파이프라인 확장 및 최적화하기

스트리밍 엔진 지정

대용량 테이블이 스트리밍되는 이미지.

  • 스트리밍을 지원하지 않는 단계는 자동으로 기본 엔진으로 전환됩니다
Polars로 데이터 파이프라인 확장 및 최적화하기

스트리밍 엔진 지정

import polars as pl
pl.Config.set_engine_affinity("streaming")
Polars로 데이터 파이프라인 확장 및 최적화하기

GPU 엔진 지정

query.collect(
    engine="gpu"
)

$$

  • NVIDIA GPU 하드웨어 필요
Polars로 데이터 파이프라인 확장 및 최적화하기

엔진 선택

$$

  • 우선 기본 엔진 사용; 메모리 부족 시에만 전환

$$

  • 스트리밍: 메모리 초과 대용량 쿼리에 사용

$$

  • GPU: 속도 우선, NVIDIA 하드웨어 필요

엔진 선택

Polars로 데이터 파이프라인 확장 및 최적화하기

배치 처리

completed_requests = (
    requests
    .filter(pl.col("STATUS") == "Completed")
    .select("SR_NUMBER", "TYPE", "DEPARTMENT", "CREATED_DATE")
)
Polars로 데이터 파이프라인 확장 및 최적화하기

각 배치 처리

for batch in completed_requests.collect_batches(

):


Polars로 데이터 파이프라인 확장 및 최적화하기

각 배치 처리

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


Polars로 데이터 파이프라인 확장 및 최적화하기

각 배치 처리

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)
Polars로 데이터 파이프라인 확장 및 최적화하기

연습해 봅시다!

Polars로 데이터 파이프라인 확장 및 최적화하기

Preparing Video For Download...