Ghép kênh các Parquet sink
Nhóm hiện cần hai bản trích xuất từ cùng một truy vấn đã làm sạch: một cho lượt mượn số (digital) và một cho lượt mượn vật lý. Hãy xây dựng cả hai thao tác ghi theo kiểu lười (lazy) để Polars có thể lập kế hoạch quét chung một lần, rồi chạy cả hai trong một lượt.
clean_checkouts đã được nạp sẵn, cùng với DIGITAL_EXPORT_PATH và PHYSICAL_EXPORT_PATH.
Bài tập này là một phần của khóa học
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars
Hướng dẫn bài tập
- Xây dựng cả hai sink theo kiểu lười để chúng không thực thi ngay.
- Chạy cả hai sink cùng lúc trên engine streaming.
Bài tập tương tác thực hành trực tiếp
Hãy thử làm bài tập này bằng cách hoàn thành đoạn mã mẫu này.
# 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)