CREATE STREAMING TABLE (конвейерные линии)

Потоковая таблица — это таблица с поддержкой потоковой или добавочной обработки данных. Потоковые таблицы поддерживаются конвейерами. При каждом обновлении потоковой таблицы новоиспечённые данные в исходных таблицах присоединяются к потоковой таблице. Таблицы потоковой передачи можно обновлять вручную или по расписанию.

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

Синтаксис

CREATE [OR REFRESH] [PRIVATE] STREAMING TABLE
  table_name
  [ table_specification ]
  [ table_clauses ]
  [ {flow_clause | AS query} ]

table_specification
  ( { column_identifier column_type [column_properties] } [, ...]
    [ column_constraint ] [, ...]
    [ , table_constraint ] [...] )

   column_properties
      { NOT NULL | GENERATED ALWAYS AS ( expr ) | GENERATED { ALWAYS | BY DEFAULT } AS IDENTITY [ ( [ START WITH start | INCREMENT BY step ] [ ...] ) ] | DEFAULT default_expression | COMMENT column_comment | column_constraint | MASK clause } [ ... ]

table_clauses
  { USING DELTA
    PARTITIONED BY (col [, ...]) |
    CLUSTER BY clause |
    LOCATION path |
    COMMENT view_comment |
    TBLPROPERTIES clause |
    WITH { ROW FILTER clause } } [ ... ]
   } [ ... ]

flow_clause
  FLOW { { INSERT [ONCE] BY NAME query } |
  { AUTO CDC auto_cdc_flow_spec } |
  { REPLACE WHERE predicate BY NAME query } |
  { REPLACE USING ( column_name [, ...] ) SEQUENCE BY sequence_column BY NAME query } }

