Соображения по использованию структурированной потоковой передачи в рабочей среде

Запустите рабочие нагрузки структурированной потоковой передачи в качестве запланированных заданий Lakeflow в Azure Databricks. Смотрите Задания Lakeflow.

Databricks рекомендует всегда настраивать следующие параметры:

  • Удалите ненужный код из записных книжек, который возвращает результаты, такие как display и count.
  • Не запускайте задачи Structured Streaming с использованием универсальных вычислительных ресурсов. Всегда запланируйте потоки в качестве заданий Lakeflow с помощью вычислений заданий.
  • Планируйте задания Lakeflow в режиме Continuous. Это относится к функции планирования заданий Azure Databricks, а не к интервалу срабатывания Структурированной потоковой передачи.
  • Не включите автомасштабирование для вычислений для заданий структурированной потоковой передачи.

Некоторые рабочие нагрузки пользуются следующими преимуществами:

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

Примечание.

Автоматическое масштабирование вычислений имеет ограничения, ограничивающие размер кластера для структурированных рабочих нагрузок потоковой передачи. Databricks рекомендует использовать декларативные конвейеры Spark в Lakeflow с расширенным автомасштабированием для потоковых рабочих нагрузок. См. раздел "Оптимизация использования кластера конвейера Lakeflow с автомасштабированием".

:::note Бессерверные вычисления

На бессерверных вычислительных ресурсах поддерживаются только Trigger.AvailableNow() и Trigger.Once(). Databricks рекомендует Trigger.AvailableNow().

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

См. ограничения потоковой передачи.

:::

Снижение задержки для операционного потокового потока

Операционные потоковые нагрузки принимают, преобразуют и действуют на данные почти в реальном времени. Распространённые примеры включают обнаружение мошенничества, выявление аномалий, персонализацию и мониторинг и оповещения в реальном времени, где задержка обработки напрямую влияет на бизнес-результаты. Низкая задержка для этих нагрузок обычно означает десятки или сотни миллисекунд, хотя многие команды устанавливают соглашения об уровне обслуживания (SLA) в диапазоне секунд, чтобы учесть вариабельность при более высоких процентилях.

Для минимальной сквозной задержки используйте режим реального времени, который обеспечивает сквозную задержку менее одной секунды в худшем случае и около 300 миллисекунд в типичных случаях. См. концепции режима в реальном времени.

Когда режим реального времени не подходит для вашей нагрузки, следующие лучшие практики снижают задержку для микропакетного структурированного потокового потока:

  • Режим вывода: используйте режим обновления, если ваши операторы запроса и приемник поддерживают его. Режим обновления выдаёт обновлённые строки после каждого срабатывания триггера и продолжает обновлять их до истечения срока действия watermark, поэтому конечный приёмник должен идемпотентно обрабатывать обновлённые результаты. Используйте режим добавления для рабочих нагрузок, которые режим обновления не поддерживает, например для соединений между потоками, или когда можно отбрасывать данные, поступающие с задержкой. Не используйте полный режим в сценариях с низкой задержкой. См. раздел Выбор выходного режима для структурированной потоковой передачи.
  • Триггер: Используйте processingTime триггер с 0 интервалом, который запускает следующую микропартию сразу после завершения предыдущей и появления новых данных. Это обеспечивает минимальную задержку микропакетов, но увеличивает стоимость API облачного хранилища. Не используйте AvailableNow, Once, или Continuous для операционных нагрузок. См. раздел "Настройка интервалов триггера структурированной потоковой передачи".
  • Водяная метка: Задайте достаточно большое значение водяной метки, чтобы включить данные, поступающие с задержкой, которые ваша рабочая нагрузка не должна отбрасывать. Водяная метка определяет, как долго запрос принимает данные о времени событий, поступающие не по порядку, прежде чем они будут отброшены и состояние будет вытеснено, поэтому слишком короткая водяная метка незаметно отбрасывает корректные поздно поступившие записи. В рамках этого ограничения более короткая водяная метка снижает задержку и требует хранения меньшего объёма состояния, а более длинная водяная метка допускает большее количество поздно поступающих данных ценой увеличения задержки и объёма хранимого состояния. Небольшой множитель целевого значения задержки по SLA, например 2×, — разумная отправная точка для подбора параметров. См. "Применение водяных знаков для управления порогами обработки данных".
  • Источники и приёмники: Считывайте данные из источников с низкой задержкой, таких как шины сообщений (Apache Kafka, Amazon Kinesis, Apache Pulsar или Google Cloud Pub/Sub) или потоки данных об изменениях из таблиц Delta Lake и Apache Iceberg. Записывайте в приёмники с низкой задержкой и высокой пропускной способностью, такие как шины сообщений, операционные базы данных или foreach приёмники. Проектируйте операции записи как идемпотентные, чтобы нижестоящие потребители могли обрабатывать дубликаты и данные, поступающие с задержкой.
  • Состояние и создание контрольных точек: Для запросов с сохранением состояния используйте хранилище состояний RocksDB, которое требуется как для создания контрольных точек по журналу изменений, так и для асинхронного создания контрольных точек состояния. Включите контрольные точки журнала изменений, чтобы сохранялись только инкрементальные изменения состояния. Если создание контрольных точек состояния становится узким местом, увеличивающим длительность пакетной обработки, включите асинхронное создание контрольных точек состояния, чтобы запись контрольных точек выполнялась параллельно со следующим микропакетом, предварительно ознакомившись с ограничениями, связанными с восстановлением после сбоев и масштабированием кластера. Дайте каждому запросу отдельную папку контрольных точек в надёжном облачном хранилище. См. Настройка хранилища состояния RocksDB в Azure Databricks, Асинхронное создание контрольных точек состояния для запросов с отслеживанием состояния и Контрольные точки Structured Streaming.
  • Управление смещением: Для снижения задержки при контрольных точках смещения в непрерывных потоках включите асинхронное отслеживание прогресса, которое обновляет логи смещения и фикса без блокировки обработки данных. Он несовместим с триггерами AvailableNow OR Once . См. Асинхронное отслеживание хода выполнения.
  • Хопы хранения: По возможности держите вычисления в одном потоковом конвейере. Разделение логики между несколькими задачами или конвейерами приводит к дополнительным обращениям к хранилищу, что увеличивает задержку.

