Aller au contenu principal

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.

    • column_identifier

      Les noms de colonnes doivent être uniques et correspondre aux colonnes de sortie de la query.

    • column_type

      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 STRING facultatif décrivant la colonne. Cette option doit être spécifiée avec column_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 expr spécifié.

      Le/La DEFAULT COLLATION de la table doit être UTF8_BINARY.

      expr peut ê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 :

      De plus, expr ne doit pas contenir de sous-requête.

    • GENERATED { ALWAYS | BY DEFAULT } AS IDENTITY [ ( [ START WITH start ] [ INCREMENT BY step ] ) ]

      S’applique à : coché oui Databricks SQL coché oui 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 step est 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 start et augmentent de step. 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. step ne peut pas être 0.

      Si les valeurs attribuées automatiquement dépassent la plage du type de colonne d'identité, la query échouera.

      Lorsque ALWAYS est 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 BY une colonne d'identité
      • UPDATE une colonne d'identité
remarque

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 à : coché oui Databricks SQL coché oui Databricks Runtime 11.3 LTS et versions ultérieures

    Définit une valeur DEFAULT pour la colonne utilisée sur INSERT, UPDATE et MERGE ... INSERT lorsque la colonne n'est pas spécifiée.

    Si aucun default n'est spécifié, DEFAULT NULL est appliqué aux colonnes acceptant les valeurs nulles.

    default_expression peuvent être composés de littéraux, et de fonctions ou d'opérateurs SQL intégrés sauf :

    De plus, default_expression ne doit pas contenir de sous-requête.

    DEFAULT est pris en charge pour les sources CSV, JSON, PARQUET et ORC.

  • contrainte de colonne

    Ajoute une contrainte de clé primaire informative ou de clé étrangère informative à la colonne dans une table de streaming.

  • Clause MASK

    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 UPDATE provoque l’échec du traitement lors de la création de la table ainsi que de l’actualisation de la table. Une attente DROP ROW provoque 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_expr peut ê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 :

    De plus, expr ne doit pas contenir de sous-requête.

  • 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.

remarque

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 BY au lieu de PARTITIONED BY pour 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 avec PARTITIONED BY.

      Consultez Utiliser le clustering liquide pour les tables.

    • 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 STRING facultatif 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 FILTER clause.

    • 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 FLOW n'est pas spécifié, vous pouvez utiliser AS query à la place, ou définir les flux séparément avec CREATE 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 ONCE n'est pas fournie, la requête doit être une requête de streaming. Utilisez le mot-clé STREAM pour 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.

remarque

FLOW INSERT BY NAME est équivalent à l'utilisation de AS query. Les deux instructions suivantes ont un comportement identique :

SQL
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
info

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.

info

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.

  • AS query

    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 lecture skipChangeCommits pour gérer les erreurs.

    Lorsque vous spécifiez un query et un table_specification ensemble, le schéma de table spécifié dans table_specification doit contenir toutes les colonnes renvoyées par le query, sinon vous obteniez une erreur. Toutes les colonnes spécifiées dans table_specification mais non renvoyées par query renvoient des valeurs null lorsqu'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 skipChangeCommits pour ignorer tout commit de modification dans les données source. Les options de lecture sont spécifiées sous forme de mappage dans la clause WITH de la query. Par exemple :

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

      Le =TRUE est facultatif, vous pouvez donc également spécifier une option booléenne comme ceci :

      SQL
      SELECT * FROM STREAM source_table WITH (SKIPCHANGECOMMITS)
remarque

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.

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

Autorisations requises

L'utilisateur d'exécution pour un pipeline doit avoir les autorisations suivantes :

  • SELECT privilège sur les tables de base référencées par la table de streaming.
  • USE CATALOG privilège sur le catalogue parent et le privilège USE SCHEMA sur le schéma parent.
  • CREATE MATERIALIZED VIEW privilè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 CATALOG privilège sur le catalogue parent et le privilège USE SCHEMA sur le schéma parent.
  • Propriété de la table de streaming ou privilège REFRESH sur la table de streaming.
  • Le propriétaire de la table de streaming doit avoir le privilège SELECT sur 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 CATALOG privilège sur le catalogue parent et le privilège USE SCHEMA sur le schéma parent.
  • SELECT privilè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 TABLE les commandes sont interdites sur les tables de streaming. La définition et les propriétés de la table doivent être modifiées via le CREATE OR REFRESH ou l'instruction ALTER STREAMING TABLE.

  • L'évolution du schéma de table par le biais de commandes DML telles que INSERT INTO et MERGE n'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 INTO
    • ANALYZE TABLE
    • RESTORE
    • TRUNCATE
    • GENERATE MANIFEST
    • [CREATE OR] REPLACE TABLE
  • Le renommage de la table ou la modification du propriétaire n'est pas pris en charge.

Exemples

SQL
-- 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;