Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Spark Structured Streaming — это движок инкрементальной обработки, построенный на базе Apache Spark. Он моделирует входной поток как неограниченную таблицу, к которой добавляются новые строки. Вы определяете логику, используя декларативные API Dataset и DataFrame, а затем выбираете, будет ли запрос работать непрерывно или обрабатывать доступные данные и останавливаются.
В Microsoft Fabric Structured Streaming поддерживает как непрерывную обработку событий, так и запланированные дополнительные нагрузки. Вы можете использовать тот же запрос и контрольную точку для обработки событий по мере их поступления или периодически выполнять запрос с помощью триггера доступного момента. В обоих паттернах Spark отслеживает прогресс исходника и поддерживает состояние для агрегирований, дедупликации, объединений и пользовательской логики состояния.
Основная модель обработки
Запрос Structured Streaming состоит из трёх основных частей: источник, преобразования и приёмник.
- Источник: Источник считывает новые входные записи из стриминговой системы или постепенно обнаруженные файлы.
- Преобразования: Преобразования определяют логику запроса. Можно фильтровать строки, выбирать столбцы, анализировать JSON, объединять данные, агрегировать значения и обогащать события справочными данными.
- Приёмник: Приёмник записывает выходные данные в хранилище, память или другое место.
Эта модель «источник-поглотитель» декларативна. Вы описываете желаемый результат, а Spark планирует повторное выполнение по мере появления новых данных.
Потоковые датафреймы используют привычный API Spark DataFrame, но не каждая пакетная операция поддерживает потоковые входные данные. Spark проверяет план запроса при запуске потокового запроса и сообщает о неподдерживаемых операциях.
Используйте один API для непрерывной и запланированной обработки
Structured Streaming отделяет определение запроса от того, как долго выполняется запрос. Эта конструкция позволяет использовать одинаковые readStream API writeStream для двух рабочих паттернов:
| Рисунок | Срок службы запроса | Типичное использование |
|---|---|---|
| Постоянная трансляция | Запрос остаётся активным и обрабатывает данные по мере их поступления. | Обработка событий, непрерывный прием данных, оперативные оповещения и конвейеры с низкой задержкой. |
| Запланированная инкрементальная обработка | Триггер доступного сейчас обрабатывает все доступные данные при начале запроса, а затем останавливается. Планировщик позже повторно запускает запрос с той же контрольной точкой. | Инкрементальный ETL, периодическая загрузка файлов, пакетная обработка с сохранением состояния, дозагрузка исторических данных и задачи наверстывания. |
Обе схемы используют контрольные точки потока. Для запланированного задания типа available-now чекпоинт сохраняет сведения о том, какие исходные данные уже были обработаны, и сохраняет состояние между запусками. Следующий запуск обрабатывает только новые входы, продолжая существующие вычисления состояния.
Используйте это уже сейчас, если вам нужно пакетное планирование заданий без необходимости самостоятельно заново реализовывать отслеживание поэтапного прогресса или управление состоянием. Используйте триггер по умолчанию, фиксированный интервал или режим реального времени, когда нагрузка должна оставаться активной между прибытиями. Часто можно переключаться между этими паттернами, меняя триггер при сохранении неизменной логики запроса, но использовать контрольную точку повторно только тогда, когда запрос остаётся совместимым с контрольными точками.
Sources
Structured Streaming поддерживает несколько типов исходных кодов в рабочих нагрузках Fabric Spark. Выбирайте источник, исходя из источника текущих данных и того, как нужно протестировать запрос.
| Source | Behavior | Типичное использование |
|---|---|---|
rate |
Генерирует строки с заданной скоростью. | Разработка, демонстрации и испытания нагрузки. |
| Файлы | Читает новые файлы, которые появляются в каталоге. Поддерживаемые форматы включают text, CSV, JSON, Parquetи ORC. |
Инкрементальная загрузка файлов. |
| Таблица дельта | Читает первоначальный снимок таблицы, а затем добавляет коммиты как инкрементальный поток. Используйте ленту данных изменения для обработки обновлений и удалений на уровне строк. | Инкрементальная обработка в озерных домах. |
| Apache Kafka | Читает сообщения от Kafka Topics через коннектор Spark Kafka. | Приложения, управляемые событиями, и интеграция с потоковой трансляцией. |
| Центры событий Azure | Читает события через разъём Spark Kafka и совместимую с Kafka конечную точку Event Hubs. | Управляемый прием событий в Azure. |
| Fabric eventstream | Читает события через совместимую с Кафкой конечную точку, открытую потоком событий. | Прием, преобразование и маршрутизация событий Fabric. |
Для сценариев поглощения lakehouse распространёнными источниками являются Центры событий Azure и Fabric eventstream. Для разработки запросов исходный rate код предоставляет простой поток без внешних зависимостей.
Добавляйте файлы в папку потокового исходника атомарно, чтобы Spark не обнаружил частично написанный файл. Не записывайте выход запроса в путь, который также читает источник файла.
Sinks
Приемник получает результаты потокового запроса. В Fabric таблица Delta является основным хранилищем для устойчивых потоковых данных, поскольку она обеспечивает ACID-транзакции, открытое хранение, управление схемой и совместимость с другими Fabric-опытами.
| Sink | Behavior | Типичное использование |
|---|---|---|
| Таблица дельта | Записывает потоковые данные в управляемую или внешнюю таблицу Delta в озерном доме. | Надёжные аналитические данные и интеграция с другими опытами Fabric. |
| Настраиваемая конечная точка Eventstream | Отправляет события через совместимые с Event Hubs или Kafka-совместимые конечные точки. Сведения о настройке см. в разделе Добавление пользовательского источника конечной точки в поток событий. | Обработка, маршрутизация и доставка потоков событий. |
| Центры событий Azure or Apache Kafka | Записывает ключевые сообщения через раковину Spark Kafka. | Приложения, управляемые событиями, нижестоящие потребители потоков данных и низколатентная маршрутизация. |
foreachBatch |
Запускает пользовательскую логику для каждой микропартии. Сделайте логику идемпотентной, потому что Spark может повторно обработать микропакет. Примеры см. раздел «Распространённые паттерны». | Операции MERGE в Delta, несколько мест назначения и пакетные средства записи без встроенного потокового приёмника. |
console |
Печатает строки на выход из блокнота. | Только короткие тесты разработки. |
memory |
Сохраняет выводы в таблице внутри памяти. | Интерактивная отладка в сессии в блокноте. |
Выбирайте производственный поглотилитель в зависимости от того, как потребители на следующей линии используют выход. Используйте таблицу Delta для прочных аналитических данных, или Eventstream, Центры событий Azure или Apache Kafka для дальнейшей обработки событий. Приёмники консоли и памяти не обеспечивают постоянного хранилища и не подходят для рабочих нагрузок в промышленной среде.
Модели выполнения: краткий обзор
Постоянно включённые и запланированные шаблоны описывают срок службы запроса. Отдельно, Structured Streaming использует следующие модели выполнения для обработки записей:
| Модель выполнения | Description | Оптимальный сценарий |
|---|---|---|
| Микропакетный режим | Spark делит поступающие данные на небольшие партии, запускает план запросов для каждой партии и фиксирует результаты после завершения каждой партии. | Универсальная потоковая загрузка данных, агрегации, обработка файлов и надёжная запись в таблицы Delta. |
| Режим реального времени | Spark обрабатывает записи с ультранизкой задержкой, выполняя длительные потоковые задачи. Режим реального времени требует Fabric Runtime 2.0 или новее. | Обработка событий с низкой задержкой, где время отклика важнее максимальной пакетной пропускной способности. |
Начните со стандартной модели микропакетной обработки, если только ваша рабочая нагрузка не предъявляет жёстких требований к задержке. Для получения дополнительной информации см. режим реального времени.
Надёжность и отказоустойчивость
Structured Streaming обеспечивает отказоустойчивость с помощью контрольных точек и повторно воспроизводимых входных данных. Контрольная точка хранит прогресс запроса, метаданные, смещения и информацию о состоянии в прочном хранилище. Если потоковый запрос останавливается и заново запускается с той же точкой проверки, Spark возобновляется с последнего зафиксированного прогресса вместо того, чтобы начинать заново.
Точная обработка зависит от источника, приемника, контрольной точки и логики запроса. При наличии источника с возможностью повторного воспроизведения, надёжного чекпоинта и идемпотентного или транзакционного приёмника, такого как Delta Lake, Spark может избежать фиксации дублирующихся результатов после сбоев. Держите каждую контрольную точку потокового запроса в отдельном месте и не используйте контрольную точку между разными определениями запросов.
Гарантии доставки различаются в зависимости от раковины. Delta и file sinks поддерживают запись точно один раз, тогда как Kafka и foreach sinks обеспечивают доставку как минимум один раз. Сделайте пользовательские и foreach операции записи идемпотентными, чтобы повторные попытки не создавали дубликаты.
Операции с состоянием, такие как агрегации и соединения между потоками, хранят промежуточное состояние в контрольной точке сохранения. Тщательно продумывайте хранение чекпойнтов, поскольку чекпойнт является частью гарантий восстановления запроса.
Пример Delta от конца до конца
Следующий пример PySpark считывает данные из источника rate, добавляет вычисляемый столбец и записывает поток в таблицу Delta. Используйте это как минимальный шаблон для потоковой передачи данных от источника к Delta.
from pyspark.sql.functions import col
streaming_df = (
spark.readStream
.format("rate")
.option("rowsPerSecond", 1000)
.load()
.select(
col("timestamp"),
col("value"),
(col("value") % 10).alias("bucket")
)
)
query = (
streaming_df.writeStream
.format("delta")
.option("checkpointLocation", "Files/checkpoints/rate_to_delta")
.outputMode("append")
.toTable("rate_stream_delta")
)
query.awaitTermination()
В этом примере источник rate создаёт непрерывный входный поток, select() определяет преобразование и toTable("rate_stream_delta") записывает вывод в таблицу Delta. Местоположение контрольной точки позволяет Spark восстановить прогресс запроса, если задание перезапускается.
Запрос работает непрерывно, потому что не указывает триггер. Чтобы выполнить ту же логику, что и ограниченное инкрементальное задание, добавьте .trigger(availableNow=True) до toTable(). Spark обрабатывает доступные данные, фиксирует контрольную точку и состояние, и останавливается. Последующий забег с той же контрольной точкой продолжается из этого прогресса.
Important
start() и toTable() немедленно возвращают дескриптор StreamingQuery. Запрос работает асинхронно и не мешает триггерному блокноту или определению задания Spark достигать терминального состояния. Вызывайте query.awaitTermination() для каждого запроса, завершения которого должна дождаться запущенная задача. Для нескольких постоянно активных запросов используйте spark.streams.awaitAnyTermination(), чтобы определить, когда один из запросов прекращается.
Дальнейшие действия
Используйте связанные статьи, чтобы перейти от концепций к реализации. Начните с пошаговых руководств по приему данных в Lakehouse, если вам нужен готовый шаблон процесса от источника данных до Delta. Проверьте триггеры и режимы вывода, прежде чем настраивать момент вывода результатов. Изучите потоковую обработку с сохранением состояния, если ваш запрос использует агрегации, дедупликацию или соединения. Используйте статью о лучших практиках, прежде чем переносить работу по стримингу в продакшн.
Связанные материалы
- Режим реального времени
- Триггеры и режимы вывода
- Потоковая обработка с сохранением состояния
- Распространенные шаблоны
- Лучшие практики структурированного стриминга
- Потоковая передача данных в lakehouse с использованием Spark
- Получите потоковые данные в хранилище данных и доступ к ним через конечную точку аналитики SQL