Aller au contenu principal

Ingérer des données depuis PostgreSQL

info

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 CONNECTION privilè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 CONNECTION ou ALL PRIVILEGES sur la connexion.

    • Vous disposez de USE CATALOG privilèges sur le catalogue cible.

    • Vous disposez de USE SCHEMA, CREATE TABLE et CREATE VOLUME privilèges sur un schéma existant ou de CREATE SCHEMA privilè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

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

  2. Sur la page **Ajouter des données**, sous **Connecteurs Databricks**, cliquez sur **PostgreSQL**.

  3. 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 CONNECTION sur le métastore, vous pouvez cliquer sur Icône Plus. Créer une connexion pour créer une nouvelle connexion avec les détails d'authentification dans Créer une connexion PostgreSQL.

  4. Cliquez sur Suivant .

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

  6. 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 CATALOG et CREATE SCHEMA sur le catalogue, vous pouvez cliquer sur Icône Plus. Créer un schéma dans le menu déroulant pour créer un nouveau schéma.

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

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

  9. 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 CATALOG et CREATE SCHEMA sur le catalogue, vous pouvez cliquer sur Icône Plus. Créer un schéma dans le menu déroulant pour créer un nouveau schéma.

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

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

  12. Cliquez sur Suivant , puis sur Enregistrer et continuer .

  13. 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 CATALOG et CREATE SCHEMA sur le catalogue, vous pouvez cliquer sur Icône Plus. Créer un schéma dans le menu déroulant pour créer un nouveau schéma.

  14. Cliquez sur Enregistrer et continuer .

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

  16. (Facultatif) Sur la page Schedules and notifications , cliquez sur Icône Plus. Créer un calendrier . Définissez la fréquence pour refresh les tables de destination.

  17. (Facultatif) Cliquez sur Icône Plus. 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.

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

  1. Créez un bundle à l'aide de la CLI Databricks :

    Bash
    databricks bundle init
  2. 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 :

    YAML
    variables:
    # 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 :

    YAML
    resources:
    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}
  3. Déployez le pipeline à l'aide de la CLI Databricks :

    Bash
    databricks 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 :

Bash
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 :

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

Vérifier la réplication

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

Ressources supplémentaires