Aller au contenu principal

Utiliser des tables de streaming autonomes

Une *table de streaming* autonome est une table enregistrée dans Unity Catalog avec un support supplémentaire pour le streaming ou le traitement incrémental des données, définie en dehors d'un LakeFlow pipeline. Un pipeline est automatiquement créé pour chaque table de streaming. Vous pouvez utiliser des tables de streaming pour le chargement incrémental des données depuis Kafka et le stockage d'objets cloud.

Vous pouvez créer et refresh des tables de streaming autonomes à partir d’un warehouse Databricks SQL, ou d’un Notebook exécuté sur un compute général serverless. Pour plus de détails sur les différences entre les deux options de compute, consultez Exigences pour les pipelines autonomes.

Pour créer et refresh des tables de streaming autonomes avec Python à partir d'un Notebook, consultez Utiliser Python avec des pipelines autonomes.

remarque

Pour savoir comment utiliser les tables Delta Lake comme sources et récepteurs de streaming, consultez Lectures et écritures en streaming des tables Delta Lake.

Exigences

Pour les options de compute, les autorisations et les autres exigences pour la création, l'actualisation et l'interrogation de tables de streaming autonomes, consultez Exigences pour les pipelines autonomes.

Créer des tables de streaming

Une table de streaming est définie par une requête SQL dans Databricks SQL. Lorsque vous créez une table de streaming, les données actuellement présentes dans les tables sources sont utilisées pour construire la table de streaming. Après cela, vous refresh la table, généralement selon un calendrier, pour extraire toutes les données ajoutées dans les tables sources afin de les ajouter à la table de streaming.

Lorsque vous créez une table de streaming, vous êtes considéré comme le propriétaire de la table.

Pour créer une table de streaming à partir d'une table existante, utilisez l'instruction CREATE STREAMING TABLE, comme dans l'exemple suivant :

SQL
CREATE OR REFRESH STREAMING TABLE sales
SCHEDULE EVERY 1 hour
AS SELECT product, price FROM STREAM raw_data;

Dans ce cas, la table de streaming sales est créée à partir de colonnes spécifiques de la table raw_data, avec une planification de refresh toutes les heures. La query utilisée doit être une query streaming . Utilisez le mot-clé STREAM pour utiliser la sémantique streaming afin de lire à partir de la source.

Compute utilisé pour le refresh

Lorsque vous créez une table de streaming à l'aide de l'instruction CREATE OR REFRESH STREAMING TABLE, le refresh initial des données et leur population commencent immédiatement. Ces opérations ne consomment pas le compute des warehouses Databricks SQL. En revanche, les tables de streaming s’appuient sur des pipelines serverless pour la création et le refresh. Un pipeline serverless dédié est automatiquement créé et géré par le système pour chaque table de streaming.

Chargez des fichiers avec Auto Loader.

Pour créer une table de streaming à partir de fichiers dans un volume, vous utilisez Auto Loader. Utilisez Auto Loader pour la plupart des tâches d'ingestion de données à partir du stockage d'objets cloud. Auto Loader et les pipelines sont conçus pour charger de manière incrémentielle et idempotente les données en constante augmentation au fur et à mesure qu'elles arrivent dans le stockage cloud.

Pour utiliser Auto Loader dans Databricks SQL, utilisez la fonction read_files. L'exemple suivant montre comment utiliser Auto Loader pour lire un volume de fichiers JSON dans une table de streaming :

SQL
CREATE OR REFRESH STREAMING TABLE sales
SCHEDULE EVERY 1 hour
AS SELECT * FROM STREAM read_files(
"/Volumes/my_catalog/my_schema/my_volume/path/to/data",
format => "json"
);

Pour lire des données depuis le stockage cloud, vous pouvez également utiliser Auto Loader :

SQL
CREATE OR REFRESH STREAMING TABLE sales
SCHEDULE EVERY 1 hour
AS SELECT *
FROM STREAM read_files(
's3://mybucket/analysis/*/*/*.json',
format => "json"
);

Pour en savoir plus sur Auto Loader, consultez Qu’est-ce qu’Auto Loader ?. Pour en savoir plus sur l’utilisation d’Auto Loader dans SQL, avec des exemples, consultez Charger les données à partir du stockage d’objets.

Ingestion en streaming à partir d'autres sources

Pour un exemple d'ingestion à partir d'autres sources, y compris Kafka, consultez Charger des données dans les pipelines.

Appliquez la capture de données de changement (CDC) avec les flux CDC automatiques.

