Что такое конвейеры Lakeflow?

Конвейеры Lakeflow предоставляют декларативную платформу для создания конвейеров пакетных и потоковых данных в SQL и Python. Их основные понятия — это конвейеры, потоки, потоковые таблицы, материализованные представления и приёмники, которые совместно используются для обработки данных с автоматической оркестрацией и инкрементальными обновлениями.

Конвейеры Lakeflow расширяют декларативные конвейеры Apache Spark™ (SDP). Дополнительные сведения о SDP и сравнении с конвейерами Lakeflow см. в статье Apache Spark Декларативные конвейеры.

Tip

Новичок в конвейерах? Начните с How to use Lakeflow pipelines, чтобы понять, как использовать конвейеры Lakeflow на всех этапах их жизненного цикла, а также почему, со ссылками на задачи для каждого этапа.

Note

Для конвейеров Lakeflow требуется тарифный план Premium. Чтобы получить дополнительные сведения, обратитесь к группе учетной записи Databricks.

Каковы преимущества конвейеров?

В отличие от разработки процессов обработки данных с использованием API Apache Spark и Spark Structured Streaming в Databricks Runtime при ручной оркестрации через Lakeflow Jobs, декларативный характер конвейеров обеспечивает следующие преимущества:

  • Автоматическая оркестрация: конвейеры выполняют этапы обработки (так называемые "потоки выполнения") в нужном порядке с максимально возможным параллелизмом и повторяют попытки при временных сбоях с постепенным повышением уровня — от задачи Spark к потоку, а затем ко всему конвейеру.
  • Декларативная обработка: декларативные функции сокращают сотни строк ручного кода Spark и структурированной потоковой передачи до нескольких. AUTO CDC API обрабатывает события захвата изменений данных (CDC), включая SCD типа 1 и типа 2, без необходимости вручную писать код для событий, поступающих не по порядку, или использовать концепции потоковой обработки, такие как водяные знаки.
  • Инкрементальная обработка: механизм инкрементальной обработки поддерживает материализованные представления в актуальном состоянии: вы задаёте логику преобразования с пакетной семантикой, а механизм по возможности повторно обрабатывает только новые или изменённые исходные данные.

Основные понятия

На схеме ниже показаны наиболее важные понятия конвейеров.

Схема, показывающая, как основные понятия конвейеров связаны друг с другом на очень высоком уровне

Наборы данных

Конвейер создает три типа наборов данных, каждый из которых имеет разные семантики обработки:

Тип набора данных Обработка записей
Потоковая таблица Каждая запись обрабатывается ровно один раз, если источник допускает только добавление данных. Потоковые таблицы подходят для приема и добавочной обработки постоянно растущих данных.
Материализованное представление Результаты перекомпилируются по мере необходимости, чтобы отразить текущее состояние данных. Материализованные представления подходят для преобразований, агрегатов или предварительного вычисления результатов, потребляемых несколькими нижестоящими наборами данных.
Вид Вычисляется по запросу, а не сохраняется. Используйте представления для промежуточных преобразований и проверки, которые не нужно публиковать в каталоге.

Потоковая таблица — это форма управляемой таблицы каталога Unity, которая также является целевым объектом потоковой передачи. В таблице потоковой передачи может быть записано одно или несколько потоковых потоков (добавление, AUTO CDC). Потоки потоковой передачи можно определять явно и отдельно от целевой таблицы потоковой передачи или неявно в рамках определения таблицы потоковой передачи.

Материализованное представление также является формой таблицы, управляемой каталогом Unity, и служит целевым объектом для пакетной обработки. Материализованное представление может содержать один или несколько материализованных потоков представления, записанных в него. Материализованные представления отличаются от потоковых таблиц тем, что потоки всегда определяются неявно в рамках определения материализованного представления.

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

Когда следует использовать представления, материализованные представления и потоковые таблицы

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

Рассмотрите возможность использования представления:

  • Разорвать большой или сложный запрос на более удобные запросы.
  • Проверьте промежуточные результаты с помощью ожиданий.
  • Уменьшите затраты на хранение и вычислительные ресурсы для результатов, которые не нужно сохранять. Так как таблицы материализованы, они требуют дополнительных вычислительных ресурсов и ресурсов хранилища.

