Выполнение нескольких структурированных запросов потоковой передачи в одном кластере

Многие клиенты выполняют несколько запросов структурированной потоковой передачи в одном кластере Azure Databricks. Хотя этот шаблон поддерживается, Databricks рекомендует ограничить количество запросов на кластер, чтобы избежать проблем масштабирования и узких мест производительности. При бессерверных вычислениях Azure Databricks автоматически управляет масштабированием, поэтому эти аспекты берет на себя. Если вы используете классические вычисления, где вы управляете размерами драйвера и исполнителя, на этой странице описываются узкие места, которые следует учитывать и способы их устранения.

Примечание.

Databricks рекомендует использовать конвейеры Lakeflow для новых рабочих нагрузок потоковой передачи, которые автоматически управляют сложностью инфраструктуры. См. декларативные конвейеры Spark.

Использование нескольких запросов в одном кластере

Выполнение нескольких потоковых запросов в одном кластере снижает затраты на инфраструктуру, особенно если у вас много небольших потоков, которые не требуют выделенных вычислений. Основной компромисс заключается в общей точке отказа: если кластер выходит из строя, все потоки в нем тоже выходят из строя. Для критически важных конвейеров общий режим сбоя часто неприемлем.

Для рабочих нагрузок, смешивающих критически важные и некритичные потоки, Databricks рекомендует следующее:

  • Назначьте каждый поток приоритету в зависимости от его влияния на бизнес.
  • Размещайте потоки, критичные для бизнеса, в выделенных кластерах, даже несмотря на более высокую стоимость.
  • Совместно размещайте потоки с более низким приоритетом, чтобы совместно использовать вычислительные ресурсы и снизить затраты.

Подбор размера драйвера

Драйвер является общим ресурсом. Несколько запросов используют один и тот же ЦП, память, планировщик DAG, планировщик задач и выполнение UDF на стороне драйвера (например, foreachBatch). При запуске множества одновременных потоков обратите внимание на следующие узкие места помимо стандартного выделения ресурсов ЦП и памяти:

  • Накладные расходы Auto Loader: если ваши потоки используют Auto Loader, обнаружение файлов и перечисление каталогов увеличивают нагрузку на драйвер.
  • Ограничения ресурсов на уровне ОС (открытые файлы): выполнение большого объема потоков на основе файлов (например FileStreamSource , автозагрузчика) одновременно на одном драйвере может исчерпать ограничения дескриптора файлов на уровне пользователя, что может привести к сбоям случайного потока.
  • Перегрузка шины событий: большое количество одновременных потоковых запросов может вызвать перегрузку в шине событий одного сеанса StreamingQueryListener Spark. Все события (включая onQueryIdle) отправляются в эту общую шину, и большая очередь событий может существенно замедлить асинхронные обработчики onQueryProgress и повлиять на стабильность кластера.
  • Дорогостоящие операции на драйвере: Избегайте вызова collect() или других дорогостоящих операций DataFrame на драйвере без крайней необходимости, чтобы не материализовывать большие наборы результатов и не вызывать ошибки нехватки памяти (OOM).

Устранение неполадок, связанных с драйвером

Если у вас возникают сбои в работе драйвера из-за нехватки памяти (OOM) или конкуренции за ресурсы:

  1. Отслеживайте метрики драйверов в пользовательском интерфейсе Spark. Если вы видите высокий размер ЦП, памяти или диска, измените размер драйвера в параметрах вычислений кластера.
  2. Если проблемы сохраняются, убедитесь, что ваш код не выполняет на драйвере ресурсоемкие по памяти операции или пользовательские функции (UDF).
  3. Если вы больше не можете масштабировать драйвер по вертикали, Databricks настоятельно рекомендует распределить ваши задания между несколькими кластерами, чтобы обойти узкие места, связанные с масштабированием общего узла.

Размер исполнителя

При выполнении нескольких запросов в одном и том же кластере все запросы делят между собой слоты задач на исполнителях. Этапы одного запроса могут занимать доступные слоты, что приводит к задержкам и нехватке других запросов. Spark использует соотношение 1:1 между слотами задач и доступными ядрами. Убедитесь, что достаточно ядер доступны, если запросы должны выполняться одновременно.