Проектирование нагрузок стриминга с расчетом на сбои

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

Некоторые операции, такие как foreachBatch, предоставляют гарантию выполнения по крайней мере один раз, а не точно один раз. Для этих операций убедитесь, что конвейер процессинга идемпотентен. См. раздел Использование foreachBatch для записи в произвольные приемники данных.

Примечание.

При перезапуске запроса обрабатывается микропакет, запланированный во время предыдущего выполнения. Если задание завершилось сбоем из-за ошибки нехватки памяти или вы вручную отменили задание из-за слишком большого микро-пакета, может потребоваться увеличить масштаб вычислений, чтобы успешно обработать микро-пакет.

При изменении конфигураций между запусками эти конфигурации применяются к первому запланированному пакету. См. раздел «Восстановление после изменений в запросе структурированного потокового вещания».

При повторном выполнении задания

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

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

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

Вы также можете объединить эти стратегии. В следующей таблице сравниваются эти подходы.

Стратегия Множественные задачи Несколько запросов
Как используется общий доступ к вычислительным ресурсам? Databricks рекомендует развертывать вычислительные ресурсы, соответствующие каждой задаче потоковой обработки. Вы также можете совместно использовать вычислительные ресурсы между задачами. Все запросы используют одинаковые вычислительные ресурсы. При необходимости можно назначать запросы пулам планировщика.
Как обрабатываются повторные попытки? Все задачи должны завершиться сбоем перед повторными попытками выполнения задания. Задача повторяется, если любой запрос завершается ошибкой.

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

Настройка заданий структурированной потоковой передачи для перезапуска запросов потоковой передачи при сбое

Databricks рекомендует настраивать все задачи потоковой обработки с использованием непрерывного триггера. См . статью "Непрерывное выполнение заданий".

По умолчанию непрерывный триггер имеет следующее поведение:

  • Предотвращает более одного одновременного выполнения задачи.
  • Запускает новый запуск при сбое предыдущего запуска.
  • Использует экспоненциальную задержку для повторных попыток.

Databricks рекомендует всегда использовать вычислительные ресурсы для заданий вместо универсальных вычислений при планировании рабочих процессов. При сбое задания и повторных попытках развертываются новые вычислительные ресурсы.

Примечание.

Databricks рекомендует не использовать streamingQuery.awaitTermination() или spark.streams.awaitAnyTermination(). См. раздел "Когда следует использовать awaitTermination()".

Когда следует использовать awaitTermination()

streamingQuery.awaitTermination() и spark.streams.awaitAnyTermination() блокируйте текущий поток до завершения потокового запроса. Использование этих функций зависит от вашей среды выполнения.

Для заданий Lakeflow не используйте streamingQuery.awaitTermination() или spark.streams.awaitAnyTermination(). Эти функции не нужны, так как служба заданий автоматически предотвращает завершение выполнения, когда активен запрос потоковой передачи. Обе функции блокируют заполнение ячеек записной книжки и не позволяют службе заданий отслеживать потоковый запрос, что нарушает метрики невыполненной работы и уведомления о задании.

Используйте awaitTermination() в следующих случаях:

Сценарий использования Поведение
Интерактивные тетради для универсальных вычислений awaitTermination() поддерживает работу ячейки, позволяет наблюдать за состоянием запроса и обеспечивает отображение сбоев в выходных данных записной книжки.
Локальные среды и среды разработки При локальном запуске программы Spark процесс завершается после завершения основного потока. Вызов awaitTermination() для поддержания работоспособности программы до завершения или сбоя потокового запроса.
Распространение сбоя на драйвер Без awaitTermination()этого сбой потокового запроса в контексте, отличном от задания, может не распространяться в вызывающий поток. Запрос может завершиться сбоем без уведомления, что затрудняет обнаружение и диагностику проблемы. При вызове awaitTermination() в драйвере снова возникает исключение при выполнении запроса.