Параметры

  • REFRESH

    Если задано, создает таблицу или обновляет существующую таблицу и его содержимое.

  • ЧАСТНЫЙ

    Создает частную потоковую таблицу.

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

    Ранее частные потоковые таблицы создавались с параметром TEMPORARY.

  • table_name

    Имя только что созданной таблицы. Полное имя таблицы должно быть уникальным.

  • спецификация таблицы

    Это необязательное предложение определяет список столбцов, их типов, свойств, описаний и ограничений столбцов.

    • column_identifier

      Имена столбцов должны быть уникальными и сопоставляться с выходными столбцами запроса.

    • тип столбца

      Указывает тип данных столбца. Не все типы данных, поддерживаемые Azure Databricks, поддерживаются стриминговыми таблицами.

    • column_comment

      Необязательное литеральное значение STRING, описывающее столбец. Этот параметр должен быть указан вместе с column_type. Если тип столбца не указан, комментарий столбца пропускается.

    • ГЕНЕРИРУЕТСЯ ВСЕГДА КАК ( expr )

      При указании этого условия значение этого столбца определяется задаваемым expr.

      DEFAULT COLLATION таблицы должен быть UTF8_BINARY.

      expr могут состоять из литералов, идентификаторов столбцов в таблице и детерминированных встроенных функций ИЛИ операторов SQL, кроме следующих:

      Кроме того, expr не должен содержать какой-либо вложенный запрос.

    • ГЕНЕРИРУЕТСЯ { ВСЕГДА | ПО УМОЛЧАНИЮ } КАК ИДЕНТИФИКАТОР [ ( [ НАЧИНАЯ С start ] [ УВЕЛИЧИВАЯ НА step ] ) ]

      Применяется к:отмечено да Databricks SQL отмечено да Databricks Runtime 10.4 LTS и выше

      Определяет столбец идентификаторов. При записи в таблицу, если не предоставлено значение для столбца с идентификатором, ему автоматически назначается уникальное и статистически увеличивающееся (или уменьшающееся, если step отрицательное) значение. Это предложение поддерживается только для таблиц Delta. Это предложение можно использовать только для столбцов с типом данных BIGINT.

      Автоматически назначаемые значения начинаются с start и увеличиваются на step. Назначенные значения являются уникальными, но не гарантируют, что они являются смежными. Оба параметра являются необязательными, а значение по умолчанию — 1. step не может иметь значение 0.

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

      При использовании ALWAYS нельзя указать собственные значения для столбца идентификаторов.

      Следующие операции не поддерживаются:

      • PARTITIONED BY идентификационный столбец
      • UPDATE идентификационный столбец

      Замечание

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

    • DEFAULT выражение_по_умолчанию

      Область применения:отмечено как да Databricks SQL отмечено как да Databricks Runtime 11.3 LTS и выше

      Определяет значение DEFAULT для столбца, которое применяется в INSERT, UPDATEи MERGE ... INSERT, если столбец не указан.

      Если значение по умолчанию не указано, DEFAULT NULL применяется для столбцов, допускающих значение NULL.

      default_expression может состоять из литералов и встроенных функций SQL или операторов за исключением следующих:

      Кроме того, default_expression не должен содержать какой-либо вложенный запрос.

      DEFAULT поддерживается для источников CSV, JSON, PARQUET и ORC.

    • column_constraint

      Добавляет ограничение на информационный первичный ключ или информационный внешний ключ в столбец в таблице потоковой передачи.

    • Клаузула MASK

      Добавляет функцию маски столбца для анонимизации конфиденциальных данных.

      См. фильтры строк и маски столбцов.

    • CONSTRAINT EXPECTATION_NAME ОЖИДАТЬ (expectation_expr) [ ON VIOLATION { FAIL UPDATE | DROP ROW } ]

      Добавляет ожидания качества данных в таблицу потоковой передачи. Эти ожидания касаемо качества данных можно отслеживать на протяжении времени и получать к ним доступ через журнал событий потоковой таблицы. Ожидание FAIL UPDATE приводит к сбою обработки при создании таблицы, а также обновлении таблицы. Ожидание DROP ROW приводит к тому, что вся строка будет удалена, если ожидание не выполнено. См. Управление качеством данных, используя ожидания конвейера.

      expectation_expr могут состоять из литералов, идентификаторов столбцов в таблице и детерминированных встроенных функций ИЛИ операторов SQL, кроме следующих:

      Кроме того, expr не должен содержать какой-либо вложенный запрос.

  • ограничение_таблицы

    При указании схемы можно определить первичные и внешние ключи. Ограничения являются информационными и не применяются. См. пункт CONSTRAINT в справочной информации по языку SQL.

    Замечание

    Чтобы определить ограничения таблицы, конвейер должен поддерживать Unity Catalog.

  • таблица_условий

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

    • ИСПОЛЬЗОВАНИЕ DELTA

      Задает формат данных. Единственным вариантом является DELTA.

      Это предложение является необязательным и по умолчанию используется delta.

    • РАЗДЕЛЁННО ПО

      Необязательный список одного или нескольких столбцов, используемых для секционирования в таблице. Взаимоисключающ с CLUSTER BY.

      Liquid Clustering предоставляет гибкое и оптимизированное решение для кластеризации. Рекомендуется использовать CLUSTER BY вместо PARTITIONED BY для конвейеров.

    • CLUSTER BY

      Включите кластеризацию жидкости в таблице и определите столбцы, используемые в качестве ключей кластеризации. Используйте автоматическую кластеризацию с CLUSTER BY AUTO, и Databricks интеллектуально выбирает ключи кластеризации для оптимизации производительности запросов. Взаимоисключающ с PARTITIONED BY.

      См. раздел "Использование кластеризации жидкости" для таблиц.

    • МЕСТОПОЛОЖЕНИЕ

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

    • КОММЕНТАРИЙ

      Необязательный STRING литерал для описания таблицы.

    • TBLPROPERTIES

      Необязательный список свойств таблицы.

    • С ROW FILTER

    Добавляет функцию фильтра строк в таблицу. Будущие запросы для этой таблицы получают подмножество строк, для которых функция оценивается как TRUE. Это полезно для точного управления доступом, так как это позволяет функции проверять личность и членства в группах вызывающего пользователя, чтобы решить, следует ли фильтровать определенные строки.

    См. ROW FILTER пункт.

    • ПОТОКА

      При необходимости определяет поток , встроенный с созданием таблицы. Поток — это запрос с отслеживанием состояния, который обновляет содержимое таблицы. Если FLOW это не указано, можно использовать AS query вместо этого или определять потоки отдельно с CREATE FLOWпомощью. Можно указать один из следующих типов потоков:

      • INSERT ПО ИМЕНИ

        Вставляет данные в таблицу по имени столбца. ONCE Если параметр не указан, запрос должен быть потоковым запросом. Используйте ключевое слово STREAM для применения семантики потоковой передачи при чтении из источника. Если чтение сталкивается с изменением или удалением существующей записи, возникает ошибка. Самое безопасное — читать из статических или источников только для добавления.

        Замечание

        FLOW INSERT BY NAME эквивалентен использованию AS query. Следующие два оператора имеют одинаковое поведение:

        CREATE OR REFRESH STREAMING TABLE raw_data
        AS SELECT * FROM STREAM read_files('abfss://my_path');
        
        CREATE OR REFRESH STREAMING TABLE raw_data
        FLOW INSERT BY NAME SELECT * FROM STREAM read_files('abfss://my_path');
        
      • ОДИН РАЗ

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

      • AUTO CDC

        Это важно

        Доступно в Databricks Runtime 17.3 и выше и PREVIEW канале Pipelines.

        Определяет AUTO CDC поток, который обрабатывает записи отслеживания измененных данных (CDC) из источника в таблицу. Используется AUTO CDC , если исходные данные включают семантику CDC. См . API AUTO CDC: упрощение отслеживания изменений с помощью конвейеров.

      • Запрос REPLACE WHEREpredicate BY NAME

        REPLACE WHERE Определяет поток, который перекомпьютерирует и перезаписывает только соответствующие predicateстроки, оставляя все остальные строки неуправляемыми. Используется REPLACE WHERE для добавочной пакетной обработки соединений и агрегатов, поздних поступающих данных, эволюции схемы и обратной заполнения. BY NAME является обязательным. См. раздел "Пакетная обработка с помощью потоков REPLACEWHERE".

      • ЗАМЕНА С ( column_name [, ...] ) ПОСЛЕДОВАТЕЛЬНОСТЬ sequence_column ПО ИМЕНИ ЗАПРОСА

        Это важно

        Эта функция доступна в бета-версии. Требуется Databricks Runtime 18.2 и выше.

        Определяет REPLACE USING поток, который заменяет все строки, совпадающие с указанными ключевыми столбцами, и оставляет все остальные строки нетронутыми. Используйте REPLACE USING тогда, когда исходный код — серия частичных снимков, ключевых по столбцам. SEQUENCE BY Упорядочивает обновления так, чтобы побеждает самая высокая последовательность для ключа, даже если обновления приходят не по порядку. Источник должен быть источником потоковой передачи. BY NAME является обязательным. См. раздел Частичная замена снимков с ЗАМЕНОЙ ИСПОЛЬЗОВАНИЕМ потоков.

  • ЗАПРОС AS

    Это предложение заполняет таблицу с помощью данных из query. Этот запрос должен быть потоковым запросом. Используйте ключевое слово STREAM для использования семантики потоковой передачи для чтения из источника. Если чтение сталкивается с изменением или удалением существующей записи, возникает ошибка. Самое безопасное — читать из статических или источников только для добавления. Чтобы получать данные с фиксациями изменений, можно добавить skipChangeCommits параметр чтения для обработки ошибок.

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

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

    • Параметры чтения

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

      SELECT * FROM STREAM source_table WITH (SKIPCHANGECOMMITS=TRUE, STARTINGVERSION=X)
      

      Необязательный =TRUE , поэтому можно также указать логический параметр, как показано ниже:

      SELECT * FROM STREAM source_table WITH (SKIPCHANGECOMMITS)
      

      Замечание

      Параметры чтения поддерживаются только для Databricks Runtime 17.3 и более поздних версий.

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

      • maxFilesPerTrigger
      • maxBytesPerTrigger
      • startingVersion
      • startingTimestamp
      • readChangeFeed
      • withEventTimeOrder
      • skipChangeCommits