Рекомендуется использовать материализованное представление, когда:

  • Несколько последующих запросов обрабатывают таблицу. Так как материализованное представление кэширует свои результаты, подчиненные запросы считывают предварительно вычисляемые результаты вместо повторного вычисления запроса для каждого доступа.
  • Другие каналы, задания или запросы потребляют таблицу. Поскольку материализованное представление сохраняется как таблица в Unity Catalog, пользователи вне конвейера, в котором оно определяется, могут выполнять запросы к нему. Представления не материализованы, поэтому их можно использовать только в рамках того же конвейера.
  • Вы хотите проверить результаты запроса во время разработки. Поскольку материализованное представление существует как отдельный объект и к нему можно выполнять запросы вне конвейера, вы можете проверять правильность вычислений в процессе разработки. После проверки преобразуйте запросы, которые не требуют материализации в представления.
  • Запрос выполняет агрегацию или соединения таблиц, либо исходные данные могут изменяться из-за обновлений и удалений, а не только увеличиваться. Материализованное представление сохраняет результаты в соответствии с текущим состоянием исходных данных, в то время как потоковая таблица предназначена для источников только для добавления и обрабатывает каждую запись за один раз.

Подумайте об использовании потоковой таблицы, когда:

  • Запрос определяется для источника данных, который постоянно или постепенно растет.
  • Результаты запроса должны вычисляться постепенно.
  • Конвейеру требуется высокая пропускная способность и низкая задержка.

Note

Таблицы потоковых данных всегда определяются по отношению к потоковым источникам. Вы также можете использовать источники потоковой передачи с AUTO CDC ... INTO для применения обновлений из потоков данных CDC. См . API AUTO CDC: упрощение отслеживания изменений с помощью конвейеров.

Flows

Поток — это базовая концепция обработки данных в конвейерах и поддерживает потоковую и пакетную семантику. Поток считывает данные из источника, применяет определяемую пользователем логику обработки и записывает результат в целевой объект. Конвейеры используют тот же тип потока потоковой передачи (добавление, обновление, завершение), что и структурированная потоковая передача Spark. (В настоящее время предоставляются только потоки добавления и обновления .) Дополнительные сведения см. в режимах вывода в структурированной потоковой передаче.

Конвейеры также предоставляют дополнительные типы потоков:

  • AUTO CDC — это уникальный поток обработки данных в конвейерах данных Lakeflow, который обрабатывает события CDC, поступающие не по порядку, и поддерживает как SCD Type 1, так и SCD Type 2. Автоматический CDC недоступен в SDP.
  • Материализованное представление — это пакетный поток в конвейерах, который обрабатывает только новые данные и изменения в исходных таблицах по возможности.

Дополнительные сведения см. в разделе Инкрементальная загрузка и обработка данных с помощью потоков конвейера Lakeflow.

Sinks

Синк — это потоковая целевая точка для конвейера, поддерживающая таблицы Delta, темы Apache Kafka, темы Azure EventHubs и пользовательские источники данных на Python. Приемник может записывать в него один или несколько потоковых потоков (добавление, обновление).

Подробнее см. в разделе «Приемники в конвейерах Lakeflow».

Трубопроводы

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

Дополнительные сведения см. в разделе "Что такое конвейеры?".

Вы также можете определить автономные материализованные представления и потоковые таблицы за пределами конвейера Lakeflow, где Azure Databricks управляет конвейером. Чтобы сравнить два подхода, см. автономные конвейеры и конвейеры Lakeflow.

Пайплайн работает либо в триггерном, либо в непрерывном режиме; этот режим определяет, обновляет ли он доступные данные и прекращает обновление таблиц или поддерживает их в актуальном состоянии по мере поступления новых данных. Для сравнения двух режимов см. Триггерный и непрерывный режим конвейера.

Прием данных

Конвейеры поддерживают все источники данных, доступные в Azure Databricks. Databricks рекомендует использовать потоковые таблицы для большинства сценариев загрузки данных. Для файлов, хранящихся в облачном объектном хранилище, Auto Loader обеспечивает инкрементальную идемпотентную загрузку. Конвейеры для потоковых данных могут принимать данные непосредственно из шин сообщений, таких как Apache Kafka, Центры событий Azure, Amazon Kinesis и Google Pub/Sub. См. сведения о загрузке данных в конвейерах.

Качество данных

Ожидания — это необязательные условия в наборах данных, которые проверяют данные по мере их прохождения через конвейер обработки. Вы определяете ожидание как логическое ограничение SQL и указываете, что происходит при сбое записи: предупреждение, удаление записи или сбой обновления. См. Управление качеством данных, используя ожидания конвейера.

Дельта-интеграция

Все таблицы, создаваемые конвейерами и управляемые ими, являются таблицами Delta. Они имеют те же гарантии, что и Delta Lake, включая транзакции ACID, поездки по времени и принудительное применение схемы. Конвейеры добавляют дополнительные свойства таблиц и выполняют автоматическое обслуживание с помощью предиктивной оптимизации, включая операции OPTIMIZE и VACUUM. См. статью "Что такое Delta Lake в Azure Databricks?".

Дополнительные ресурсы