Aller au contenu principal

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 CONNECTION privilè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 CREATE et USE SCHEMA sur 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.

L'interface utilisateur Databricks déploie des pipelines basés sur des query vers le compute serverless.

  1. Dans la barre latérale du Databricks workspace, cliquez sur Ingestion des données .

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

  3. Sur la page Pipeline d'ingestion , saisissez un nom pour le pipeline.

  4. Pour le catalogue de destination , sélectionnez un catalogue Unity Catalog pour stocker les données ingérées.

  5. 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 CONNECTION sur le métastore.

  6. Cliquez sur Créer un pipeline et continuer .

  7. Sur la page Source , sélectionnez les schémas et les tables à ingérer.

  8. 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_at ou row_id). Si vous ne sélectionnez pas de colonne de curseur à augmentation monotone, le connecteur effectuera un chargement complet à chaque exécution.

  9. 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).

  10. Cliquez sur Suivant .

  11. 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 CATALOG et CREATE SCHEMA sur le catalogue parent.

  12. Cliquez sur Enregistrer et continuer .

  13. (Facultatif) Sur la page **Paramètres**, cliquez sur **Créer une planification** et définissez la fréquence de refresh.

  14. (Facultatif) Configurez les notifications par e-mail pour le succès ou l'échec du pipeline.

  15. Cliquez sur Enregistrer et exécuter le pipeline .

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.

L'interface utilisateur Databricks déploie des pipelines basés sur des query vers le compute serverless.

  1. Dans la barre latérale du Databricks workspace, cliquez sur Ingestion des données .

  2. Sur la page Ajouter des données , sous Connecteurs Databricks , cliquez sur votre source. L'assistant d'ingestion s'ouvre.

  3. Sur la page Pipeline d'ingestion , saisissez un nom pour le pipeline.

  4. Pour le catalogue de destination , sélectionnez un catalogue Unity Catalog pour stocker les données ingérées.

  5. Pour le Type de connexion , sélectionnez Catalogue étranger , puis choisissez le catalogue étranger enregistré dans Lakehouse Federation.

  6. Cliquez sur Créer un pipeline et continuer .

  7. Sur la page Source , sélectionnez les schémas et les tables à ingérer.

  8. 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_at ou row_id).

  9. 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).

  10. Cliquez sur Suivant .

  11. 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 CATALOG et CREATE SCHEMA sur le catalogue parent.

  12. Cliquez sur Enregistrer et continuer .

  13. (Facultatif) Sur la page **Paramètres**, cliquez sur **Créer une planification** et définissez la fréquence de refresh.

  14. (Facultatif) Configurez les notifications par e-mail pour le succès ou l'échec du pipeline.

  15. Cliquez sur Enregistrer et exécuter le pipeline .

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_at ou last_modified sont 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 id ou row_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.

Ressources supplémentaires