Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Important
Режим реального времени доступен только на Microsoft Fabric Spark Runtime 2.0 (Spark 4.1) и выше. Этот новый режим стриминга отсутствовал в ранних версиях Fabric.
Apache Spark 4.1 вводит Trigger.RealTime API для режима реального времени. Fabric поддерживает режим выполнения структурированной потоковой передачи с ультранизкой задержкой в Fabric Runtime 2.0. Он использует длительно выполняющиеся задачи, которые непрерывно считывают данные из потоковых источников, применяют поддерживаемые преобразования и записывают данные в поддерживаемые приёмники. Эта модель снижает задержку в расписании, потому что Spark не ждёт отдельной микропартии для начала каждой единицы работы.
Используйте режим реального времени, когда вашей рабочей нагрузке нужны новые события для быстрого прохождения запроса. Используйте стандартный режим микропакета, когда ставите приоритет на пропускную способность, широкую поддержку разъёмов или пакетно-ориентированные операции.
Режим реального времени против микропакетного режима
По умолчанию Structured Streaming использует микропакетную обработку. В микропакетном режиме Spark планирует, планирует и выполняет серию небольших партий заданий. Интервал триггера определяет, как часто Spark проверяет новые данные и запускает следующую партию. Этот подход хорошо подходит для большинства задач загрузки данных в lakehouse, обслуживания таблиц, агрегации и файловых приёмников.
Режим реального времени меняет паттерн выполнения. Spark запускает длительные задачи, которые остаются активными во время выполнения запроса. Процесс задач ведёт записи непрерывно, вместо того чтобы ждать следующей границы микропакета. Этот шаблон может снизить сквозную задержку для поддерживаемых запросов, особенно когда события должны перемещаться между стриминговыми системами с минимальной задержкой.
Компромисс заключается в том, что первоначальная реализация режима реального времени поддерживает меньший набор источников, поглотителей и операторов, чем микропакетный режим.
Как режим реального времени обрабатывает данные
Режим реального времени заменяет цикл коротких микропакетов одним длительно выполняющимся пакетом, который остаётся активным в течение всего интервала срабатывания триггера и обрабатывает записи по мере их поступления. Два варианта дизайна обеспечивают низкую задержку. Во-первых, Spark запускает все стадии выполнения запроса одновременно и связывает их с помощью shuffle в потоковом режиме, так что запись может проходить от источника через все преобразования к приёмнику без остановки на границе между стадиями. Во-вторых, поскольку задачи остаются в памяти, Spark не несет накладных расходов на планирование и диспетчеризацию каждого пакета при каждом срабатывании триггера.
Spark всё ещё нужна точка для сохранения прогресса. Он фиксирует смещения источника, обновляет хранилище состояния и публикует метрики потоковой обработки на границе между одним длительно выполняемым пакетом и следующим. Именно поэтому интервал срабатывания триггера ведёт себя иначе, чем интервал микропакетов: он определяет, как долго обрабатывается каждый пакет, прежде чем Spark создаст контрольную точку, а не сколько времени Spark ждёт перед началом обработки.
Настройте интервал с учётом такого поведения. Более длинный интервал контроля и метрики реже отображаются, поэтому сбой воспроизводит больше данных, а мониторинговый вид обновляется медленнее. При более коротком интервале контрольные точки создаются чаще, что ускоряет восстановление и делает метрики более актуальными, но дополнительная работа по созданию контрольных точек может увеличивать задержку. Универсального лучшего соотношения не существует, поэтому сравнивайте несколько интервалов с вашей собственной нагрузкой. Примеры в этой статье начинаются с интервала в пять минут.
Необходимые условия
Перед тем как включить режим реального времени, выполните следующие требования:
- Используйте Fabric Spark Runtime 2.0 (Spark 4.1) или новее. Режим реального времени недоступен в более ранних режимах.
- Запустите рабочую нагрузку в среде Fabric Spark, которая поддерживает Runtime 2.0, например в записной книжке или в определении задания Spark.
- Используйте потоковые источник и приемник, которые поддерживают режим реального времени в выбранной среде выполнения.
- Настройте надёжную точку контроля для запроса.
- Предоставьте достаточно слотов задач для запроса. Поскольку режим реального времени планирует все этапы одновременно, пулу нужно как минимум столько же доступных слотов задач, сколько общее количество задач на каждом этапе запроса. Планируйте по одному запросу в режиме реального времени на каждый пул, если только не убедитесь, что в пуле есть свободные слоты для дополнительных запросов.
Сведения о шагах по настройке Runtime 2.0 см. в разделе Runtime 2.0 в Fabric.
Включить режим реального времени
Включите режим реального времени, установив триггер потоковой передачи на Trigger.RealTime("<interval>") в writeStream. Этот интервал — это длительность длительно выполняемого пакетного задания, которая определяет частоту контрольных точек и сбора метрик, как описано в Как режим реального времени обрабатывает данные.
Important
Режим реального времени поддерживает только update режим вывода. Режимы вывода append и complete не поддерживаются. Установите .outputMode("update") для запроса.
Note
Триггер в реальном времени доступен в API потоковой передачи JVM. В первоначальной реализации Spark 4.1 PySpark не экспонирует его нативно через DataStreamWriter.trigger().
Примеры предполагают, что input_options, output_options, и c_path содержат ваши исходные опции, опции sink и путь контрольной точки. Запускайте эти примеры только на Fabric Spark Runtime 2.0 (Spark 4.1) или новее.
Сначала создайте потоковый DataFrame на Python.
passthrough = (
spark.readStream
.format("kafka")
.options(**input_options)
.load()
.selectExpr("CAST(key AS STRING) AS key", "CAST(value AS STRING) AS value")
)
Затем примените временный обходной путь для моста JVM. PySpark пока не предоставляет доступа к Trigger.RealTime в DataStreamWriter.trigger(), поэтому это обходное решение создаёт триггер JVM через py4j и применяет его к базовому объекту Java DataStreamWriter. Замените этот код на родной API триггера PySpark после завершения поддержки.
# Create RTM trigger object via JVM bridge
rtm_trigger = (
spark._jvm.org.apache.spark.sql.streaming
.Trigger.RealTime("5 minutes")
)
# Define streaming query
query = (
passthrough.writeStream
.format("kafka")
.options(**output_options)
.option("checkpointLocation", c_path)
.outputMode("update")
)
# Apply RTM trigger via Java DataStreamWriter
query._jwrite = query._jwrite.trigger(rtm_trigger)
# Start streaming query
query = query.start()
query.awaitTermination()
Запросы к проектированию для режима реального времени
Режим реального времени проверяет источник, погрузку, выходной режим и план запроса при начале запроса. В начальной реализации он поддерживает сфокусированный набор потоковых паттернов.
Лучшая ментальная модель — это путь с низкой задержкой, который перемещает события между системами потоковых сообщений и применяет лёгкую обработку по пути. Постройте свой запрос вокруг этой модели:
- Системы чтения и записи потоковых сообщений. Режим реального времени создан для конечных точек, совместимых с Kafka, включая Apache Kafka и Центры событий Azure через разъём Kafka. Чтобы отправить результаты напрямую во внешнюю систему, вызовите
.foreach(...)с помощьюForeachWriter. - Оставьте lakehouse и файловый ввод-вывод на стандартном триггере. Таблицы Delta, таблицы lakehouse, а также файловые источники и приемники не поддерживаются в качестве целевых объектов в режиме реального времени. Вместо этого загружайте данные в Lakehouse, поддерживайте таблицы или записывайте файлы с триггером микропакетов, а режим Real-time оставьте для чувствительного к задержкам этапа передачи между системами обмена сообщениями.
- Предпочитайте лёгкие трансформации. Проекции, отбросы, фильтры, выражения столбцов и простое обогащение обеспечивают минимальную задержку. Чем больше состояний и перемешивания данных включает запрос, тем важнее сначала протестировать его в режиме реального времени.
Когда вы используете конечную точку Event Hubs, совместимую с Kafka, согласуйте интервалы простоя и обновления метаданных клиента Kafka с тайм-аутом простоя Event Hubs. Несоответствие значений тайм-аута может вызвать задержку при повторном подключении вблизи границы длительно выполняющейся пакетной операции. Для рекомендуемых вариантов источника и назначения см. «Обеспечение работоспособности подключений, совместимых с Kafka».
Режим реального времени поддерживает множество операторов с состоянием, включая оконные агрегации, дедупликацию и объединения, но поддерживаемые варианты в нём более ограничены, чем в микропакетном режиме. Используйте следующие матрицы в качестве краткого справочника при разработке запроса и подтверждайте текущее поведение в используемой версии среды выполнения.
Матрица коннекторов наглядно показывает, какое место режим реального времени занимает в конвейере:
| Connector | В качестве источника | Как раковина |
|---|---|---|
| Конечные точки, совместимые с Kafka (Apache Kafka, Центры событий Azure через соединитель Kafka) | Поддерживается | Поддерживается |
Кастомный провал сквозь .foreach(...) и ForeachWriter |
Неприменимо | Поддерживается |
| Таблицы Delta и lakehouse | Не поддерживаются | Не поддерживаются |
| Форматы на основе файлов | Не поддерживаются | Не поддерживаются |
Операторная матрица группирует преобразования по тому, насколько хорошо они соответствуют модели с низкой задержкой:
| Transformation | Настройка режима реального времени | Что делать |
|---|---|---|
| Операции без сохранения состояния: выборка, фильтрация, приведение типов, проецирование, скалярные пользовательские функции | Поддерживается | Предпочтительный путь; Даёт минимальную задержку |
| Агрегации и кувыркающиеся или раздвижные окна | Поддерживается | Добавить водяной знак для связанного состояния |
| Дедупликация с водяным знаком или внутри него | Поддерживается | Выбирайте ключи, которые определяют логическое событие |
| Окна сессии (с перерывами) | Не поддерживаются | Запустите запрос на микропакетном триггере |
| Соединение потока со справочником | Поддерживается, когда опорная сторона выполняет широковещательную передачу | Держите эталонный набор небольшим |
| Присоединение потока к потоку | Только внутреннее объединение с дополнительной настройкой | Избегайте внешних соединений между двумя ручьями |
Обычная ситуация с transformWithState |
Поддерживается разной семантикой | См. примечание, приведённое ниже |
Операторы разбиения за раз: mapPartitions, mapInPandas, mapInArrow |
Не поддерживаются | Переписывайте с помощью UDF на уровне строк, фильтров или выражений комплексного типа |
| Самообъединение, объединение с пакетным источником или объединение после оператора состояния | Не поддерживаются | Используйте независимые входные потоки и применяйте union перед операциями с сохранением состояния |
Поддержка не означает, что операция сохраняет минимальную задержку. Агрегации, объединения, дедупликация, пользовательское состояние и Python UDF добавляют управление состоянием, тасовки, сериализацию или буферизацию. Проводите тестирование производительности всего запроса с реалистичной интенсивностью входного потока и размерами состояния, а не измеряйте запрос без сохранения состояния, который только пропускает данные.
Important
Если Spark не поддерживает источник, приёмник, операцию или режим вывода в режиме реального времени, вместо этого запустите этот запрос со стандартным триггером микропакетной обработки. Не думайте, что запрос, работающий в микропакетном режиме, работает и в режиме реального времени.
Если вы создаёте пользовательский процессор с сохранением состояния с transformWithState, имейте в виду, что модель построчного выполнения изменит своё поведение. Режим реального времени передаёт события процессору по мере их поступления, а не по клавишам, поэтому записывайте процессор, не предполагая, что он видит каждую строку ключа в одном вызове. Таймеры по времени событий не поддерживаются, а таймеры обработки могут быть отложены до прибытия данных или окончания длительной партии. Вариант API на основе pandas недоступен, поэтому используйте строковый API. Общая модель программирования описана в статье Обработка потоков с сохранением состояния.
Случаи использования
Режим реального времени работает лучше всего, когда задержка важнее максимальной пропускной способности или широкого покрытия функций. Рассмотрите его для таких нагрузок, как маршрутизация событий, операционные оповещения, обогащение лёгких потоков и перемещение с низкой задержкой между потоковыми системами.
Стандартный режим микропакетов обычно лучший выбор, когда вам нужно:
- Высокопроизводительная загрузка данных в таблицу Delta в lakehouse.
- Большие партии для оптимизации размера файлов и меньших накладных расходов на запись.
- Сложные агрегации, соединения или обработка состояний.
- Широкая совместимость с разъёмами и приёмниками.
- Ниже затраты за счёт меньшей постоянной вычислительной нагрузки.
Выбирайте режим для зависимости от нагрузки. Рабочее пространство может использовать режим реального времени для потоков событий, чувствительных к задержке, и микропакетный режим для загрузки данных в Lakehouse или для аналитических конвейеров.
Мониторинг запросов в реальном времени
Отслеживайте запросы в режиме реального времени в центре мониторинга Fabric. Откройте приложение Spark и используйте вкладку Structured Streaming , чтобы просмотреть метрики потока, такие как скорость ввода, скорость обработки, входные строки, статус запроса и продолжительность работы.
Режим реального времени публикует ход выполнения запросов и контрольные точки между длительно выполняемыми пакетами. Вкладка Structured Streaming и lastProgress обновляются после завершения пакета, поэтому их актуальность зависит от интервала срабатывания триггера режима реального времени. Не интерпретируйте длительность пакетной обработки или возраст lastProgress как задержку обработки отдельной записи.
Используйте дескриптор запроса, чтобы проверить, активен ли запрос, проверить его текущее состояние и получить самые свежие метрики завершённой партии:
import json
monitoring_snapshot = {
"isActive": query.isActive,
"status": query.status,
"lastProgress": query.lastProgress,
}
print(json.dumps(monitoring_snapshot, indent=2, default=str))
if not query.isActive and query.exception() is not None:
raise query.exception()
На странице сведений о приложении Spark просмотрите активные этапы и экзекьюторы во время выполнения длительного пакетного задания. Ищите неудачные задачи, недоступных исполнителей задач, длительную высокую нагрузку на ЦП или память, а также нехватку слотов для задач.
Одних метрик прогресса Spark недостаточно для оперативного мониторинга, так как они обновляются на границе пакета. Отслеживайте отставание группы потребителей Kafka или эквивалентное отставание источника, скорости входящего и исходящего потоков, дросселирование и ошибки в исходной и целевой системах. Для Центры событий Azure используйте Azure Monitor вместе с Fabric monitoring hub.
Измерьте сквозную задержку вне события прогресса Spark. Включите в каждую запись временную метку события, заданную производителем, затем сравните её со временем, в которое потребитель в последующем звене получает результат. Этот показатель учитывает задержку источника, обработку в Spark, передачу по сети и задержку приёмника. Оповещение о завершении запроса, устаревании контрольных точек, растущем отставании источника, ошибках приёмника и сквозной задержке, превышающей целевое значение для рабочей нагрузки.
Связанные материалы
- Прочитайте обзор структурированного стриминга.
- Проверьте триггеры и режимы вывода.
- Применяйте лучшие практики структурированного стриминга.
- Используйте распространённые шаблоны.
- Настройте Runtime 2.0 в Fabric.
- Узнайте о потоке данных в домик на озере с помощью Spark.