Utilisez la clause FLOW AUTO CDC pour traiter les enregistrements de capture des changements de données (CDC) d'une source dans une table de streaming. Auparavant, l'instruction MERGE INTO était couramment utilisée pour traiter les enregistrements CDC sur Databricks. Cependant, MERGE INTO peut produire des résultats incorrects en raison d'enregistrements hors séquence ou nécessiter une logique complexe pour réordonner les enregistrements. Consultez la capture des changements de données et les instantanés.

AUTO CDC simplifie la CDC en gérant automatiquement les enregistrements désordonnés. Vous spécifiez les clés pour identifier les enregistrements, une colonne de séquence pour l'ordonnancement, et si les résultats doivent être stockés en tant que SCD de type 1 (mises à jour directes) ou SCD de type 2 (suivi de l'historique).

L'exemple suivant crée une table en streaming qui applique les modifications CDC en utilisant le SCD de type 1 :

SQL
CREATE OR REFRESH STREAMING TABLE target
FLOW AUTO CDC
FROM stream(cdc_data.users)
KEYS (userId)
SEQUENCE BY sequenceNum
STORED AS SCD TYPE 1;

L'exemple suivant utilise SCD de type 2 pour conserver un historique des modifications :

SQL
CREATE OR REFRESH STREAMING TABLE target
FLOW AUTO CDC
FROM stream(cdc_data.users)
KEYS (userId)
APPLY AS DELETE WHEN operation = "DELETE"
SEQUENCE BY sequenceNum
COLUMNS * EXCEPT (operation, sequenceNum)
STORED AS SCD TYPE 2;

Pour plus de détails sur les options et le comportement de l’API CDC, consultez Les APIs AUTO CDC : Simplifiez la capture des données de modification avec les pipelines. Pour la référence syntaxique complète, consultez CREATE STREAMING TABLE.

Appliquer le remplacement de batch sélectif avec les flux REPLACE WHERE

info

Bêta

Cette fonctionnalité est en Bêta.

Utilisez la clause FLOW REPLACE WHERE pour recalculer et écraser un sous-ensemble ciblé d'une table de streaming sans retraiter l'historique complet de votre table. Les flux REPLACE WHERE sont bien adaptés au traitement par batch incrémentiel des jointures et des agrégations, aux données à arrivée tardive, au retraitement en amont, à l'évolution des schémas et aux remplissages rétrospectifs.

Pour plus de détails sur les flux REPLACE WHERE, y compris les exigences, les remplacements de prédicats et l'incremental refresh, consultez Flux REPLACE WHERE pour les tables de streaming autonomes.

Ingérer uniquement les nouvelles données

Par default, la fonction read_files lit toutes les données existantes dans le dossier source lors de la création de la table, puis traite les nouveaux enregistrements à chaque refresh.

Pour éviter d'ingérer des données qui existent déjà dans le dossier source au moment de la création de la table, définissez l'option includeExistingFiles sur false. Cela signifie que seules les données qui arrivent dans le dossier après la création de la table sont traitées. Par exemple :

SQL
CREATE OR REFRESH STREAMING TABLE sales
SCHEDULE EVERY 1 hour
AS SELECT *
FROM STREAM read_files(
'/path/to/files',
includeExistingFiles => false
);

Définir le canal d’exécution

Les tables de streaming créées à l'aide de SQL warehouses sont automatiquement actualisées à l'aide d'un pipeline. Les pipelines utilisent le runtime dans le canal current par default. Pour en savoir plus sur le processus de publication, consultez les notes de version de LakeFlow Pipelines et le processus de mise à niveau de la publication.

Databricks recommande d'utiliser le current canal de distribution pour les workloads de production. Les nouvelles fonctionnalités sont d'abord publiées sur le preview canal de distribution. Vous pouvez définir un pipeline sur le canal de prévisualisation pour tester de nouvelles fonctionnalités en spécifiant preview comme propriété de table à l'aide d'une instruction CREATE OR REFRESH STREAMING TABLE. Pour mettre à jour le canal de distribution d'une table de streaming existante, vous devez exécuter CREATE OR REFRESH STREAMING TABLE avec le TBLPROPERTIES mis à jour.

L'exemple de code suivant montre comment définir le Canal de distribution sur aperçu :

SQL
CREATE OR REFRESH STREAMING TABLE sales
TBLPROPERTIES ('pipelines.channel' = 'preview')
SCHEDULE EVERY 1 hour
AS SELECT *
FROM STREAM raw_data;

Masquer les données sensibles

