НачатьНачать бесплатно

Мультиплексирование 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)
Редактировать и запускать код