Aller au contenu principal

Créer un pipeline d’ingestion MySQL

info

Aperçu

Le connecteur MySQL est en préversion publique. Contactez votre équipe de compte Databricks pour demander l'accès.

Découvrez comment importer des données de MySQL dans Databricks à l'aide de Lakeflow Connect. Le connecteur MySQL prend en charge Amazon RDS pour MySQL, Aurora MySQL, Azure Database pour MySQL, Google Cloud SQL pour MySQL et MySQL exécuté sur EC2.

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.

    • 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
      }
      }
      • La politique de compute suivante permet à Databricks de monter en charge la passerelle d'ingestion pour répondre aux besoins de votre workload. La configuration minimale requise est de 4 cœurs. Cependant, pour de meilleures performances d'extraction d'instantanés, Databricks recommande d'utiliser des types d'instances plus grands avec plus de mémoire et de cœurs de CPU.
      Python
      {
      "driver_node_type_id": {
      "type": "fixed",
      "value": "r5n.2xlarge"
      },
      "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 MySQL, vous devez également compléter la configuration de la source.

Option 1 : Interface utilisateur de Databricks

Les utilisateurs administrateurs peuvent créer une connexion et un pipeline en même temps dans l'interface utilisateur. C'est le moyen le plus simple de créer des pipelines d'ingestion gérés.

  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 **MySQL**. L'assistant d'ingestion s'ouvre.

  3. Dans la page Passerelle d'ingestion de l'assistant, saisissez un nom unique pour la passerelle.

  4. Sélectionnez un catalogue et un schéma pour les données d'ingestion de préproduction, puis cliquez sur Suivant .

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

  6. Pour le **Catalogue de destination**, sélectionnez un catalogue pour stocker les données ingérées.

  7. Sélectionnez la connexion Unity Catalog qui stocke les identifiants requis pour accéder aux données source.

    S'il n'y a pas de connexions existantes à la source, cliquez sur Créer une connexion pour créer une nouvelle connexion. Pour obtenir des instructions, consultez Create a MySQL connection. Vous devez disposer de CREATE CONNECTION privilèges sur le métastore.

remarque

Le bouton **Tester la connexion** peut échouer pour les utilisateurs MySQL utilisant sha256_password caching_sha2_password l'authentification ou. Il s'agit d'une limitation connue. Vous pouvez toujours poursuivre la création de la connexion.

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

  2. Sur la page Source , sélectionnez les bases de données et les tables à ingérer.

  3. Modifiez éventuellement le paramètre de suivi de l’historique par default. Pour plus d'informations, consultez Activer le suivi de l'historique (SCD de type 2).

  4. Cliquez sur Suivant .

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

  6. Cliquez sur Enregistrer et continuer .

  7. (Facultatif) Sur la page Paramètres , cliquez sur Créer un planning . Définissez la fréquence de refresh des tables de destination.

  8. (Facultatif) Définissez les notifications par e-mail pour la réussite ou l'échec des opérations de pipeline.

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

Option 2 : Interfaces programmatiques

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

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_mysql_connection"
export TARGET_CATALOG="main"
export TARGET_SCHEMA="lakeflow_mysql_connector"
export STAGING_CATALOG=$TARGET_CATALOG
export STAGING_SCHEMA=$TARGET_SCHEMA
export DB_HOST="mysql-instance.region.rds.amazonaws.com"
export DB_PORT="3306"
export DB_USER="databricks_replication"
export DB_PASSWORD="your_secure_password"

output=$(databricks connections create --json '{
"name": "'"$CONNECTION_NAME"'",
"connection_type": "MYSQL",
"options": {
"host": "'"$DB_HOST"'",
"port": "'"$DB_PORT"'",
"user": "'"$DB_USER"'",
"password": "'"$DB_PASSWORD"'"
}
}')

export CONNECTION_ID=$(echo $output | jq -r '.connection_id')

Créer la passerelle et le pipeline d'ingestion

La passerelle d’ingestion extrait les données d’instantané et de modification de la base de données source et les stocke dans un volume de staging Unity Catalog. Vous devez exécuter la passerelle en tant que pipeline continu. Cela permet de respecter les politiques de rétention des binlogs sur la base de données source.

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

Cet tab décrit comment déployer un pipeline d'ingestion à l'aide des 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/mysql_pipeline.yml).
    • Un fichier de définition de Job qui contrôle la fréquence d'ingestion de données (par exemple, resources/mysql_job.yml).

    Voici un exemple de fichier resources/mysql_pipeline.yml :

    YAML
    variables:
    # Common variables used multiple places in the DAB definition.
    gateway_name:
    default: mysql-gateway
    dest_catalog:
    default: main
    dest_schema:
    default: ingest-destination-schema

    resources:
    pipelines:
    gateway:
    name: ${var.gateway_name}
    gateway_definition:
    connection_name: <mysql-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_mysql:
    name: mysql-ingestion-pipeline
    ingestion_definition:
    ingestion_gateway_id: ${resources.pipelines.gateway.id}
    objects:
    # Modify this with your tables!
    - table:
    # Ingest the table mydb.customers to dest_catalog.dest_schema.customers
    source_schema: public
    source_table: customers
    destination_catalog: ${var.dest_catalog}
    destination_schema: ${var.dest_schema}
    - schema:
    # Ingest all tables in the mydb.sales schema to dest_catalog.dest_schema
    # The destination table name will be the same as it is on the source
    source_schema: sales
    destination_catalog: ${var.dest_catalog}
    destination_schema: ${var.dest_schema}
    catalog: ${var.dest_catalog}
    schema: ${var.dest_schema}

    Voici un exemple de fichier resources/mysql_job.yml :

    YAML
    resources:
    jobs:
    mysql_dab_job:
    name: mysql_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_mysql.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
export GATEWAY_PIPELINE_NAME="mysql-gateway"

output=$(databricks pipelines create --json '{
"name": "'"$GATEWAY_PIPELINE_NAME"'",
"catalog": "'"$STAGING_CATALOG"'",
"schema": "'"$STAGING_SCHEMA"'",
"gateway_definition": {
"connection_id": "'"$CONNECTION_ID"'",
"gateway_storage_catalog": "'"$STAGING_CATALOG"'",
"gateway_storage_schema": "'"$STAGING_SCHEMA"'",
"gateway_storage_name": "'"$GATEWAY_PIPELINE_NAME"'"
}
}')

export GATEWAY_PIPELINE_ID=$(echo $output | jq -r '.pipeline_id')

Pour créer le pipeline d’ingestion :

Bash
export INGESTION_PIPELINE_NAME="mysql-ingestion-pipeline"

databricks pipelines create --json '{
"name": "'"$INGESTION_PIPELINE_NAME"'",
"catalog": "'"$TARGET_CATALOG"'",
"schema": "'"$TARGET_SCHEMA"'",
"ingestion_definition": {
"ingestion_gateway_id": "'"$GATEWAY_PIPELINE_ID"'",
"objects": [
{"table": {
"source_schema": "public",
"source_table": "customers",
"destination_catalog": "'"$TARGET_CATALOG"'",
"destination_schema": "'"$TARGET_SCHEMA"'",
"destination_table": "customers"
}},
{"schema": {
"source_schema": "sales",
"destination_catalog": "'"$TARGET_CATALOG"'",
"destination_schema": "'"$TARGET_SCHEMA"'"
}}
]
}
}'

Terraform

Vous pouvez utiliser Terraform pour déployer et gérer des pipelines d'ingestion MySQL. 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.

Ressources supplémentaires