Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
На этой странице описываются запросы состояние-ориентированной структурированной потоковой передачи, включая состояние-ориентированные операции, рекомендации по оптимизации, цепочку нескольких состояние-ориентированных операторов и перебалансирование состояния.
Запрос структурированной потоковой передачи с отслеживанием состояния требует добавочных обновлений сведений о промежуточном состоянии, тогда как запрос структурированной потоковой передачи без отслеживания состояния отслеживает только сведения о том, какие записи были обработаны из источника в приемник. Функции оптимизации, доступные для запросов без отслеживания состояния, см. в разделе "Оптимизация потоковых запросов без отслеживания состояния".
Операции с сохранением состояния
Операции с отслеживанием состояния включают агрегирование потоков, distinct, dropDuplicates соединения между потоками и пользовательские приложения с отслеживанием состояния.
Промежуточная информация о состоянии, необходимая для запросов структурированной потоковой передачи с отслеживанием состояния, может привести к непредвиденным задержкам и проблемам в производстве при неправильной настройке.
В Databricks Runtime 13.3 LTS или более поздних версиях вы можете включить контрольную точку журнала изменений с использованием RocksDB, чтобы уменьшить длительность контрольной точки и общую задержку в рабочих нагрузках с использованием Structured Streaming. Databricks рекомендует включать контрольные точки журнала изменений для всех запросов структурированных потоков с отслеживанием состояния. См. раздел "Включить контрольную точку журнала изменений".
Оптимизация запросов структурированной потоковой передачи с отслеживанием состояния
Databricks рекомендует следующее для запросов структурированной потоковой передачи с сохранением состояния:
- Используйте оптимизированные для вычислений экземпляры в качестве работников.
- Установите количество секций перетасовки в 1–2 раза от числа ядер в кластере.
Внимание
Число разделов перетасовки фиксировано на момент создания контрольной точки. Изменение spark.sql.shuffle.partitions не влияет на потоковый запрос, который уже имеет контрольную точку, — запрос продолжает использовать исходное число секций. Чтобы применить новое число секций, необходимо запустить запрос с новым расположением контрольной точки.
В Databricks Runtime 18.0 и более поздних версиях запросы потоковой передачи без состояния поддерживают изменения динамического перераспределения разделов без необходимости создания новой контрольной точки.
В Databricks Runtime 18 LTS и выше можно изменить количество разделов для запросов с сохранением состояния без потери состояния контрольной точки. См. Переразбиение состояния по запросу для потоковых запросов с сохранением состояния.
- Задайте для конфигурации
spark.sql.streaming.noDataMicroBatches.enabledзначениеfalseв SparkSession. Это предотвращает потоковый микробатч-движок от обработки микробатчей без данных. Настройка этой конфигурацииfalseможет также привести к тому, что операции с отслеживанием состояния, которые используют водяные знаки или время ожидания обработки, не будут выводить данные до появления новых данных.
Databricks рекомендует использовать RocksDB с контрольными точками журнала изменений для управления состоянием потоков. См. статью Настройка хранилища состояний RocksDB в Azure Databricks.
Примечание.
Невозможно изменить схему управления состоянием между перезапусками запросов. Если запрос был запущен с помощью управления по умолчанию, необходимо перезапустить его с нуля с новым расположением контрольной точки, чтобы изменить хранилище состояний.
Работа с несколькими операторами с состоянием в структурированной потоковой передаче
В Databricks Runtime 13.3 LTS или более поздней версии Azure Databricks предлагает расширенную поддержку операторов с сохранением состояния в рабочих нагрузках структурированной потоковой передачи. Можно объединить несколько операторов с отслеживанием состояния, то есть вы можете передать выходные данные операции, например, окно агрегирования, в другую операцию с отслеживанием состояния, например, соединение.
В Databricks Runtime 16.2 или более поздней версии можно использовать transformWithState в рабочих нагрузках с несколькими операторами с отслеживанием состояния. См. Создать пользовательское приложение с состоянием с transformWithState.
В следующих примерах показано несколько шаблонов, которые можно использовать.
Внимание
При работе с несколькими операторами с отслеживанием состояния существуют следующие ограничения:
- Устаревшие пользовательские сохраняющие состояние операторы (
FlatMapGroupWithStateиapplyInPandasWithState) не поддерживаются. - Поддерживается только режим добавления в вывод.
Агрегирование связанных временных окон
Питон
words = ... # streaming DataFrame of schema { timestamp: Timestamp, word: String }
# Group the data by window and word and compute the count of each group
windowedCounts = words.groupBy(
window(words.timestamp, "10 minutes", "5 minutes"),
words.word
).count()
# Group the windowed data by another window and word and compute the count of each group
anotherWindowedCounts = windowedCounts.groupBy(
window(window_time(windowedCounts.window), "1 hour"),
windowedCounts.word
).count()
язык программирования Scala
import spark.implicits._
val words = ... // streaming DataFrame of schema { timestamp: Timestamp, word: String }
// Group the data by window and word and compute the count of each group
val windowedCounts = words.groupBy(
window($"timestamp", "10 minutes", "5 minutes"),
$"word"
).count()
// Group the windowed data by another window and word and compute the count of each group
val anotherWindowedCounts = windowedCounts.groupBy(
window($"window", "1 hour"),
$"word"
).count()
Агрегирование временного окна в двух разных потоках с последующим соединением окон между потоками
Питон
clicksWindow = clicksWithWatermark.groupBy(
clicksWithWatermark.clickAdId,
window(clicksWithWatermark.clickTime, "1 hour")
).count()
impressionsWindow = impressionsWithWatermark.groupBy(
impressionsWithWatermark.impressionAdId,
window(impressionsWithWatermark.impressionTime, "1 hour")
).count()
clicksWindow.join(impressionsWindow, "window", "inner")
язык программирования Scala
val clicksWindow = clicksWithWatermark
.groupBy(window("clickTime", "1 hour"))
.count()
val impressionsWindow = impressionsWithWatermark
.groupBy(window("impressionTime", "1 hour"))
.count()
clicksWindow.join(impressionsWindow, "window", "inner")
Соединение потоков по временным интервалам, за которым следует агрегирование по временным окнам
Питон
joined = impressionsWithWatermark.join(
clicksWithWatermark,
expr("""
clickAdId = impressionAdId AND
clickTime >= impressionTime AND
clickTime <= impressionTime + interval 1 hour
"""),
"leftOuter" # can be "inner", "leftOuter", "rightOuter", "fullOuter", "leftSemi"
)
joined.groupBy(
joined.clickAdId,
window(joined.clickTime, "1 hour")
).count()
язык программирования Scala
val joined = impressionsWithWatermark.join(
clicksWithWatermark,
expr("""
clickAdId = impressionAdId AND
clickTime >= impressionTime AND
clickTime <= impressionTime + interval 1 hour
"""),
joinType = "leftOuter" // can be "inner", "leftOuter", "rightOuter", "fullOuter", "leftSemi"
)
joined
.groupBy($"clickAdId", window($"clickTime", "1 hour"))
.count()
Перебалансирование состояния для структурированной потоковой передачи
Перебалансировка состояния включена по умолчанию для всех заданий потоковой обработки в конвейерах Lakeflow. В Databricks Runtime 11.3 LTS или более поздней версии можно задать следующий параметр конфигурации в конфигурации кластера Spark, чтобы включить перебалансирование состояния:
spark.sql.streaming.statefulOperator.stateRebalancing.enabled true
Перебалансировка состояния приносит пользу конвейерам управления состоянием в структурированной потоковой передаче, которые проходят через изменения размера кластера. Операции потоковой передачи без отслеживания состояния не получают преимущества независимо от изменения размеров кластера.
Примечание.
Автоматическое масштабирование вычислений имеет ограничения, ограничивающие размер кластера для структурированных рабочих нагрузок потоковой передачи. Databricks рекомендует использовать декларативные конвейеры Spark в Lakeflow с расширенным автомасштабированием для потоковых рабочих нагрузок. См. раздел "Оптимизация использования кластера конвейера Lakeflow с автомасштабированием".
События изменения размера кластера приводят к перебалансировке состояния. Микропакеты могут иметь более высокую задержку во время перебалансировки событий, так как состояние загружается из облачного хранилища в новые исполнители.