Необходимые разрешения

Пользователь с правами выполнения для конвейера должен иметь следующие разрешения:

  • SELECT привилегии над базовыми таблицами, на которые ссылается потоковая таблица.
  • USE CATALOG привилегия в отношении родительского каталога и привилегия USE SCHEMA в отношении родительской схемы.
  • CREATE MATERIALIZED VIEW привилегии в схеме для потоковой таблицы.

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

  • USE CATALOG привилегия в отношении родительского каталога и привилегия USE SCHEMA в отношении родительской схемы.
  • Владение таблицей потоковой передачи или REFRESH привилегиями в потоковой таблице.
  • Владелец потоковой таблицы должен иметь SELECT привилегии над базовыми таблицами, на которые ссылается потоковая таблица.

Для того чтобы пользователь мог запрашивать полученную потоковую таблицу, им требуется:

  • USE CATALOG привилегия в отношении родительского каталога и привилегия USE SCHEMA в отношении родительской схемы.
  • SELECT привилегия на потоковую таблицу.

Ограничения

  • Только владельцы таблиц могут обновлять потоковые таблицы, чтобы получить последние данные.
  • Использование команд ALTER TABLE запрещено в потоковых таблицах. Определение и свойства таблицы должны быть изменены с помощью инструкции CREATE OR REFRESH или ALTER STREAMING TABLE.
  • Эволюционирование схемы таблицы с помощью команд DML, таких как INSERT INTO, и MERGE не поддерживается.
  • Следующие команды не поддерживаются в таблицах потоковой передачи:
    • CREATE TABLE ... CLONE <streaming_table>
    • COPY INTO
    • ANALYZE TABLE
    • RESTORE
    • TRUNCATE
    • GENERATE MANIFEST
    • [CREATE OR] REPLACE TABLE
  • Переименование таблицы или изменение владельца не поддерживается.

