Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Spark Structured Streaming поддерживает как операции без состояния, так и операции с состоянием. Операции без сохранения состояния, такие как map, filter и select, обрабатывают каждую входную запись независимо. Операции с состоянием, такие как агрегации, соединения и дедупликация, сохраняют информацию между микропакетами, чтобы Spark мог обновлять результаты при поступлении новых событий.
В Microsoft Fabric вы используете те же концепции Apache Spark Structured Streaming с открытым исходным кодом в сочетании со средами выполнения Fabric Spark, записными книжками, определениями заданий Spark, таблицами lakehouse и другими возможностями Fabric для data engineering.
Государственный магазин и контрольно-пропускные пункты
Потоковые запросы с состоянием используют хранилище состояний для хранения промежуточных данных между микропартиями. Например, текущий подсчёт для каждого устройства должен запоминать текущее значение счётчика для каждого ключа устройства, прежде чем сможет обработать следующую микропартию.
Spark сохраняет состояние с помощью контрольной точки потока. Контрольная точка отслеживает смещения, ход запросов и данные хранилища состояний. Если запрос останавливается и перезапускается с той же точкой проверки, Spark перезагружает состояние и возобновляет процесс с последнего закреплённого прогресса.
Important
Используйте надёжную точку контроля и не делите одну и ту же точку между разными потоковыми запросами. Схема состояния и план запроса являются частью состояния контрольной точки, поэтому значительные изменения логики состояния часто требуют новой контрольной точки.
Состояние растёт по мере того, как Spark отслеживает всё больше ключей, окон или буферизованных строк объединения. Ограничьте объём состояния с помощью водяных меток, параметра time-to-live (TTL) и тщательно подобранных ключей, чтобы чекпоинт не разрастался без ограничений.
Включите хранилище состояний RocksDB
По умолчанию Spark использует хранилище состояний с поддержкой HDFS, рабочее состояние которого поддерживается в памяти, управляемой JVM, а надёжные версии хранятся в контрольной точке. Включите провайдера хранилища состояний RocksDB, чтобы Spark управлял активным состоянием в нативной памяти и на локальном диске каждого исполнителя вместо кучи JVM. RocksDB снижает нагрузку на кучу JVM и обеспечивает более предсказуемое использование памяти. Это напрямую полезно для запросов с сохранением состояния и является разумным вариантом по умолчанию даже для запросов с небольшим объёмом состояния.
spark.conf.set(
"spark.sql.streaming.stateStore.providerClass",
"org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider"
)
Установите провайдера state store перед первым запуском запроса и поддерживайте его согласованность между перезапусками, потому что это часть запроса с контрольной точкой.
Включить контрольные точки журнала изменений RocksDB
Контрольная точка журнала изменений — это оптимизация Apache Spark для хранилища состояния RocksDB. Он выгружает инкрементальные изменения состояния и периодически создаёт снимки состояния, что может снизить задержку при создании контрольных точек для больших объёмов состояния.
spark.conf.set(
"spark.sql.streaming.stateStore.rocksdb.changelogCheckpointing.enabled",
True
)
Создание контрольных точек на основе журнала изменений обратно совместимо с контрольными точками на основе снимков RocksDB. Вы можете включить эту функцию для существующего запроса RocksDB, не удаляя состояние, но чтобы настройка вступила в силу, необходимо перезапустить запрос.
Оконные агрегации
Оконные агрегации группируют события по времени событий, а не по времени обработки. Они помогают рассчитывать такие метрики, как подсчёты, суммы и средние значения в промежутках времени.
Распространённые типы окон включают:
-
Падающие окна. Каскадное окно имеет фиксированный размер и не перекрывает другие окна. Например, пятиминутное кувыркающееся окно подсчитывает события от
10:00до10:05, а затем начинает новый подсчёт для10:05до10:10. - Раздвижные окна. Раздвижное окно имеет фиксированный размер и начинается с регулярного интервала слайда. Например, 10-минутное окно, которое скользит каждую минуту, даёт скользящий счёт за 10 минут.
- Окна сессий. Сеансовое окно объединяет события, которые поступают с небольшим интервалом и имеют один и тот же ключ. Например, окно сессии с перерывом в 15 минут закрывает пользовательскую сессию после того, как у этого пользователя нет событий в течение 15 минут.
Следующий пример подсчитывает события устройств в пятиминутных неперекрывающихся окнах.
from pyspark.sql import functions as F
windowed_counts = (
events
.withWatermark("event_time", "10 minutes")
.groupBy(
F.window("event_time", "5 minutes"),
F.col("device_id")
)
.count()
)
Для раздвижных окон добавьте длительность слайда. Например, window("event_time", "10 minutes", "1 minute") рассчитывает окно в 10 минут каждую минуту. Для окон сессий используйте session_window(event_time, '15 minutes') в Spark SQL или session_window("event_time", "15 minutes") в API DataFrame.
Агрегации окна сессии не поддерживают update режим вывода. Используйте режим вывода, поддерживаемый запросом и приемником, и проверьте полученное время излучения с поздними данными.
Watermarks
Водяной знак показывает Spark, насколько поздно вы ожидаете поступления данных о времени событий. Spark использует водяной знак, чтобы определить, когда можно удалить старое состояние, а поздние строки слишком старые для обновления результата.
Используйте withWatermark перед операцией со временем состояния события:
Этот пример позволяет данным поступать с опозданием до 10 минут в зависимости от колонки event_time . Искра может со временем потерять состояние для окон, которые старше водяного знака.
Note
Spark не отбрасывает данные, поступающие в пределах настроенной задержки водяного знака. Данные старше водяной метки всё ещё могут быть обработаны, но Spark не гарантирует этого. Очистка состояния также происходит асинхронно, а не сразу при продвижении водяной метки.
Выбирайте задержки временных меток на основе фактического поведения источника. Слишком короткая задержка может привести к потере корректных данных, поступивших с опозданием. Слишком длинная задержка сохраняет больше состояния и увеличивает размер контрольной точки.
Водяные метки обновляются, когда Spark обрабатывает новые входные данные. Если данные не поступают, водяной знак может не продвинуться, поэтому вывод оконных результатов, несопоставленных строк внешнего соединения и очистка состояния могут быть отложены до тех пор, пока в одной из последующих микропартий не появятся данные.
Соединения потоков
Операция соединения двух потоков объединяет два входных потока данных. Spark буферизует строки с обеих сторон, пока не сможет определить, поступят ли совпадающие строки.
Внутренние соединения могут работать без водяных знаков, но неограниченное состояние может расти бесконечно. Добавьте водяные метки и условие диапазона времени события, чтобы Spark мог удалить старые буферизованные строки.
joined_events = (
impressions.withWatermark("impression_time", "10 minutes")
.join(
clicks.withWatermark("click_time", "10 minutes"),
"""
impression_id = click_impression_id AND
click_time >= impression_time AND
click_time <= impression_time + interval 5 minutes
""",
"inner"
)
)
Внешние присоединения требуют достаточного количества информации о событиях, чтобы Спарк знал, когда будущий матч не может появиться. Используйте водяные знаки и условие диапазона времени. Для левого внешнего соединения поставьте водяной знак на правой стороне, которая может привести к нулируемому совпадению. Для правого внешнего соединения поставьте водяной знак на левой стороне. Для полного внешнего соединения задайте водяную метку хотя бы для одной стороны; задайте водяные метки для обеих сторон, когда Spark требуется очищать состояние для обоих входных потоков.
Important
Избегайте широких соединений между потоками без временного диапазона. Spark должен хранить в буфере больше строк, а состояние может расти без чётко определённого момента очистки.
Несовпадающие ряды внешних соединений не испускаются сразу. Спарк ждёт, пока водяной знак и условие временного диапазона не докажут, что будущий матч больше не может состояться.
Когда запрос объединяет несколько входных потоков, Spark выводит один водяной знак запроса из входных водяных знаков. По умолчанию самый медленный ввод контролирует прогресс, чтобы Spark не сбрасывал данные из этого потока преждевременно. Таким образом, зависший входной поток может задержать очистку состояния и выдачу результатов по всему запросу.
Дедупликация в потоках
Дедупликация потока удаляет повторяющиеся события, запоминая ключи, которые Spark уже обработала. Используйте его, когда источники могут повторно отправить одно и то же событие, например при повторных попытках доставки из системы обмена сообщениями.
dropDuplicates сохраняет значения столбцов, используемые для обнаружения дублиратов. Без водяного знака Spark должен помнить эти значения на протяжении всего времени выполнения запроса. Чтобы Spark удалял устаревшее состояние dropDuplicates, добавьте столбец с водяным знаком в ключ дедупликации.
deduplicated_events = (
events
.withWatermark("event_time", "10 minutes")
.dropDuplicates(["event_id", "event_time"])
)
Включение времени события означает, что две записи с одинаковым идентификатором события, но разные временные метки не являются дубликатами. Используйте dropDuplicatesWithinWatermark тогда, когда ваша среда выполнения поддерживает это и только идентификатор события определяет дубликат. Spark затем сохраняет идентификатор каждого события в пределах горизонта водяного знака:
deduplicated_events = (
events
.withWatermark("event_time", "10 minutes")
.dropDuplicatesWithinWatermark(["event_id"])
)
Выбирайте ключи, которые уникально определяют логическое событие. Если ключ слишком широк, Spark может убрать отдельные события. Если ключ слишком узкий, через него могут пройти дубликаторы.
Пользовательская логика состояния
Встроенные агрегации, объединения и дедупликация подходят для многих задач с хранением состояния. Используйте произвольную обработку с сохранением состояния, когда требуется собственная логика для каждого ключа, например автоматы состояний, подавление оповещений, обработка тайм-аутов или многоэтапное обогащение.
Оператор transformWithState заменяет устаревшие API mapGroupsWithState и flatMapGroupsWithState для пользовательской логики состояния. Вы группируете строки по ключу, реализуете процессор с состоянием и управляете переменными состояния, таймерами, выходным режимом и временем. Apache Spark представил transformWithState в Spark 4.0, а Fabric Runtime 2.0 включает эту возможность вплоть до Spark 4.1.
Note
transformWithState это более новый API. Проверьте свой код с помощью версии Fabric runtime, которая выполняет вашу производственную нагрузку.
Используйте TTL для пользовательского состояния, когда бизнес-логика позволяет истечь срока действия. TTL помогает Spark удалять устаревшие ключи без ожидания конкретного события ввода.
Инструменты отладки состояния
Fabric Runtime 2.0 включает Apache Spark 4.1 и State Data Source for Structured Streaming, который Apache Spark внедрил в Spark 4.0. Используйте его, чтобы проверить содержимое хранилища состояния по контрольной точке с помощью отдельного пакетного запроса.
Note
Государственный источник данных является экспериментальным в Apache Spark. Его опции и схема вывода могут меняться.
Средство чтения statestore может проверять совместимое состояние в существующей контрольной точке, включая контрольную точку, созданную в более ранней версии среды выполнения Spark. Запустите отдельный пакетный считыватель с Fabric Runtime 2.0 или новее; вам не нужно сначала создавать контрольную точку или перезапускать исходный запрос.
state_df = (
spark.read
.format("statestore")
.load("Files/checkpoints/device-counts")
)
state_df.printSchema()
display(state_df)
Исходник state-metadata — это отдельное удобство для поиска идентификаторов операторов, названий магазинов и доступных пакетных идентификаторов. Spark создаёт эти метаданные только во время запуска потокового запроса на Spark 4.0 или новее. Для старой контрольной точки возобновите исходный запрос к этой точке в Fabric Runtime 2.0 перед использованием state-metadata.
state_metadata_df = (
spark.read
.format("state-metadata")
.load("Files/checkpoints/device-counts")
)
display(state_metadata_df)
Для запросов с несколькими операторами состояния используйте метаданные, чтобы определить оператор, который вы хотите проверить. Если вы уже знаете ID оператора и название магазина, можно пропустить источник метаданных и передать эти опции напрямую ридеру statestore . Для transformWithState, укажите имя переменной состояния при чтении состояния.
Important
Рассматривайте состояние контрольной точки как операционные данные. Не редактируйте файлы контрольных точек вручную. Используйте источник данных состояния для проверки и используйте изменения кода запроса, водяные знаки или TTL для контроля поведения состояния.
Лучшие практики для ограниченного состояния
Следуйте этим рекомендациям, чтобы обеспечить надёжность запросов потоковой обработки с сохранением состояния:
- Добавьте водяные знаки по времени события для оконных агрегирований, соединений поток–поток и устранения дубликатов, если это допускается правилами обработки поздних данных.
- Используйте условия диапазона времени для соединений поток-поток, чтобы Spark мог удалять буферизованные строки.
- Выбирайте ключи с подходящей кардинальностью. Ключи с очень высокой кардинальностью увеличивают объём состояния, тогда как чрезмерно широкие ключи могут смешивать несвязанные события.
- Включите поставщик хранилища состояния RocksDB и контрольные точки журнала изменений, чтобы снизить нагрузку на кучу, паузы сборки мусора и задержку создания контрольных точек.
- Используйте TTL с
transformWithState, если настраиваемое состояние не должно храниться вечно. - Сохраняйте позиции контрольных точек стабильными на протяжении всего срока действия запроса и используйте новую контрольную точку для изменений схем состояния несовместимых состояний.
- Отслеживайте показатели, связанные с состоянием, скорость входа, скорость обработки и рост контрольных точек во время производственных циклов.
- Протестируйте запрос с сохранением состояния на реалистичных сценариях с поздними данными, дублирующимися данными и перезапуском перед его развертыванием.