Créer un pipeline d'ingestion basé sur une query
Cette page explique comment créer un pipeline d'ingestion basé sur une query dans Lakeflow Connect.
Exigences
Avant de créer un pipeline d'ingestion basé sur une query, vous devez satisfaire aux exigences suivantes :
- Unity Catalog est activé pour votre workspace Databricks.
- Votre environnement de compute serverless permet la connectivité réseau à la base de données source. Consultez Réseau et Recommandations réseau pour Lakehouse Federation.
- Pour l'**ingestion de connexion externe** : vous disposez d'une connexion existante à la base de données source ou de
CREATE CONNECTIONprivilèges sur le métastore. Voir Connecter des sources d'ingestion gérées. - Pour l' ingestion de catalogue étranger : vous disposez d'un catalogue étranger existant enregistré dans Lakehouse Federation ou des privilèges pour en créer un.
- Vous disposez des privilèges
CREATEetUSE SCHEMAsur le catalogue et le schéma de destination.
Option 1 : Ingestion de connexions externes
Utilisez cette approche lorsque vous avez une connexion qui stocke les identifiants d'authentification pour la base de données source. Les sources prises en charge sont Oracle, Teradata, SQL Server, MySQL, MariaDB et PostgreSQL.
- Databricks UI
- Declarative Automation Bundles
L'interface utilisateur Databricks déploie des pipelines basés sur des query vers le compute serverless.
-
Dans la barre latérale du Databricks workspace, cliquez sur Ingestion des données .
-
Sur la page Ajouter des données , sous Connecteurs Databricks , cliquez sur votre source (par exemple, Oracle ou SQL Server ). L'assistant d'ingestion s'ouvre.
-
Sur la page Pipeline d'ingestion , saisissez un nom pour le pipeline.
-
Pour le catalogue de destination , sélectionnez un catalogue Unity Catalog pour stocker les données ingérées.
-
Sélectionnez la connexion Unity Catalog qui stocke les identifiants requis pour accéder à la base de données source.
S’il n’y a pas de connexion existante, cliquez sur Créer une connexion et saisissez les détails de la connexion. Vous devez disposer des privilèges
CREATE CONNECTIONsur le métastore. -
Cliquez sur Créer un pipeline et continuer .
-
Sur la page Source , sélectionnez les schémas et les tables à ingérer.
-
Pour chaque table, spécifiez la colonne de curseur . Il doit s'agir d'une seule colonne avec des valeurs qui augmentent de manière monotone (par exemple,
updated_atourow_id). Si vous ne sélectionnez pas de colonne de curseur à augmentation monotone, le connecteur effectuera un chargement complet à chaque exécution. -
Vous pouvez éventuellement modifier le paramètre de suivi de l'historique default. Pour plus d'informations, consultez Activer le suivi de l'historique (SCD de type 2).
-
Cliquez sur Suivant .
-
Sur la page Destination , sélectionnez le catalogue et le schéma Unity Catalog dans lesquels écrire.
Si vous ne souhaitez pas utiliser de schéma existant, cliquez sur Créer un schéma . Vous devez disposer des privilèges
USE CATALOGetCREATE SCHEMAsur le catalogue parent. -
Cliquez sur Enregistrer et continuer .
-
(Facultatif) Sur la page **Paramètres**, cliquez sur **Créer une planification** et définissez la fréquence de refresh.
-
(Facultatif) Configurez les notifications par e-mail pour le succès ou l'échec du pipeline.
-
Cliquez sur Enregistrer et exécuter le pipeline .
Déployez un pipeline d'ingestion basé sur une query à l'aide de Declarative Automation Bundles. Les bundles contiennent des définitions YAML de pipelines et de Job, sont gérés avec la CLI Databricks et peuvent être déployés sur plusieurs Workspace cibles. Pour plus d'informations, consultez What are Declarative Automation Bundles?.
Cet exemple déploie le pipeline vers le compute Serverless (default). Pour déployer sur un compute classique à la place, consultez l'exemple de compute classique.
-
Créer un bundle :
Bashdatabricks bundle init -
Ajouter un fichier de définition de pipeline au bundle (par exemple,
resources/query_based_pipeline.yml) :YAMLvariables:
dest_catalog:
default: main
dest_schema:
default: ingest_destination_schema
resources:
pipelines:
pipeline_query_based:
name: query-based-ingestion-pipeline
ingestion_definition:
connection_name: <your-uc-connection-name>
objects:
- table:
source_catalog: <source-catalog>
source_schema: <source-schema>
source_table: <source-table>
table_configuration:
query_based_connector_config:
cursor_columns:
- updated_at
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
target: ${var.dest_schema}
catalog: ${var.dest_catalog} -
Ajouter un fichier de définition de Job qui contrôle le planning d’ingestion (par exemple,
resources/query_based_job.yml) :YAMLresources:
jobs:
query_based_job:
name: query_based_job
trigger:
periodic:
interval: 1
unit: HOURS
email_notifications:
on_failure:
- <email-address>
tasks:
- task_key: refresh_pipeline
pipeline_task:
pipeline_id: ${resources.pipelines.pipeline_query_based.id} -
Déployez le bundle :
Bashdatabricks bundle deploy
Exemple de Classic compute (Bêta)
Beta
Classic compute for query-based ingestion pipelines is in Beta. Databricks recommends serverless compute for most workloads.
Pour déployer sur un compute classique, définissez serverless: false et ajoutez un bloc clusters à la définition du pipeline. Pour connaître l'ensemble complet des champs de cluster pris en charge, consultez Configurer le compute classique pour les pipelines.
-
Créer un bundle :
Bashdatabricks bundle init -
Ajouter un fichier de définition de pipeline au bundle (par exemple,
resources/query_based_pipeline.yml) :YAMLvariables:
dest_catalog:
default: main
dest_schema:
default: ingest_destination_schema
resources:
pipelines:
pipeline_query_based:
name: query-based-ingestion-pipeline
serverless: false
clusters:
- label: default
node_type_id: r6i.xlarge
driver_node_type_id: i3.large
autoscale:
min_workers: 1
max_workers: 5
ingestion_definition:
connection_name: <your-uc-connection-name>
objects:
- table:
source_catalog: <source-catalog>
source_schema: <source-schema>
source_table: <source-table>
table_configuration:
query_based_connector_config:
cursor_columns:
- updated_at
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
target: ${var.dest_schema}
catalog: ${var.dest_catalog} -
Ajouter un fichier de définition de Job qui contrôle le planning d’ingestion (par exemple,
resources/query_based_job.yml) :YAMLresources:
jobs:
query_based_job:
name: query_based_job
trigger:
periodic:
interval: 1
unit: HOURS
email_notifications:
on_failure:
- <email-address>
tasks:
- task_key: refresh_pipeline
pipeline_task:
pipeline_id: ${resources.pipelines.pipeline_query_based.id} -
Déployez le bundle :
Bashdatabricks bundle deploy
Option 2 : Ingestion de catalogue étranger
Utilisez cette approche lorsque vous souhaitez ingérer des données à partir d'un catalogue étranger enregistré dans Lakehouse Federation. L'ingestion de catalogue externe prend en charge toutes les sources de données Lakehouse Federation et le suivi des suppressions.
- Databricks UI
- Declarative Automation Bundles
L'interface utilisateur Databricks déploie des pipelines basés sur des query vers le compute serverless.
-
Dans la barre latérale du Databricks workspace, cliquez sur Ingestion des données .
-
Sur la page Ajouter des données , sous Connecteurs Databricks , cliquez sur votre source. L'assistant d'ingestion s'ouvre.
-
Sur la page Pipeline d'ingestion , saisissez un nom pour le pipeline.
-
Pour le catalogue de destination , sélectionnez un catalogue Unity Catalog pour stocker les données ingérées.
-
Pour le Type de connexion , sélectionnez Catalogue étranger , puis choisissez le catalogue étranger enregistré dans Lakehouse Federation.
-
Cliquez sur Créer un pipeline et continuer .
-
Sur la page Source , sélectionnez les schémas et les tables à ingérer.
-
Pour chaque table, spécifiez la colonne de curseur . Il doit s'agir d'une seule colonne avec des valeurs qui augmentent de manière monotone (par exemple,
updated_atourow_id). -
Vous pouvez éventuellement modifier le paramètre de suivi de l'historique default. Pour plus d'informations, consultez Activer le suivi de l'historique (SCD de type 2).
-
Cliquez sur Suivant .
-
Sur la page Destination , sélectionnez le catalogue et le schéma Unity Catalog dans lesquels écrire.
Si vous ne souhaitez pas utiliser de schéma existant, cliquez sur Créer un schéma . Vous devez disposer des privilèges
USE CATALOGetCREATE SCHEMAsur le catalogue parent. -
Cliquez sur Enregistrer et continuer .
-
(Facultatif) Sur la page **Paramètres**, cliquez sur **Créer une planification** et définissez la fréquence de refresh.
-
(Facultatif) Configurez les notifications par e-mail pour le succès ou l'échec du pipeline.
-
Cliquez sur Enregistrer et exécuter le pipeline .
Déployer un pipeline d'ingestion de catalogue étranger à l'aide de Bundles d'automatisation déclaratifs. Les bundles contiennent des définitions YAML de pipelines et de Job, sont gérés avec la CLI Databricks et peuvent être déployés sur plusieurs Workspace cibles. Pour plus d'informations, consultez What are Declarative Automation Bundles?.
Cet exemple déploie le pipeline vers le compute Serverless (default). Pour déployer sur un compute classique à la place, consultez l'exemple de compute classique.
-
Créer un bundle :
Bashdatabricks bundle init -
Ajouter un fichier de définition de pipeline au bundle (par exemple,
resources/foreign_catalog_pipeline.yml) :YAMLvariables:
dest_catalog:
default: main
dest_schema:
default: ingest_destination_schema
resources:
pipelines:
pipeline_foreign_catalog:
name: foreign-catalog-ingestion-pipeline
ingestion_definition:
ingest_from_uc_foreign_catalog: true
objects:
- table:
source_catalog: <foreign-catalog-name>
source_schema: <source-schema>
source_table: <source-table>
cursor_columns:
- updated_at
primary_keys:
- id
deletion_condition: 'deleted_at IS NOT NULL'
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
target: ${var.dest_schema}
catalog: ${var.dest_catalog} -
Ajoutez un fichier de définition de Job (par exemple,
resources/foreign_catalog_job.yml) :YAMLresources:
jobs:
foreign_catalog_job:
name: foreign_catalog_job
trigger:
periodic:
interval: 1
unit: HOURS
email_notifications:
on_failure:
- <email-address>
tasks:
- task_key: refresh_pipeline
pipeline_task:
pipeline_id: ${resources.pipelines.pipeline_foreign_catalog.id} -
Déployez le bundle :
Bashdatabricks bundle deploy
Exemple de Classic compute (Bêta)
Beta
Classic compute for query-based ingestion pipelines is in Beta. Databricks recommends serverless compute for most workloads.
Pour déployer sur un compute classique, définissez serverless: false et ajoutez un bloc clusters à la définition du pipeline. Pour connaître l'ensemble complet des champs de cluster pris en charge, consultez Configurer le compute classique pour les pipelines.
-
Créer un bundle :
Bashdatabricks bundle init -
Ajouter un fichier de définition de pipeline au bundle (par exemple,
resources/foreign_catalog_pipeline.yml) :YAMLvariables:
dest_catalog:
default: main
dest_schema:
default: ingest_destination_schema
resources:
pipelines:
pipeline_foreign_catalog:
name: foreign-catalog-ingestion-pipeline
serverless: false
clusters:
- label: default
node_type_id: r6i.xlarge
driver_node_type_id: i3.large
autoscale:
min_workers: 1
max_workers: 5
ingestion_definition:
ingest_from_uc_foreign_catalog: true
objects:
- table:
source_catalog: <foreign-catalog-name>
source_schema: <source-schema>
source_table: <source-table>
cursor_columns:
- updated_at
primary_keys:
- id
deletion_condition: 'deleted_at IS NOT NULL'
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
target: ${var.dest_schema}
catalog: ${var.dest_catalog} -
Ajoutez un fichier de définition de Job (par exemple,
resources/foreign_catalog_job.yml) :YAMLresources:
jobs:
foreign_catalog_job:
name: foreign_catalog_job
trigger:
periodic:
interval: 1
unit: HOURS
email_notifications:
on_failure:
- <email-address>
tasks:
- task_key: refresh_pipeline
pipeline_task:
pipeline_id: ${resources.pipelines.pipeline_foreign_catalog.id} -
Déployez le bundle :
Bashdatabricks bundle deploy
Configurer le suivi incrémentiel
Les connecteurs basés sur les requêtes utilisent une colonne de curseur pour déterminer quelles lignes sont nouvelles ou mises à jour après la dernière exécution du pipeline. Votre choix de colonne de curseur est crucial pour une ingestion incrémentale efficace.
Tenez compte des éléments suivants lorsque vous sélectionnez une colonne de curseur :
- Utilisez une colonne de Timestamp, si possible. Des colonnes comme
updated_atoulast_modifiedsont idéales car elles reflètent directement le moment où une ligne a été modifiée pour la dernière fois. - Les ID Integer fonctionnent pour les sources à ajout uniquement. Si les lignes ne sont jamais mises à jour, vous pouvez utiliser une colonne d'ID à incrémentation automatique (tels que
idourow_id) comme curseur. N'utilisez pas un ID entier comme curseur si les lignes peuvent être mises à jour sans changer l'ID. - La colonne doit augmenter de manière monotone. Les valeurs ne doivent jamais diminuer. Si un processus, tel qu'un backfill, définit la colonne sur une valeur passée, le connecteur ne réingère pas les lignes écrites avant le précédent seuil de sécurité.
- Vous ne pouvez spécifier qu’une seule colonne de curseur. Vous ne pouvez pas spécifier plusieurs colonnes comme curseur composite.
Une fois que le connecteur a stocké la marque de seuil supérieur du curseur, il utilise cette marque comme filtre de limite inférieure (cursor_column > last_value) lors de la prochaine exécution. Les lignes dont la valeur du curseur est NULLE ne sont pas ingérées.
Configurer le suivi de l'historique (SCD)
Pour suivre l'historique complet des modifications de lignes dans les tables de destination, configurez le SCD de type 2. Voir Activer le suivi de l'historique (SCD de type 2).
Modèles courants
Pour les configurations de pipeline avancées, consultez Modèles courants pour les pipelines d'ingestion gérés.