Parquet sink 다중화
이번에는 동일한 정제된 쿼리에서 두 가지 추출본이 필요합니다. 하나는 디지털 대출용이고, 다른 하나는 실물 대출용입니다. Polars가 공유 스캔을 한 번만 계획할 수 있도록 두 쓰기 작업을 지연 방식으로 구성한 후, 한 번의 패스로 함께 실행하세요.
clean_checkouts와 DIGITAL_EXPORT_PATH, PHYSICAL_EXPORT_PATH는 미리 로드되어 있습니다.
이 연습은 강의의 일부입니다
Polars로 데이터 파이프라인 확장 및 최적화하기
연습 안내
- 두 sink가 즉시 실행되지 않도록 지연 방식으로 구성하세요.
- 스트리밍 엔진에서 두 sink를 함께 실행하세요.
실습형 인터랙티브 연습
이 예제를 이 샘플 코드를 완성하여 풀어보세요.
# Lazy sink for digital
digital_sink = (
clean_checkouts
.filter(pl.col("use") == "Digital")
.sink_parquet(DIGITAL_EXPORT_PATH, lazy=____)
)
# Lazy sink for physical
physical_sink = (
clean_checkouts
.filter(pl.col("use") == "Physical")
.sink_parquet(PHYSICAL_EXPORT_PATH, lazy=____)
)
# Run both sinks together
pl.____([digital_sink, physical_sink], engine="streaming")
# Check the row counts of each extract
result = pl.DataFrame(
{
"extract": ["digital", "physical"],
"rows": [
pl.scan_parquet(DIGITAL_EXPORT_PATH).select(pl.len()).collect().item(),
pl.scan_parquet(PHYSICAL_EXPORT_PATH).select(pl.len()).collect().item(),
],
}
)
print(result)