Как правило, исполнители могут выполнять более интенсивные операции с памятью, чем узел драйвера. При необходимости настройте параметры выделения памяти JVM исполнителя и памяти вне кучи, чтобы справиться с нагрузкой приложения. Убедитесь, что узлы-исполнители правильно подобраны по объёму ресурсов — ЦП, памяти и дискового пространства — и при необходимости вертикально масштабируются. Если вертикальное масштабирование невозможно, рассмотрите возможность добавления дополнительных рабочих узлов в кластер.

Примечание.

Для некоторых из этих изменений может потребоваться перезапуск кластера.

Используйте пулы планировщика

Пулы планировщиков можно настроить для назначения вычислительных ресурсов запросам при выполнении нескольких потоковых запросов из одного исходного кода.

По умолчанию все запросы, запущенные в ноутбуке, выполняются в том же пуле планирования. Задания Apache Spark, создаваемые триггерами всех потоковых запросов в блокноте, запускаются один за другим в порядке «первым поступил — первым обработан» (FIFO). Это может привести к ненужным задержкам в запросах, так как они не эффективно обмениваются ресурсами кластера.

Пулы планировщиков позволяют объявлять, какие запросы структурированной потоковой передачи совместно используют вычислительные ресурсы.

В следующем примере query1 назначается выделенному пулу, тогда как query2 и query3 используют общий пул планировщика.

:::примечание о совместимости бессерверных серверов

Databricks рекомендует отказаться от использования spark.sparkContext, так как этот компонент несовместим с бессерверной вычислительной архитектурой Databricks. Вместо этого используйте напрямую spark (SparkSession). Пулы планировщика — это классическая концепция вычислений; при бессерверном режиме Databricks автоматически управляет масштабированием и выделением ресурсов.

:::

# Run streaming query1 in scheduler pool1
spark.sparkContext.setLocalProperty("spark.scheduler.pool", "pool1")
df.writeStream.queryName("query1").toTable("table1")

# Run streaming query2 in scheduler pool2
spark.sparkContext.setLocalProperty("spark.scheduler.pool", "pool2")
df.writeStream.queryName("query2").toTable("table2")

# Run streaming query3 in scheduler pool2
spark.sparkContext.setLocalProperty("spark.scheduler.pool", "pool2")
df.writeStream.queryName("query3").toTable("table3")

Примечание.

Конфигурация локального свойства должна находиться в той же ячейке записной книжки, где запускается запрос потоковой передачи.

Дополнительные сведения о пулах справедливого планировщика см. в документации Apache Spark по справедливому планировщику.

Рекомендации по запросу с отслеживанием состояния

Для запросов с отслеживанием состояния, выполняемых в одном кластере, помните следующее:

  • Используйте RocksDB в качестве поставщика хранилища состояний , чтобы избежать проблем OOM и приостановки GC. RocksDB — это поставщик хранилища состояний по умолчанию в Databricks Runtime 17.3 и выше. См. статью Настройка хранилища состояний RocksDB в Azure Databricks.
  • Настройте разделы перемешивания в соответствии с требованиями вашего приложения. Для этапов с сохранением состояния Spark планирует задачи пропорционально количеству разделов shuffle.
  • Ограничьте использование памяти RocksDB для каждого узла, чтобы избежать ошибок OOM из-за использования памяти вне кучи. Это обрабатывается автоматически в Databricks Runtime 17.3 и более поздних версиях, но требует ручной настройки в предыдущих выпусках. См. сведения об использовании памяти Cap RocksDB.
  • Избегайте размещения слишком большого количества разделов на одном узле-исполнителе. Операции обслуживания хранилища состояний, включая загрузку снимков и очистку, выполняются для каждого узла отдельно. Назначение слишком большого числа секций одному узлу-исполнителю может привести к дефициту ресурсов на обслуживание и увеличению времени восстановления из-за меньшего количества доступных полных снимков.