Мультиплексирование Parquet-приёмников
Команде нужно получить два среза из одного и того же очищенного запроса: один для цифровых выдач, другой — для физических. Создайте обе операции записи в ленивом режиме, чтобы Polars мог спланировать общее сканирование один раз, а затем выполнить их за один проход.
clean_checkouts уже загружен, а также DIGITAL_EXPORT_PATH и PHYSICAL_EXPORT_PATH.
Это упражнение является частью курса
Масштабирование и оптимизация конвейеров данных с Polars
Инструкции к упражнению
- Создайте оба приёмника в ленивом режиме, чтобы они не выполнялись немедленно.
- Запустите оба приёмника совместно на потоковом движке.
Интерактивное практическое упражнение
Попробуйте выполнить это упражнение, дополнив этот пример кода.
# 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)