Использование каталога Unity с структурированной потоковой передачей

На этой странице показано, как использовать структурированную потоковую передачу с каталогом Unity для управления данными для добавочных и потоковых рабочих нагрузок на Azure Databricks.

Какие функции структурированной потоковой передачи поддерживают каталог Unity?

Unity Catalog не добавляет каких-либо явных ограничений для источников и приёмников данных Structured Streaming, доступных в Azure Databricks.

С помощью каталога Unity и структурированной потоковой передачи вы можете:

  • Передавайте данные в потоковом режиме как из управляемых, так и из внешних таблиц. См. управляемые таблицы Unity Catalog для Delta Lake и Apache Iceberg.
  • Используйте внешние расположения, управляемые каталогом Unity, для взаимодействия с данными с помощью URI хранилища объектов.
  • Записывайте данные во внешние таблицы, используя имена таблиц или пути к файлам. Для взаимодействия с управляемыми таблицами необходимо использовать имя таблицы.

Для структурированных контрольных точек потоковой передачи необходимо использовать пути во внешних расположениях, управляемых каталогом Unity. Дополнительные сведения о безопасном подключении хранилища к каталогу Unity см. в статье "Подключение к облачному хранилищу объектов" с помощью каталога Unity.

чтение представления каталога Unity в виде потока

В Databricks Runtime 14.3 LTS и более поздних версиях можно использовать структурированную потоковую передачу для чтения из представлений, зарегистрированных в каталоге Unity. Базовые таблицы должны использовать формат Delta Lake. Другие ограничения см. в разделе "Ограничения".

Чтобы прочитать представление со структурированной потоковой передачей, используйте .table() метод с идентификатором представления:

df = (spark.readStream
  .table("demoView")
)

Пользователи должны иметь SELECT привилегии целевого представления.

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

Поддерживаемые параметры потоковой передачи

Средство чтения потоковой передачи применяет параметры к файлам и метаданным базовых таблиц Delta Lake для указанного представления.

Поддерживаются следующие параметры:

  • maxFilesPerTrigger
  • maxBytesPerTrigger
  • ignoreDeletes
  • skipChangeCommits
  • withEventTimeOrder
  • startingTimestamp
  • startingVersion

Операции чтения для представлений с UNION ALL не поддерживают параметры withEventTimeOrder и startingVersion.

Если вы предоставляете неподдерживаемые параметры, например readChangeFeed, Spark вызывает это исключение:

AnalysisException: [UNSUPPORTED_STREAMING_OPTIONS_FOR_VIEW.UNSUPPORTED_OPTION] Unsupported for streaming a view. Reason: option <option> is not supported.

Поддерживаемые операции потоковой передачи

К поддерживаемым операциям относятся:

Операция Описание Operator Example
Проект Управление разрешениями на уровне столбцов SELECT... FROM... CREATE VIEW project_view AS SELECT id, value FROM source_table
Фильтр Управление разрешениями на уровне строк WHERE... CREATE VIEW filter_view AS SELECT * FROM source_table WHERE value > 100
Объединение всех Результаты из нескольких таблиц UNION ALL CREATE VIEW union_view AS SELECT id, value FROM source_table1 UNION ALL SELECT * FROM source_table2

Неподдерживаемые операции включают агрегирование, сортировку и табличные функции, такие как table_changes(). Дополнительные сведения о функциях с табличным значением см. в разделе вызов функции с табличным значением (TVF).

При потоковой передаче из представления с неподдерживаемой операцией Spark вызывает это исключение:

UnsupportedOperationException: [UNEXPECTED_OPERATOR_IN_STREAMING_VIEW] Unexpected operator <operator> in the CREATE VIEW statement as a streaming source. A streaming view query must consist only of SELECT, WHERE, and UNION ALL operations.

Ограничения