Ingérer des données depuis Celigo
Bêta
Cette fonctionnalité est en bêta. Les administrateurs du workspace peuvent contrôler l’accès à cette fonctionnalité à partir de la page Previews . See Manage Databricks previews.
Cette page indique comment créer un pipeline d’ingestion Celigo géré à l’aide de Lakeflow Connect.
Exigences
-
Pour créer un pipeline d'ingestion, veuillez d'abord respecter les conditions requises suivantes :
-
Votre Workspace doit être activé pour Unity Catalog.
-
Le compute serverless doit être activé pour votre workspace. Voir Exigences relatives au compute serverless.
-
Pour créer une nouvelle connexion, vous devez disposer des privilèges
CREATE CONNECTIONsur le métastore. Consultez la page Gérer les privilèges dans Unity Catalog.If the connector supports UI-based pipeline authoring, an admin can create the connection and the pipeline at the same time by completing the steps on this page. However, if the users who create pipelines use API-based pipeline authoring or are non-admin users, an admin must first create the connection in Catalog Explorer. See Connect to managed ingestion sources.
-
Pour utiliser une connexion existante, vous devez disposer des privilèges
USE CONNECTIONou deALL PRIVILEGESsur l'objet de connexion. -
Vous devez disposer des privilèges
USE CATALOGsur le catalogue cible. -
Vous devez disposer des privilèges
USE SCHEMAetCREATE TABLEsur un schéma existant ou du privilègeCREATE SCHEMAsur le catalogue cible.
-
-
Pour ingérer à partir de Celigo, commencez par configurer l’authentification depuis Databricks et créez une connexion. Consultez Configurer l’authentification pour Celigo et Créer une connexion Celigo.
Options du connecteur
Définissez les options de portée du pipeline dans source_configurations et les options de portée de la table sur l’objet individuel. Consultez la section Exemples pour connaître l’utilisation.
Option | Périmètre | Obligatoire | S'applique à | Description |
|---|---|---|---|---|
| Table | Non |
| Date-heure UTC ISO-8601 pour le start du backfill de la première synchronisation. default to 365 days before the first sync. |
Créer un pipeline d’ingestion
Pour consulter la liste des tables sources prises en charge, voir Tables sources prises en charge.
- Databricks UI
- Declarative Automation Bundles
- Databricks notebook
- Dans la barre latérale du workspace Databricks, cliquez sur Data Ingestion .
- Sur la page Add data , sous Databricks connectors , cliquez sur Celigo .
- Sur la page Connexion de l'assistant d'ingestion, sélectionnez la connexion qui stocke vos identifiants Celigo. Si vous disposez du privilège
CREATE CONNECTIONsur le métastore, cliquez surCréer une connexion pour créer une connexion avec les identifiants de la section Configurer l'authentification auprès de Celigo.
- Cliquez sur Suivant .
- Sur la page Ingestion setup , saisissez un nom pour le pipeline.
- Sélectionnez un catalogue et un schéma dans lesquels écrire les event logs. Si vous disposez des privilèges
USE CATALOGetCREATE SCHEMAsur le catalogue, cliquez surCreate schema dans le menu déroulant pour créer un schéma.
- Cliquez sur Créer un pipeline et continuer .
- Sur la page Source , sélectionnez les tables à ingérer.
- Cliquez sur Enregistrer et continuer .
- Sur la page Destination , sélectionnez un catalogue et un schéma dans lesquels charger les données. Si vous disposez des privilèges
USE CATALOGetCREATE SCHEMAsur le catalogue, cliquez surCréer un schéma dans le menu déroulant pour créer un schéma.
- Cliquez sur Enregistrer et continuer .
- (Facultatif) Sur la page Schedules and notifications , cliquez sur
Create schedule . Définissez la fréquence de refresh des tables de destination.
- (Facultatif) Cliquez sur
Ajouter une notification pour configurer les notifications par e-mail en cas de succès ou d’échec de l’opération du pipeline, puis cliquez sur Enregistrer et exécuter le pipeline .
Utilisez les Declarative Automation Bundles pour gérer les pipelines Celigo en tant que code. Les bundles peuvent contenir des définitions YAML de jobs et de tâches, sont gérés à l'aide de la CLI Databricks et peuvent être partagés et exécutés dans différents workspaces cibles (tels que le développement, la pré-production et la production). Pour plus d'informations, consultez Que sont les Declarative Automation Bundles ?.
-
Créer un bundle à l’aide de la CLI Databricks :
Bashdatabricks bundle init -
Ajoutez deux nouveaux fichiers de ressources au bundle :
- Un fichier de définition de pipeline (par exemple,
resources/celigo_pipeline.yml). Voir pipeline.ingestion_definition et Exemples. - Fichier de définition de job qui contrôle la fréquence de l'ingestion des données (par exemple,
resources/celigo_job.yml).
- Un fichier de définition de pipeline (par exemple,
-
Déployez le pipeline à l'aide de la CLI Databricks :
Bashdatabricks bundle deploy
- Importez le notebook suivant dans votre workspace Databricks :
-
Laissez les cellules un et deux telles quelles. Ne pas modifier.
-
Modifiez la troisième cellule avec les détails de configuration de votre pipeline. Voir pipeline.ingestion_definition et Exemples.
-
Vous pouvez configurer des paramètres avancés du pipeline (facultatif). Consultez la page Modèles courants pour les pipelines d’ingestion gérés.
-
Cliquez sur Tout exécuter .
Exemples
Le connecteur Celigo met à disposition la table source audit_logs dans le schéma source default. Ingérez la table directement ou ingérez le schéma entier.
Ingérer des tables spécifiques
Utilisez cette option pour ingérer un sous-ensemble spécifique de tables ou pour personnaliser la dénomination de la destination par table. Définissez le start du remplissage facultatif start_datetime sur l’objet audit_logs.
- Declarative Automation Bundles
- Databricks notebook
Le fichier de définition de pipeline suivant ingère des tables Celigo individuelles :
resources:
pipelines:
celigo_pipeline:
name: celigo_pipeline
catalog: 'main'
target: 'celigo_data'
ingestion_definition:
connection_name: celigo_connection
objects:
- table:
source_schema: 'default'
source_table: 'audit_logs'
destination_catalog: 'main'
destination_schema: 'celigo_data'
destination_table: 'audit_logs'
connector_options:
api_source_connector_options:
options:
start_datetime: '<start-datetime>'
La spécification de pipeline suivante ingère des tables Celigo individuelles :
pipeline_name = "celigo_pipeline"
connection_name = "<celigo-connection>"
pipeline_spec = {
"name": pipeline_name,
"ingestion_definition": {
"connection_name": connection_name,
"objects": [
{
"table": {
"source_schema": "default",
"source_table": "audit_logs",
"destination_catalog": "main",
"destination_schema": "celigo_data",
"destination_table": "audit_logs",
"connector_options": {
"api_source_connector_options": {
"options": {
"start_datetime": "<start-datetime>"
}
}
}
}
}
]
}
}
json_payload = json.dumps(pipeline_spec, indent=2)
create_pipeline(json_payload)
Ingérer l’ensemble du schéma
Utilisez cette option pour ingérer toutes les tables source Celigo dans un schéma de destination unique avec une seule déclaration.
- Declarative Automation Bundles
- Databricks notebook
Le fichier de définition de pipeline suivant ingère toutes les tables Celigo prises en charge dans un schéma de destination :
resources:
pipelines:
celigo_pipeline:
name: celigo_pipeline
catalog: 'main'
target: 'celigo_data'
ingestion_definition:
connection_name: celigo_connection
objects:
- schema:
source_schema: 'default'
destination_catalog: 'main'
destination_schema: 'celigo_data'
La spécification de pipeline suivante ingère toutes les tables Celigo prises en charge dans un schéma de destination :
pipeline_name = "celigo_pipeline"
connection_name = "<celigo-connection>"
pipeline_spec = {
"name": pipeline_name,
"ingestion_definition": {
"connection_name": connection_name,
"objects": [
{
"schema": {
"source_schema": "default",
"destination_catalog": "main",
"destination_schema": "celigo_data"
}
}
]
}
}
json_payload = json.dumps(pipeline_spec, indent=2)
create_pipeline(json_payload)
Fichier de définition de job Declarative Automation Bundles
Voici un exemple de fichier de définition de job à utiliser avec les Declarative Automation Bundles. Le job s’exécute quotidiennement.
- Declarative Automation Bundles
resources:
jobs:
celigo_job:
name: celigo_job
schedule:
quartz_cron_expression: '0 0 0 * * ?'
timezone_id: 'UTC'
tasks:
- task_key: celigo_ingestion
pipeline_task:
pipeline_id: ${resources.pipelines.celigo_pipeline.id}
Modèles courants
Pour des configurations de pipeline avancées, consultez Common patterns for managed ingestion pipelines.
Étapes suivantes
start, planifiez et configurez des alertes sur votre pipeline. Consultez la section Tâches courantes de maintenance des pipelines.