Créer un pipeline d’ingestion MySQL
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 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. -
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.
-
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 **MySQL**. L'assistant d'ingestion s'ouvre.
-
Dans la page Passerelle d'ingestion de l'assistant, saisissez un nom unique pour la passerelle.
-
Sélectionnez un catalogue et un schéma pour les données d'ingestion de préproduction, puis cliquez sur Suivant .
-
Sur la page Pipeline d'ingestion , saisissez un nom unique pour le pipeline.
-
Pour le **Catalogue de destination**, sélectionnez un catalogue pour stocker les données ingérées.
-
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 CONNECTIONprivilèges sur le métastore.
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.
-
Cliquez sur Créer un pipeline et continuer .
-
Sur la page Source , sélectionnez les bases de données et les tables à ingérer.
-
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).
-
Cliquez sur Suivant .
-
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 CATALOGetCREATE SCHEMAsur le catalogue parent. -
Cliquez sur Enregistrer et continuer .
-
(Facultatif) Sur la page Paramètres , cliquez sur Créer un planning . Définissez la fréquence de refresh des tables de destination.
-
(Facultatif) Définissez les notifications par e-mail pour la réussite ou l'échec des opérations de pipeline.
-
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.
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.
-
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/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:YAMLvariables:
# 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:YAMLresources:
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} - 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 :
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 :
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.