Ingérer des données depuis PostgreSQL
Aperçu
Le connecteur PostgreSQL pour Lakeflow Connect est en aperçu public. Contactez votre équipe de compte Databricks pour vous inscrire à l'aperçu public.
Cette page décrit comment importer des données de PostgreSQL et les charger dans Databricks à l'aide de LakeFlow Connect. Le connecteur PostgreSQL prend en charge AWS RDS PostgreSQL, Aurora PostgreSQL, Amazon EC2, Azure Database for PostgreSQL, les machines virtuelles Azure, GCP Cloud SQL for PostgreSQL et les bases de données PostgreSQL on-premise utilisant Azure ExpressRoute, AWS Direct Connect ou le réseau VPN.
Avant de commencer
-
Pour créer une passerelle d'ingestion et un pipeline d'ingestion, vous devez remplir les conditions suivantes :
-
Votre workspace est activé pour Unity Catalog.
-
Le compute Serverless est activé pour votre workspace. Consultez les exigences du compute Serverless.
-
Si vous prévoyez de créer une connexion : vous avez
CREATE CONNECTIONprivilèges sur le metastore. Voir Gérer les privilèges dans Unity Catalog.Si votre connecteur prend en charge la création de pipelines basée sur l'interface utilisateur, vous pouvez créer la connexion et le pipeline simultanément en suivant les étapes de cette page. Cependant, si vous utilisez l’édition de pipelines basée sur l’API, vous devez créer la connexion dans l’Explorateur de catalogues avant de suivre les étapes de cette page. Voir Connexion aux sources d'ingestion gérées.
-
Si vous prévoyez d'utiliser une connexion existante : vous disposez des privilèges
USE CONNECTIONouALL PRIVILEGESsur la connexion. -
Vous disposez de
USE CATALOGprivilèges sur le catalogue cible. -
Vous disposez de
USE SCHEMA,CREATE TABLEetCREATE VOLUMEprivilèges sur un schéma existant ou deCREATE SCHEMAprivilèges sur le catalogue cible. -
Vous avez accès à une instance PostgreSQL principale. La réplication logique est uniquement prise en charge sur les instances primaires et non sur les réplicas en lecture.
-
Autorisations illimitées pour créer des clusters, ou une politique personnalisée (API seulement). Une politique personnalisée pour la passerelle doit répondre aux exigences suivantes :
-
Famille : Job compute
-
Dérogations de la famille de règles :
{
"cluster_type": {
"type": "fixed",
"value": "dlt"
},
"num_workers": {
"type": "unlimited",
"defaultValue": 1,
"isOptional": true
},
"runtime_engine": {
"type": "fixed",
"value": "STANDARD",
"hidden": true
}
}- Databricks recommande de spécifier les plus petits nœuds Worker possibles pour les passerelles d'ingestion car ils n'affectent pas les performances de la passerelle. La politique de compute suivante permet à Databricks de faire évoluer la passerelle d'ingestion pour répondre aux besoins de votre charge de travail. L'exigence minimale est de 8 cœurs pour permettre une extraction de données efficace et performante à partir de votre base de données source.
Python{
"driver_node_type_id": {
"type": "fixed",
"value": "r5n.16xlarge"
},
"node_type_id": {
"type": "fixed",
"value": "m5n.large"
}
}Pour plus d'information sur les politiques de clusters, consultez Sélectionner une politique de compute.
-
-
-
Pour ingérer depuis PostgreSQL, vous devez également effectuer la configuration de la source.
Créer une passerelle et un pipeline d'ingestion
Interface utilisateur Databricks
-
Dans la barre latérale du workspace Databricks, cliquez sur Ingestion de données .
-
Sur la page **Ajouter des données**, sous **Connecteurs Databricks**, cliquez sur **PostgreSQL**.
-
Sur la page Connexion de l'assistant d'ingestion, sélectionnez la connexion qui stocke vos identifiants d'accès PostgreSQL. Si vous avez le privilège
CREATE CONNECTIONsur le métastore, vous pouvez cliquer surCréer une connexion pour créer une nouvelle connexion avec les détails d'authentification dans Créer une connexion PostgreSQL.
-
Cliquez sur Suivant .
-
Sur la page Configuration de l’ingestion , saisissez un nom unique pour le pipeline d'ingestion. Ce pipeline transfère les données de l'emplacement intermédiaire vers la destination.
-
Sélectionnez un catalogue et un schéma pour écrire les Logs d'événements. Le journal des événements contient les Logs d'audit, les contrôles de qualité des données, la progression du pipeline et les erreurs. Si vous disposez des privilèges
USE CATALOGetCREATE SCHEMAsur le catalogue, vous pouvez cliquer surCréer un schéma dans le menu déroulant pour créer un nouveau schéma.
-
(Facultatif) Définissez **Refresh automatique de toutes les tables** sur **Activé**. Lorsque le refresh automatique est activé, le pipeline tente automatiquement de résoudre les problèmes tels que les événements de nettoyage des logs et certains types d'évolution des schémas en effectuant un refresh complet de la table impactée. Si le suivi de l'historique est activé, un refresh complet efface cet historique.
-
Saisissez un nom unique pour la passerelle d'ingestion. La passerelle est un pipeline qui extrait les modifications de la source et les prépare pour que le pipeline d'ingestion les charge.
-
Sélectionnez un catalogue et un schéma pour l’ emplacement de préproduction . Un volume est créé à cet emplacement pour extraire des données. Si vous disposez des privilèges
USE CATALOGetCREATE SCHEMAsur le catalogue, vous pouvez cliquer surCréer un schéma dans le menu déroulant pour créer un nouveau schéma.
-
Cliquez sur Créer un pipeline et continuer .
-
Sur la page Source , sélectionnez les tables à ingérer. Si vous sélectionnez des tables spécifiques, vous pouvez configurer les paramètres de la table :
a. (Facultatif) Dans l'onglet Paramètres , spécifiez un Nom de destination pour chaque table ingérée. Ceci est utile pour différencier les tables de destination lorsque vous ingérez un objet dans le même schéma plusieurs fois. Consultez Nommer une table de destination.
a. (Facultatif) Modifiez le paramètre default du suivi de l’historique . Voir Activer le suivi de l'historique (SCD de type 2).
-
Cliquez sur Suivant , puis sur Enregistrer et continuer .
-
Sur la page Destination , sélectionnez un catalogue et un schéma dans lesquels charger des données. Si vous disposez des privilèges
USE CATALOGetCREATE SCHEMAsur le catalogue, vous pouvez cliquer surCréer un schéma dans le menu déroulant pour créer un nouveau schéma.
-
Cliquez sur Enregistrer et continuer .
-
Sur la page Configuration de la base de données , saisissez le nom de l'emplacement de réplication et le nom de la publication pour chaque base de données à partir de laquelle vous souhaitez ingérer des données.
-
(Facultatif) Sur la page Schedules and notifications , cliquez sur
Créer un calendrier . Définissez la fréquence pour refresh les tables de destination.
-
(Facultatif) Cliquez sur
Ajouter une notification pour définir 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 .
Avant d'ingérer des données en utilisant les Bundles d'automatisation déclarative, les APIs Databricks, les SDK Databricks, la CLI Databricks ou Terraform, vous devez avoir accès à une connexion Unity Catalog existante. Pour obtenir des instructions, consultez Créer une connexion PostgreSQL.
Créer le catalogue et le schéma de staging
Le catalogue et le schéma de préproduction peuvent être les mêmes que le catalogue et le schéma de destination. Le catalogue de préproduction ne peut pas être un catalogue étranger.
- CLI
export CONNECTION_NAME="my_postgresql_connection"
export TARGET_CATALOG="main"
export TARGET_SCHEMA="lakeflow_postgresql_connector_cdc"
export STAGING_CATALOG=$TARGET_CATALOG
export STAGING_SCHEMA=$TARGET_SCHEMA
export DB_HOST="postgresql-instance.example.com"
export DB_PORT="5432"
export DB_DATABASE="your_database"
export DB_USER="databricks_replication"
export DB_PASSWORD="your_secure_password"
output=$(databricks connections create --json '{
"name": "'"$CONNECTION_NAME"'",
"connection_type": "POSTGRESQL",
"options": {
"host": "'"$DB_HOST"'",
"port": "'"$DB_PORT"'",
"database": "'"$DB_DATABASE"'",
"user": "'"$DB_USER"'",
"password": "'"$DB_PASSWORD"'"
}
}')
export CONNECTION_ID=$(echo $output | jq -r '.connection_id')
La passerelle d'ingestion extrait les données d'instantané et de changement de la base de données source et les stocke dans le volume de staging d'Unity Catalog. Vous devez exécuter la passerelle en tant que pipeline continu. Ceci est essentiel pour PostgreSQL afin d'éviter l'engorgement du journal d'écriture anticipée (WAL) et de garantir que les slots de réplication n'accumulent pas de modifications non consommées.
Le pipeline d'ingestion applique les données d'instantané et de modification du volume de staging aux tables de streaming de destination.
Declarative Automation Bundles
Vous pouvez déployer un pipeline d'ingestion à l'aide de Declarative Automation Bundles. Les bundles peuvent contenir des définitions YAML de jobs et de tâches, sont gérés à l’aide de la Databricks CLI, et peuvent être partagés et exécutés dans différents Workspace cibles (comme le développement, la pré-production et la production). Pour plus d'informations, consultez Declarative Automation Bundles.
-
Créez 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/postgresql_pipeline.yml). - Un fichier de définition de Job qui contrôle la fréquence d'ingestion de données (par exemple,
resources/postgresql_job.yml).
Voici un exemple de fichier
resources/postgresql_pipeline.yml:YAMLvariables:
# Common variables used multiple places in the DAB definition.
gateway_name:
default: postgresql-gateway
dest_catalog:
default: main
dest_schema:
default: ingest-destination-schema
resources:
pipelines:
gateway:
name: ${var.gateway_name}
gateway_definition:
connection_name: <postgresql-connection>
gateway_storage_catalog: main
gateway_storage_schema: ${var.dest_schema}
gateway_storage_name: ${var.gateway_name}
catalog: ${var.dest_catalog}
schema: ${var.dest_schema}
pipeline_postgresql:
name: postgresql-ingestion-pipeline
ingestion_definition:
ingestion_gateway_id: ${resources.pipelines.gateway.id}
source_type: POSTGRESQL
objects:
# Modify this with your tables!
- table:
# Ingest the table public.orders to dest_catalog.dest_schema.orders.
source_catalog: your_database
source_schema: public
source_table: orders
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
- schema:
# Ingest all tables in the public schema to dest_catalog.dest_schema. The destination
# table name will be the same as it is on the source.
source_catalog: your_database
source_schema: public
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
source_configurations:
- catalog:
source_catalog: your_database
postgres:
slot_config:
slot_name: databricks_slot
publication_name: databricks_publication
catalog: ${var.dest_catalog}
schema: ${var.dest_schema}Voici un exemple de fichier
resources/postgresql_job.yml:YAMLresources:
jobs:
postgresql_dab_job:
name: postgresql_dab_job
trigger:
# Run this job every day, exactly one day from the last run
# See https://docs.databricks.com/api/workspace/jobs/create#trigger
periodic:
interval: 1
unit: DAYS
email_notifications:
on_failure:
- <email-address>
tasks:
- task_key: refresh_pipeline
pipeline_task:
pipeline_id: ${resources.pipelines.pipeline_postgresql.id} - Un fichier de définition de pipeline (par exemple,
-
Déployez le pipeline à l'aide de la CLI Databricks :
Bashdatabricks bundle deploy
Notebook Databricks
Mettez à jour la cellule Configuration du Notebook suivant avec la connexion source, le catalogue cible, le schéma cible et les tables à ingérer depuis la source.
Create gateway and ingestion pipeline
Databricks CLI
Pour créer la passerelle :
gateway_json=$(cat <<EOF
{
"name": "$GATEWAY_PIPELINE_NAME",
"catalog": "$STAGING_CATALOG",
"schema": "$STAGING_SCHEMA",
"gateway_definition": {
"connection_name": "$CONNECTION_NAME",
"gateway_storage_catalog": "$STAGING_CATALOG",
"gateway_storage_schema": "$STAGING_SCHEMA",
"gateway_storage_name": "$GATEWAY_PIPELINE_NAME"
}
}
EOF
)
output=$(databricks pipelines create --json "$gateway_json")
echo $output
export GATEWAY_PIPELINE_ID=$(echo $output | jq -r '.pipeline_id')
Pour créer le pipeline d’ingestion :
pipeline_json=$(cat <<EOF
{
"name": "$INGESTION_PIPELINE_NAME",
"catalog": "$TARGET_CATALOG",
"schema": "$TARGET_SCHEMA",
"ingestion_definition": {
"ingestion_gateway_id": "$GATEWAY_PIPELINE_ID",
"source_type": "POSTGRESQL",
"objects": [
{
"table": {
"source_catalog": "your_database",
"source_schema": "public",
"source_table": "orders",
"destination_catalog": "$TARGET_CATALOG",
"destination_schema": "$TARGET_SCHEMA",
"destination_table": "orders"
}
},
{
"schema": {
"source_catalog": "your_database",
"source_schema": "public",
"destination_catalog": "$TARGET_CATALOG",
"destination_schema": "$TARGET_SCHEMA"
}
}
],
"source_configurations": [
{
"catalog": {
"source_catalog": "your_database",
"postgres": {
"slot_config": {
"slot_name": "databricks_slot",
"publication_name": "databricks_publication"
}
}
}
}
]
}
}
EOF
)
databricks pipelines create --json "$pipeline_json"
Nécessite Databricks CLI v0.276.0 ou ultérieure.
Terraform
Vous pouvez utiliser Terraform pour déployer et gérer des pipelines d'ingestion PostgreSQL. Pour un framework d'exemple complet, y compris les configurations Terraform pour la création de passerelles et de pipelines d'ingestion, consultez le repository Lakeflow Connect Terraform examples sur GitHub.
Start, planifiez et définissez des alertes sur votre pipeline
Pour obtenir des informations sur le démarrage, la planification et la configuration des alertes de votre pipeline, consultez Tâches courantes de maintenance de pipeline.
Vérifier l'ingestion des données
La vue liste sur la page des détails du pipeline affiche le nombre d'enregistrements traités au fur et à mesure que les données sont ingérées. Ces nombres refresh automatiquement.

Les colonnes Upserted records et Deleted records ne sont pas affichées par default. Vous pouvez les activer en cliquant sur le bouton de configuration des colonnes et en les sélectionnant.