Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Потоковая таблица — это таблица с поддержкой потоковой или добавочной обработки данных. Потоковые таблицы поддерживаются конвейерами. При каждом обновлении потоковой таблицы новоиспечённые данные в исходных таблицах присоединяются к потоковой таблице. Таблицы потоковой передачи можно обновлять вручную или по расписанию.
Дополнительные сведения о выполнении или планировании обновлений см. в разделе "Запуск обновления конвейера".
Синтаксис
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
Имя только что созданной таблицы. Полное имя таблицы должно быть уникальным.
спецификация таблицы
Это необязательное предложение определяет список столбцов, их типов, свойств, описаний и ограничений столбцов.
-
Имена столбцов должны быть уникальными и сопоставляться с выходными столбцами запроса.
-
Указывает тип данных столбца. Не все типы данных, поддерживаемые Azure Databricks, поддерживаются стриминговыми таблицами.
column_comment
Необязательное литеральное значение
STRING, описывающее столбец. Этот параметр должен быть указан вместе сcolumn_type. Если тип столбца не указан, комментарий столбца пропускается.ГЕНЕРИРУЕТСЯ ВСЕГДА КАК ( expr )
При указании этого условия значение этого столбца определяется задаваемым
expr.DEFAULT COLLATIONтаблицы должен бытьUTF8_BINARY.exprмогут состоять из литералов, идентификаторов столбцов в таблице и детерминированных встроенных функций ИЛИ операторов SQL, кроме следующих:- Агрегатные функции
- Функции окна аналитики
- Ранжирование функций окна
- Функции, генерирующие табличные значения
- Столбцы с кодировкой, отличной от
UTF8_BINARY
Кроме того,
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.-
Добавляет ограничение на информационный первичный ключ или информационный внешний ключ в столбец в таблице потоковой передачи.
-
Добавляет функцию маски столбца для анонимизации конфиденциальных данных.
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является обязательным. См. раздел Частичная замена снимков с ЗАМЕНОЙ ИСПОЛЬЗОВАНИЕМ потоков.
-
Это предложение заполняет таблицу с помощью данных из
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 для чтения и записи.
maxFilesPerTriggermaxBytesPerTriggerstartingVersionstartingTimestampreadChangeFeedwithEventTimeOrderskipChangeCommits
Необходимые разрешения
Пользователь с правами выполнения для конвейера должен иметь следующие разрешения:
-
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 INTOANALYZE TABLERESTORETRUNCATEGENERATE 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);