Vous pouvez utiliser des tables de streaming pour masquer les données sensibles aux utilisateurs accédant à la table. Une approche consiste à définir la query de manière à ce qu’elle exclue entièrement les colonnes ou les lignes sensibles. Vous pouvez également appliquer des masques de colonne ou des filtres de ligne en fonction des autorisations de l'utilisateur qui effectue la requête. Par exemple, vous pourriez masquer la colonne tax_id pour les utilisateurs qui ne font pas partie du groupe HumanResourcesDept. Pour ce faire, utilisez la syntaxe ROW FILTER et MASK lors de la création de la table de streaming. Pour en savoir plus, consultez Filtres de ligne et masques de colonne.

Refresh une table de streaming

Les tables de streaming créent et utilisent automatiquement des pipelines serverless pour traiter les opérations de refresh. Le refresh est géré par le pipeline et la mise à jour est surveillée par le warehouse Databricks SQL utilisé pour créer la table de streaming. Les tables de streaming peuvent être mises à jour à l'aide d'un pipeline qui s'exécute selon une planification.

Même si vous avez une refresh planifiée, vous pouvez déclencher une refresh manuelle à tout moment. Les actualisations sont gérées par le même pipeline qui a été automatiquement créé avec la table de streaming.

Pour refresh une table de streaming :

SQL
REFRESH STREAMING TABLE sales;

Vous pouvez vérifier l'état du dernier refresh avec DESCRIBE TABLE EXTENDED.

remarque

Vous devrez peut-être refresh votre table de streaming avant d'utiliser les queries time travel.

Pour savoir comment planifier une actualisation, consultez Planifier les refresh. Les **refresh** planifiés peuvent avoir des notifications de mise à jour, et vous pouvez définir le mode de performance pour le **refresh**.

Fonctionnement du refresh

Un refresh de table de streaming évalue uniquement les nouvelles lignes arrivées après la dernière mise à jour et ajoute uniquement les nouvelles données.

Chaque refresh utilise la définition actuelle de la table de streaming pour traiter ces nouvelles données. La modification d'une définition de table streaming ne recalcule pas automatiquement les données existantes. Si une modification est incompatible avec les données existantes (par exemple, la modification d'un type de données), le prochain refresh échouera avec une erreur.

Les exemples suivants expliquent comment les modifications apportées à une définition de table de streaming affectent le comportement de refresh :

  • La suppression d’un filtre ne reprocesse pas les lignes précédemment filtrées.
  • La modification des projections de colonne n'affecte pas la manière dont les données existantes ont été traitées.
  • Les jointures avec des instantanés statiques utilisent l'état de l'instantané au moment du traitement initial. Les données arrivées tardivement qui auraient correspondu à l'instantané mis à jour sont ignorées. Cela peut entraîner la perte de faits si les dimensions sont en retard.
  • La modification du CAST d'une colonne existante entraîne une erreur.

Si vos données changent d'une manière qui ne peut pas être prise en charge dans la table de streaming existante, vous pouvez effectuer une refresh complète.

Fully refresh a table de streaming

Les actualisations complètes retraitent toutes les données disponibles dans la source avec la définition la plus récente. Il n'est pas recommandé d'appeler des refresh complets sur des sources qui ne conservent pas l'historique complet des données ou qui ont de courtes périodes de rétention, telles que Kafka, car le refresh complet tronque les données existantes. Il est possible que vous ne puissiez pas récupérer les anciennes données si celles-ci ne sont plus disponibles dans la source.

Par exemple :

SQL
REFRESH STREAMING TABLE sales FULL;

Planifier et surveiller les actualisations

Vous pouvez refresh une table de streaming automatiquement selon un planning ou lorsque les données en amont changent, et vous pouvez configurer les délais de refresh, les notifications et les modes de performance. Consultez Planifier les refresh.

Contrôler l'accès aux tables de streaming

Les tables de streaming prennent en charge des contrôles d'accès enrichis pour faciliter le Data Sharing tout en évitant d'exposer des données potentiellement privées. Un propriétaire de table de streaming ou un utilisateur disposant du privilège MANAGE peut accorder les privilèges SELECT à d’autres utilisateurs. Les utilisateurs ayant un accès SELECT à la table de streaming ne nécessitent pas d'accès SELECT aux tables référencées par la table de streaming. Ce contrôle d'accès permet le Data Sharing tout en contrôlant l'accès aux données sous-jacentes.

Vous pouvez également modifier le propriétaire d’une table de streaming.