Примеры

-- Define a streaming table from a volume of files:
CREATE OR REFRESH STREAMING TABLE customers_bronze
AS SELECT * FROM STREAM read_files("/databricks-datasets/retail-org/customers/*", format => "csv")

-- Define a streaming table from a streaming source table:
CREATE OR REFRESH STREAMING TABLE customers_silver
AS SELECT * FROM STREAM(customers_bronze)

-- Use automatic liquid clustering to let Databricks choose the clustering columns:
CREATE OR REFRESH STREAMING TABLE customers_bronze_auto
CLUSTER BY AUTO
AS SELECT * FROM STREAM read_files("/databricks-datasets/retail-org/customers/*", format => "csv")

-- Define a table with a row filter and column mask:
CREATE OR REFRESH STREAMING TABLE customers_silver (
  id int COMMENT 'This is the customer ID',
  name string,
  region string,
  ssn string MASK catalog.schema.ssn_mask_fn COMMENT 'SSN masked for privacy'
)
WITH ROW FILTER catalog.schema.us_filter_fn ON (region)
AS SELECT * FROM STREAM(customers_bronze)

-- Define a streaming table with an identity column:
CREATE OR REFRESH STREAMING TABLE customers_with_id (
  customer_id BIGINT GENERATED ALWAYS AS IDENTITY,
  name string,
  region string
)
AS SELECT name, region FROM STREAM(customers_bronze)

-- Define a streaming table that you can add flows into:
CREATE OR REFRESH STREAMING TABLE orders;

-- Define a streaming table with an inline append flow:
CREATE OR REFRESH STREAMING TABLE raw_data
FLOW INSERT BY NAME SELECT * FROM STREAM read_files('abfss://my_path');

-- Define a streaming table with an inline AUTO CDC flow:
CREATE OR REFRESH STREAMING TABLE target
FLOW AUTO CDC
FROM stream(cdc_data.users)
KEYS (userId)
SEQUENCE BY sequenceNum
STORED AS SCD TYPE 1;

-- Define a streaming table with an inline REPLACE USING flow that keeps the latest
-- row for each payment_id:
CREATE OR REFRESH STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);