대용량 쿼리를 디스크에 싱크하기

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

Liam Brannigan

Data Scientist & Polars Contributor

대용량 출력 문제

레이지 쿼리가 디스크에 직접 저장되어야 하는 대용량 행 보존 추출을 생성하는 다이어그램

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

정제된 요청 데이터셋

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

clean_requests = (
    requests





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

정제된 요청 데이터셋

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

clean_requests = (
    requests
    .with_columns(
        pl.col("CREATED_DATE").dt.date(),
        pl.col("WARD").cast(pl.Int16),
        (pl.col("STATUS") == "Completed").alias("IS_COMPLETED"),
    )
)
Polars로 데이터 파이프라인 확장 및 최적화하기

수집 후 저장

clean_requests.collect(
    engine="streaming"
).write_parquet(
    "311_requests_clean.parquet"
)
Polars로 데이터 파이프라인 확장 및 최적화하기

Parquet으로 싱크하기

clean_requests.sink_parquet(
    "311_requests_clean.parquet",


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

싱크 제어하기

clean_requests.sink_parquet(
    "311_requests_clean.parquet",
    compression="zstd",
    row_group_size=100_000,
)
Polars로 데이터 파이프라인 확장 및 최적화하기

기타 싱크 메서드

df.sink_csv("output.csv")
df.sink_ndjson("output.ndjson")
Polars로 데이터 파이프라인 확장 및 최적화하기

파티션 출력

clean_requests.sink_parquet(
    pl.PartitionBy(



    ),

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

파티션 출력

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key="CREATED_DATE",
        max_rows_per_file=1_000_000,
    ),

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

파티션 출력

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key="CREATED_DATE",
        max_rows_per_file=1_000_000,
    ),
    mkdir=True,
)
Polars로 데이터 파이프라인 확장 및 최적화하기

파티션 출력

311_requests_clean/
  CREATED_DATE=2025-12-31/
    00000000.parquet
  CREATED_DATE=2026-01-01/
    00000000.parquet
Polars로 데이터 파이프라인 확장 및 최적화하기

다중 열 파티션

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key=["STATUS", "CREATED_DATE"],
    ),
    mkdir=True,
)
Polars로 데이터 파이프라인 확장 및 최적화하기

레이지 싱크 구성하기

internal_sink = clean_requests.sink_parquet(
    "311_requests_internal.parquet",
    lazy=True,
)





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

레이지 싱크 구성하기

internal_sink = clean_requests.sink_parquet(
    "311_requests_internal.parquet",
    lazy=True,
)

public_sink = clean_requests.drop("REQUEST_ID","REPORTER").sink_parquet(
    "311_requests_public.parquet",
    lazy=True,
)
Polars로 데이터 파이프라인 확장 및 최적화하기

싱크 멀티플렉싱

pl.collect_all(
    [internal_sink, public_sink],
)
Polars로 데이터 파이프라인 확장 및 최적화하기

저장된 출력 파일

311_requests_internal.parquet
311_requests_public.parquet
Polars로 데이터 파이프라인 확장 및 최적화하기

커스텀 배치 싱크

def send_batch(batch: pl.DataFrame) -> None:
    print(batch.height)
    batch_json = batch.write_json()
    # send batch_json to an API
Polars로 데이터 파이프라인 확장 및 최적화하기

커스텀 배치 싱크

clean_requests.sink_batches(
    send_batch,
    chunk_size=50_000,
)
50000
50000
...
29268
Polars로 데이터 파이프라인 확장 및 최적화하기

싱크 워크플로

  • sink_parquet()를 사용하여 대용량 쿼리 결과를 디스크에 직접 저장합니다
  • PartitionBy()를 사용하여 출력을 파티션 데이터셋으로 분할합니다
  • 레이지 싱크와 collect_all()을 함께 사용하여 여러 출력을 계획하고 실행합니다
  • 각 배치를 커스텀 함수로 처리할 때는 sink_batches()를 사용합니다

데이터를 정리된 파티션 서랍으로 분배하는 데이터 정렬 시스템의 아이소메트릭 뷰

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

실습해 봅시다!

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

Preparing Video For Download...