Accorder des privilèges à une table de streaming

Pour accorder l'accès à une table de streaming, utilisez l'instruction GRANT:

SQL
GRANT <privilege_type> ON <st_name> TO <principal>;

Le privilege_type peut être :

  • SELECT - l'utilisateur peut SELECT la table de streaming.
  • REFRESH - l'utilisateur peut REFRESH la table de streaming. Les actualisations sont exécutées en utilisant les autorisations du propriétaire.

L'exemple suivant crée une table de streaming et accorde les privilèges de sélection et de refresh aux utilisateurs :

SQL
CREATE OR REFRESH STREAMING TABLE st_name AS SELECT * FROM source_table;

-- Grant read-only access:
GRANT SELECT ON st_name TO read_only_user;

-- Grant read and refresh access:
GRANT SELECT ON st_name TO refresh_user;
GRANT REFRESH ON st_name TO refresh_user;

Pour plus d'informations sur l'octroi de privilèges sur les objets sécurisables d'Unity Catalog, veuillez consulter la référence des privilèges d'Unity Catalog.

Révoquer les privilèges d'une table de streaming

Pour révoquer l'accès d'une table de streaming, utilisez l'instruction REVOKE:

SQL
REVOKE privilege_type ON <st_name> FROM principal;

Lorsque les privilèges SELECT sur une table source sont révoqués du propriétaire de la table de streaming ou de tout autre utilisateur qui s'est vu accorder les privilèges MANAGE ou SELECT sur la table de streaming, ou que la table source est supprimée, le propriétaire de la table de streaming ou l'utilisateur ayant reçu l'accès peut toujours interroger la table de streaming. Cependant, le comportement suivant se produit :

  • Le propriétaire de la table de streaming ou d'autres personnes ayant perdu l'accès à une table de streaming ne peuvent plus REFRESH cette table de streaming, et la table de streaming devient obsolète au fil du temps.
  • Si elle est automatisée avec un calendrier, le REFRESH planifié suivant échoue ou n'est pas exécuté.

L'exemple suivant révoque le privilège SELECT de read_only_user:

SQL
REVOKE SELECT ON st_name FROM read_only_user;

Modifier le propriétaire d'une table de streaming

Un utilisateur avec des autorisations MANAGE sur une table de streaming autonome peut définir un nouveau propriétaire via l’Explorateur de catalogues. Le nouveau propriétaire peut être lui-même ou un service principal sur lequel il a le rôle d'utilisateur de service principal.

  1. Depuis votre Databricks workspace, cliquez Icône de données. sur **Catalog** pour ouvrir l'Explorateur de catalogues.

  2. Sélectionnez la table de streaming que vous souhaitez mettre à jour.

  3. Dans la barre latérale droite, sous À propos de cette table de streaming , recherchez le Propriétaire , et cliquez sur Icône de crayon. modifier.

remarque

Si vous recevez un message vous demandant de mettre à jour le propriétaire en modifiant l'utilisateur **Exécuter en tant que** dans les paramètres du pipeline, alors la table de streaming est définie dans un LakeFlow Pipelines, et non comme une table autonome. Le message inclut un Link vers les paramètres du pipeline, où vous pouvez modifier l'utilisateur **Exécuter en tant que**.

  1. Sélectionner un nouveau propriétaire pour la table de streaming.

    Les propriétaires disposent automatiquement des privilèges MANAGE et SELECT sur les tables de streaming qu'ils possèdent. Si vous définissez un Service Principal comme propriétaire d'une table de streaming que vous possédez, et que vous ne disposez pas explicitement des privilèges SELECT ou MANAGE sur la table de streaming, alors ce changement vous ferait perdre tout accès à la table de streaming. Dans ce cas, vous êtes invité à fournir explicitement ces privilèges.

    Sélectionnez à la fois les privilèges **Grant MANAGE** et **Grant SELECT** pour les appliquer lors de l'**enregistrement**.

  2. Cliquez sur Enregistrer pour modifier le propriétaire.

Le propriétaire de la table de streaming est mis à jour. Toutes les futures mises à jour sont exécutées à l'aide de la nouvelle identité du propriétaire.

Lorsque le propriétaire perd les privilèges sur les tables sources

Si vous changez le propriétaire et que le nouveau propriétaire n'a pas accès aux tables source (ou si les privilèges SELECT sont révoqués sur les tables source sous-jacentes), les utilisateurs peuvent toujours query la table de streaming. Cependant :

  • Ils ne peuvent pas REFRESH la table de streaming.
  • La prochaine refresh planifiée de la table de streaming échoue.

