CREATE STREAMING TABLE (pipelines)
Une table de streaming est une table prenant en charge le traitement de données en streaming ou incrémentiel. Les tables de streaming sont prises en charge par des pipelines. Chaque fois qu'une table de streaming est actualisée, les données ajoutées aux tables sources sont ajoutées à la table de streaming. Vous pouvez refresh les tables de streaming manuellement ou selon un calendrier.
Pour en savoir plus sur la façon d’effectuer ou de planifier des actualisations, consultez Exécuter une mise à jour de pipeline.
Syntaxe
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 } }
parameter
-
REFRESH
S'il est spécifié, crée la table, ou met à jour une table existante et son contenu.
-
PRIVÉ
Crée une table de streaming privée.
- Ils ne sont pas ajoutés au catalogue et ne sont accessibles que dans le pipeline de définition.
- Ils peuvent avoir le même nom qu'un objet existant dans le catalogue. Dans le pipeline, si une table de streaming privée et un objet du catalogue ont le même nom, les références au nom se résolvent en la table de streaming privée.
- Les tables de streaming privées sont conservées uniquement pendant la durée de vie du pipeline, et non pas pour une seule mise à jour.
Les tables de streaming privées ont été précédemment créées avec le parameter
TEMPORARY. -
nom_de_table
Le nom de la table nouvellement créée. Le nom de table entièrement qualifié doit être unique.
-
table_specification
Cette clause facultative définit la liste des colonnes, leurs types, propriétés, descriptions et contraintes de colonne.
-
Les noms de colonnes doivent être uniques et correspondre aux colonnes de sortie de la query.
-
Spécifie le type de données de la colonne. Tous les types de données pris en charge par Databricks ne sont pas pris en charge par les tables de streaming.
-
commentaire_colonne
Un littéral
STRINGfacultatif décrivant la colonne. Cette option doit être spécifiée aveccolumn_type. Si le type de colonne n'est pas spécifié, le commentaire de colonne est ignoré. -
GÉNÉRÉ TOUJOURS COMME ( expr )
Lorsque vous spécifiez cette clause, la valeur de cette colonne est déterminée par le
exprspécifié.Le/La
DEFAULT COLLATIONde la table doit êtreUTF8_BINARY.exprpeut être composé de littéraux, d'identifiants de colonnes dans la table, et de fonctions ou opérateurs SQL intégrés 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 de table
- Colonnes avec un classement autre que
UTF8_BINARY
De plus,
exprne doit pas contenir de sous-requête. -
GENERATED { ALWAYS | BY DEFAULT } AS IDENTITY [ ( [ START WITH start ] [ INCREMENT BY step ] ) ]
S’applique à :
Databricks SQL
Databricks Runtime 10.4 LTS et versions ultérieures
Définit une colonne d'identité. Lorsque vous écrivez dans la table, et que vous ne fournissez pas de valeurs pour la colonne d'identité, une valeur unique et statistiquement croissante (ou décroissante si
stepest négatif) lui sera automatiquement attribuée. Cette clause est uniquement prise en charge pour les tables Delta. Cette clause peut être utilisée uniquement pour les colonnes de type de données BIGINT.Les valeurs attribuées automatiquement start par
startet augmentent destep. Les valeurs attribuées sont uniques mais ne sont pas garanties d’être contiguës. Les deux paramètres sont facultatifs, et la valeur default est 1.stepne peut pas être0.Si les valeurs attribuées automatiquement dépassent la plage du type de colonne d'identité, la query échouera.
Lorsque
ALWAYSest utilisé, vous ne pouvez pas fournir vos propres valeurs pour la colonne d'identité.Les Opérations suivantes ne sont pas prises en charge :
PARTITIONED BYune colonne d'identitéUPDATEune colonne d'identité
-
La déclaration d'une colonne d'identité sur une table désactive les transactions concurrentes. Utilisez uniquement les colonnes d'identité dans les cas d'utilisation où les écritures concurrentes dans la table cible ne sont pas requises.
-
DEFAULT default_expression
S'applique à :
Databricks SQL
Databricks Runtime 11.3 LTS et versions ultérieures
Définit une valeur
DEFAULTpour la colonne utilisée surINSERT,UPDATEetMERGE ... INSERTlorsque la colonne n'est pas spécifiée.Si aucun default n'est spécifié,
DEFAULT NULLest appliqué aux colonnes acceptant les valeurs nulles.default_expressionpeuvent être composés de littéraux, et de fonctions ou d'opérateurs SQL intégrés sauf :- Fonctions d'agrégation
- Fonctions de fenêtre analytiques
- Fonctions de fenêtre de classement
- Fonctions de génération à valeur de table
De plus,
default_expressionne doit pas contenir de sous-requête.DEFAULTest pris en charge pour les sourcesCSV,JSON,PARQUETetORC. -
Ajoute une contrainte de clé primaire informative ou de clé étrangère informative à la colonne dans une table de streaming.
-
Ajoute une fonction de masque de colonne pour anonymiser les données sensibles.
Consultez Filtres de lignes et masques de colonne.
-
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 de streaming. Ces attentes en matière de qualité des données peuvent être suivies dans le temps et accessibles via les Logs d'événements de la table de streaming. Une attente
FAIL UPDATEprovoque l’échec du traitement lors de la création de la table ainsi que de l’actualisation de la table. Une attenteDROP ROWprovoque 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 du pipeline.expectation_exprpeut être composé de littéraux, d'identifiants de colonnes dans la table, et de fonctions ou opérateurs SQL intégrés 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 de table
De plus,
exprne doit pas contenir de sous-requête. - Fonctions d'agrégation
-
contrainte de table
Lorsque vous spécifiez un schéma, vous pouvez définir des clés primaires et étrangères. Les contraintes sont informatives et ne sont pas appliquées. Consultez la clause CONSTRAINT dans la référence du langage SQL.
Pour définir des contraintes de table, votre pipeline doit être un pipeline compatible avec Unity Catalog.
-
clauses_de_table
Spécifiez éventuellement le partitionnement, les commentaires et les propriétés définies par l'utilisateur pour la table. Chaque sous-clause ne peut être spécifiée qu'une seule fois.
-
UTILISATION DE DELTA
Spécifie le format des données. La seule option est DELTA.
Cette clause est facultative et utilise default.
-
PARTITIONNÉ PAR
Une liste facultative d'une ou plusieurs colonnes à utiliser pour le partitionnement dans la table. Exclusif mutuellement avec
CLUSTER BY.Le clustering liquide offre une solution flexible et optimisée pour le clustering. Envisagez d'utiliser
CLUSTER BYau lieu dePARTITIONED BYpour les pipelines. -
CLUSTER BY
Activez le clustering liquide sur la table et définissez les colonnes à utiliser comme clés de clustering. Utilisez le clustering liquide automatique avec
CLUSTER BY AUTO, et Databricks choisit intelligemment les clés de clustering pour optimiser les performances de la query. Mutuellement exclusif avecPARTITIONED BY. -
Emplacement
Un emplacement de stockage facultatif pour les données de table. S'il n'est pas défini, le système utilise par default l'emplacement de stockage du pipeline.
-
Commentaire
Un littéral
STRINGfacultatif pour décrire la table. -
TBLPROPERTIES
Une liste facultative de propriétés de table pour la table.
-
AVEC FILTRE DE LIGNE
Ajoute une fonction de filtre de ligne à la table. Les futures requêtes pour cette table reçoivent un sous-ensemble des lignes pour lesquelles la fonction évalue à VRAI. Ceci est utile pour le contrôle d'accès précis, car cela permet à la fonction d'inspecter l'identité et les appartenances aux groupes de l'utilisateur appelant afin de décider s'il faut filtrer certaines lignes.
Voir
ROW FILTERclause.-
FLUX
Définit facultativement un flux en ligne avec la création de table. Un flux est une requête avec état qui refresh le contenu de la table. Si
FLOWn'est pas spécifié, vous pouvez utiliserAS queryà la place, ou définir les flux séparément avecCREATE FLOW. Vous pouvez spécifier l'un des types de flux suivants :- Insérer par nom
Insère des données dans la table par nom de colonne. Si l'option
ONCEn'est pas fournie, la requête doit être une requête de streaming. Utilisez le mot-cléSTREAMpour utiliser la sémantique de streaming pour lire depuis la source. Si la lecture rencontre une modification ou une suppression d'un enregistrement existant, une erreur est générée. Il est plus sûr de lire à partir de sources statiques ou en mode ajout uniquement.
-
FLOW INSERT BY NAME est équivalent à l'utilisation de AS query. Les deux instructions suivantes ont un comportement identique :
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');
- Une fois
Définit éventuellement le flux comme un flux unique, tel qu'un remplissage. Lorsque ONCE est fourni, la requête n'est pas une requête de streaming, et le flux s'exécute une seule fois par default. Si la table est refresh avec une full refresh, le ONCE flow runs again to recreate the data. ONCE ne s'applique qu'à INSERT BY NAME flux.
AUTO CDC
Bêta
Disponible dans Databricks Runtime 17.3 et versions ultérieures et le canal de distribution PREVIEW Pipelines.
Définit un flux AUTO CDC qui traite les enregistrements de capture de données modifiées (CDC) d'une source vers la table. Utilisez AUTO CDC lorsque les données source incluent la sémantique de la CDC. Voir Les AUTO CDC APIs : Simplifier la capture de données modifiées avec les pipelines.
- REPLACE WHERE prédicat BY NAME query
Bêta
FLOW REPLACE WHERE est en bêta.
Définit un flux REPLACE WHERE qui recalcule et écrase uniquement les lignes correspondant à predicate, laissant toutes les autres lignes intactes. Utilisez REPLACE WHERE pour le traitement batch incrémentiel des jointures et des agrégations, les données à arrivée tardive, l'évolution des schémas et les rattrapages. BY NAME est requis. Voir Traitement par batch avec des flux REPLACE WHERE.
-
Cette clause remplit la table en utilisant les données de
query. Cette query doit être une requête streaming . Utilisez le mot-clé STREAM pour utiliser la sémantique de streaming afin de lire à partir de la source. Si la lecture rencontre une modification ou une suppression d'un enregistrement existant, une erreur est générée. Il est plus sûr de lire à partir de sources statiques ou à ajout seulement. Pour ingérer des données avec des commits de modification, vous pouvez ajouter l'option de lectureskipChangeCommitspour gérer les erreurs.Lorsque vous spécifiez un
queryet untable_specificationensemble, le schéma de table spécifié danstable_specificationdoit contenir toutes les colonnes renvoyées par lequery, sinon vous obteniez une erreur. Toutes les colonnes spécifiées danstable_specificationmais non renvoyées parqueryrenvoient des valeursnulllorsqu'elles sont interrogées.Pour plus d'informations sur le streaming de données, consultez Transformer des données avec des pipelines.
-
Options de lecture
Vous pouvez spécifier des options de lecture dans la query pour configurer la façon dont les données sont lues à partir de la source. Par exemple, vous pouvez spécifier
skipChangeCommitspour ignorer tout commit de modification dans les données source. Les options de lecture sont spécifiées sous forme de mappage dans la clauseWITHde la query. Par exemple :SQLSELECT * FROM STREAM source_table WITH (SKIPCHANGECOMMITS=TRUE, STARTINGVERSION=X)Le
=TRUEest facultatif, vous pouvez donc également spécifier une option booléenne comme ceci :SQLSELECT * FROM STREAM source_table WITH (SKIPCHANGECOMMITS)
-
Les options de lecture ne sont prises en charge que pour Databricks Runtime 17.3 et les versions ultérieures.
Les options de lecture ci-dessous sont prises en charge pour Delta. Pour plus de détails sur chaque option, consultez lectures et écritures en streaming de tables Delta Lake.
maxFilesPerTriggermaxBytesPerTriggerstartingVersionstartingTimestampreadChangeFeedwithEventTimeOrderskipChangeCommits
Autorisations requises
L'utilisateur d'exécution pour un pipeline doit avoir les autorisations suivantes :
SELECTprivilège sur les tables de base référencées par la table de streaming.USE CATALOGprivilège sur le catalogue parent et le privilègeUSE SCHEMAsur le schéma parent.CREATE MATERIALIZED VIEWprivilège sur le schéma pour la table de streaming.
Pour qu'un utilisateur puisse mettre à jour le pipeline dans lequel la table de streaming est définie, il doit disposer des éléments suivants :
USE CATALOGprivilège sur le catalogue parent et le privilègeUSE SCHEMAsur le schéma parent.- Propriété de la table de streaming ou privilège
REFRESHsur la table de streaming. - Le propriétaire de la table de streaming doit avoir le privilège
SELECTsur les tables de base référencées par la table de streaming.
Pour qu'un utilisateur puisse **query** la table de **streaming** résultante, il a besoin :
USE CATALOGprivilège sur le catalogue parent et le privilègeUSE SCHEMAsur le schéma parent.SELECTprivilège sur la table de streaming.
Limitations
-
Seuls les propriétaires de tables peuvent refresh les tables de streaming pour obtenir les données les plus récentes.
-
ALTER TABLEles commandes sont interdites sur les tables de streaming. La définition et les propriétés de la table doivent être modifiées via leCREATE OR REFRESHou l'instruction ALTER STREAMING TABLE. -
L'évolution du schéma de table par le biais de commandes DML telles que
INSERT INTOetMERGEn'est pas prise en charge. -
Les commandes suivantes ne sont pas prises en charge sur les tables de streaming :
CREATE TABLE ... CLONE <streaming_table>COPY INTOANALYZE TABLERESTORETRUNCATEGENERATE MANIFEST[CREATE OR] REPLACE TABLE
-
Le renommage de la table ou la modification du propriétaire n'est pas pris en charge.
Exemples
-- 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;