CREATE TABLE ... FLOW (pipelines)
Bêta
Cette fonctionnalité est en Bêta.
Utilisez l'instruction CREATE TABLE ... FLOW pour créer une table gérée dans un pipeline, rédigée par un ou plusieurs flux.
Syntaxe
CREATE TABLE
table_name
[ table_specification ]
[ table_clauses ]
[ flow_clause ]
table_specification
( { column_identifier column_type [column_properties] } [, ...]
[ CONSTRAINT expectation_name EXPECT (expectation_expr)
[ ON VIOLATION { FAIL UPDATE | DROP ROW } ] ] [, ...] )
table_clauses
{ PARTITIONED BY (col [, ...]) |
CLUSTER BY clause |
LOCATION path |
COMMENT table_comment |
TBLPROPERTIES clause |
WITH { ROW FILTER clause } } [ ... ]
flow_clause
FLOW INSERT [ONCE] BY NAME query
Pour fusionner plusieurs sources en une seule table gérée, déclarez plusieurs flux qui la ciblent avec CREATE FLOW (pipelines):
CREATE FLOW flow_name AS INSERT INTO table_name BY NAME query
parameter
-
nom de table
Le nom de la table gérée à créer. Si le nom n'est pas qualifié, la table est créée dans le schéma cible du pipeline. Le nom ne doit pas déjà appartenir à une table de streaming.
-
spécification de table
Définit facultativement les colonnes, leurs types, propriétés et descriptions. S'il est omis, le schéma est déduit de la query du flux.
-
CONSTRAINT expectation_name EXPECT (expectation_expr) [ ON VIOLATION { FAIL UPDATE | DROP ROW } ]
Ajoute des attentes en matière de qualité des données à la table gérée. Ces attentes en matière de qualité des données peuvent être suivies au fil du temps et consultées via le journal des événements du pipeline. Une attente
FAIL UPDATEentraîne l'échec du traitement lors de la création de la table ainsi que lors de son actualisation. Une attenteDROP ROWentraîne la suppression de la ligne entière si l'attente n'est pas satisfaite. Voir Gérer la qualité des données avec les attentes de pipeline.expectation_exprpeut être composé de littéraux, d’identifiants de colonne au sein de la table et de fonctions ou opérateurs SQL intégrés et déterministes, à l’exception de :- Fonctions d’agrégation
- Fonctions de fenêtre analytiques
- Fonctions de fenêtre de classement
- Fonctions de génération à valeur tabulaire
De plus,
expectation_exprne doit contenir aucune sous-requête. - Fonctions d’agrégation
-
PARTITIONNÉ PAR (col [, ...])
Partitionne facultativement la table par un sous-ensemble de colonnes.
-
Clause CLUSTER BY
Active éventuellement le liquid clustering sur la table. Vous ne pouvez pas combiner
PARTITIONED BYetCLUSTER BY. -
LOCATION path
Un emplacement de stockage facultatif pour les données de la table.
-
COMMENT table_comment
Un littéral
STRINGdécrivant la table. -
Clause TBLPROPERTIES
Définit facultativement une ou plusieurs propriétés de table définies par l’utilisateur.
-
Clause WITH ROW FILTER
Ajoute une fonction de filtre de ligne à la table. Les futures query pour cette table reçoivent un sous-ensemble des lignes pour lesquelles la fonction évalue à
TRUE. -
FLOW INSERT [ONCE] BY NAME query
Définit un flux d'ajout qui insère le résultat de
querydans la table, faisant correspondre les colonnes de résultat aux colonnes de table *par nom*.querypeut référencer des sources par batch ou en streaming.ONCEexécute le flux une seule fois (par exemple, pour un remplissage) plutôt qu'à chaque mise à jour. Chaque flux nommé traite son entrée exactement une fois par mise à jour de pipeline, de manière identique àFLOW INSERT BY NAMEsur une table de streaming.
Limitations
- Les tables gérées ne prennent pas en charge les flux de modifications CDC.
AUTO CDC INTO(SQL) ouapply_changes/apply_changes_from_snapshot(Python) sur une table gérée échoue avecMANAGED_TABLE_DOES_NOT_SUPPORT_CDC. Utilisez une CREATE STREAMING TABLE (pipelines) pour les cibles CDC. - Les tables gérées ne prennent pas en charge
FLOW ... REPLACE WHERE. SeulFLOW INSERT BY NAMEest pris en charge. - Les tables gérées ne sont prises en charge que dans les pipelines avec Unity Catalog. Le Hive metastore n'est pas pris en charge.
- Vous ne pouvez pas réutiliser le nom d’une table de streaming existante pour une table gérée. Supprimez d’abord la table de streaming, ou l’instruction échouera avec
CANNOT_SWITCH_STREAMING_TABLE_TO_MANAGED_TABLE.
Exemples
-- Create a managed table populated by an inline append flow from a streaming table
CREATE TABLE output
FLOW INSERT BY NAME SELECT * FROM STREAM(samples.tpch.orders);
-- Create a managed table that ingests files with schema inference and evolution
CREATE TABLE raw_data
FLOW INSERT BY NAME
SELECT * FROM STREAM read_files('abfss://<container-name>@<storage-account-name>.dfs.core.windows.net/base/path');
-- Create a partitioned managed table from a streaming source
CREATE TABLE events
PARTITIONED BY (bucket)
FLOW INSERT BY NAME
SELECT id, bucket FROM STREAM read_files('abfss://my_path', format => 'json');
-- Create a managed table with liquid clustering
CREATE TABLE orders_clustered
CLUSTER BY (order_date, customer_id)
FLOW INSERT BY NAME
SELECT
o_orderkey AS order_id,
o_custkey AS customer_id,
o_orderdate AS order_date,
o_totalprice AS total_price
FROM STREAM(samples.tpch.orders);
-- Create a managed table with a data quality expectation that drops violating rows
CREATE TABLE valid_events
(CONSTRAINT positive_id EXPECT (id > 0) ON VIOLATION DROP ROW)
FLOW INSERT BY NAME
SELECT id FROM STREAM read_files('s3://bucket/path', format => 'json');