Perdre l'accès aux données source empêche les mises à jour, mais n'invalide pas immédiatement la lecture de la table de streaming existante.

Supprimer définitivement les enregistrements d'une table de streaming

info

Aperçu

La prise en charge de l'instruction REORG avec les tables de streaming est en Aperçu public.

remarque
  • L'utilisation d'une instruction REORG avec une table de streaming nécessite Databricks Runtime 15.4 ou une version ultérieure.
  • Bien que vous puissiez utiliser l'instruction REORG avec n'importe quelle table de streaming, ce n'est requis que lors de la suppression d'enregistrements d'une table de streaming avec les vecteurs de suppression activés. La commande n'a aucun effet lorsqu'elle est utilisée avec une table de streaming sans vecteurs de suppression activés.

Pour supprimer physiquement les enregistrements du stockage sous-jacent d'une table de streaming avec des vecteurs de suppression activés, par exemple pour la conformité GDPR, des mesures supplémentaires doivent être prises afin de garantir qu'une opération VACUUM s'exécute sur les données de la table de streaming.

Pour supprimer physiquement les enregistrements du stockage sous-jacent :

  1. Mettre à jour ou supprimer des enregistrements de la table de streaming.
  2. Exécutez une instruction REORG sur la table de streaming, en spécifiant le parameter APPLY (PURGE). Par exemple REORG TABLE <streaming-table-name> APPLY (PURGE);.
  3. Attendez que la période de rétention des données de la table de streaming s'écoule. La période de rétention des données par default est de sept jours, mais elle peut être configurée avec la propriété de table delta.deletedFileRetentionDuration. Consultez Configurer la rétention des données pour les query Time Travel.
  4. REFRESH la table de streaming. Consultez refresh une table de streaming. Dans les 24 heures suivant l’opération REFRESH, les tâches de maintenance du pipeline, y compris l’opération VACUUM requise pour garantir la suppression permanente des enregistrements, sont exécutées automatiquement.

Supervisez les exécutions à l'aide de l'historique des query.

Vous pouvez utiliser la page d'historique des query pour accéder aux détails des query et aux profils de query qui peuvent vous aider à identifier les query peu performantes et les goulots d'étranglement dans le pipeline utilisé pour exécuter les mises à jour de vos tables de streaming. Pour un aperçu du type d'informations disponibles dans les historiques de requêtes et les profils de requêtes, consultez Historique des requêtes et Profil de requête.

info

Aperçu

Cette fonctionnalité est en Aperçu public. Les administrateurs du Workspace peuvent contrôler l'accès à cette fonctionnalité à partir de la page Previews . Consultez Gérer les aperçus Databricks.

Toutes les déclarations relatives aux tables de streaming apparaissent dans l'historique des query. Vous pouvez utiliser le filtre déroulant **Statement** pour sélectionner n'importe quelle commande et inspecter les query associées. Toutes les CREATE déclarations sont suivies d'une déclaration REFRESH qui s'exécute de manière asynchrone sur un pipeline. Les REFRESH déclarations incluent généralement des plans de query détaillés qui fournissent des aperçus pour l'optimisation des performances.

Pour accéder aux REFRESH déclarations dans l'interface utilisateur de l'historique des query, suivez les étapes suivantes :

  1. Cliquez Icône Historique. sur dans la barre latérale gauche pour ouvrir l'interface utilisateur de l'**Historique des query**.
  2. Sélectionnez la case à cocher REFRESH dans le filtre déroulant Instruction .
  3. Cliquez sur le nom de l'instruction de query pour afficher les détails du résumé, tels que la durée de la query et les métriques agrégées.
  4. Cliquez sur Voir le profil de query pour ouvrir le profil de query. Consultez le profil de query pour plus de détails sur la navigation dans le profil de query.
  5. En option, vous pouvez utiliser les liens de la section Source de la query pour ouvrir la query ou le pipeline associé.

Vous pouvez également accéder aux détails de la query en utilisant les liens dans l'éditeur SQL ou à partir d'un notebook attaché à un SQL warehouse.

Accéder aux tables de streaming à partir de clients externes

Pour accéder aux tables de streaming à partir de clients Delta Lake ou Iceberg externes qui ne prennent pas en charge les APIs ouvertes, vous pouvez utiliser le Mode de compatibilité. Le mode de compatibilité crée une version en lecture seule de votre table de streaming qui peut être consultée par tout client Delta Lake ou Iceberg.

Ressources supplémentaires