Créer un pipeline CDC intégré pour MySQL
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 . Consultez Gérer les aperçus Databricks.
Un pipeline CDC intégré ingère les données de modification de MySQL dans Databricks à l'aide d'un seul pipeline. Contrairement à l'architecture standard basée sur une passerelle, un pipeline CDC intégré exécute les étapes d'extraction et d'application dans une seule mise à jour de pipeline. L'architecture standard nécessite une passerelle d'ingestion et un pipeline d'ingestion séparés.
Quand utiliser le connecteur CDC intégré
Choisissez le pipeline CDC intégré lorsque :
- Vous souhaitez une configuration plus simple avec un pipeline au lieu d'une passerelle d'ingestion et d'un pipeline d'ingestion distincts.
- L'exécution déclenchée (planifiée) répond à vos besoins. Les pipelines CDC intégrés s'exécutent selon un planning ; l'exécution continue (toujours active) n'est pas prise en charge.
- Vous avez besoin d'un support de refresh complet automatique, qui n'est pas disponible pour les flux existants basés sur des passerelles MySQL.
Le tableau suivant compare les deux architectures en détail :
Fonctionnalité | CDC standard (basée sur une passerelle) | CDC intégré |
|---|---|---|
Nombre de pipelines | Deux (passerelle d'ingestion et pipeline d'ingestion) | Un (pipeline unifié) |
Installer | Créez une passerelle, puis créez un pipeline d'ingestion qui référence l'ID de la passerelle | Créez un pipeline unique qui référence une connexion Unity Catalog |
Mode passerelle | La passerelle s'exécute en continu en tant que processus distinct et de longue durée | L'extraction est intégrée à chaque mise à jour planifiée du pipeline. |
Référence de connexion |
|
|
Type de connecteur | Comportement CDC par default implicite | Explicite : |
Volume intermédiaire | Géré en interne par la passerelle | Créé automatiquement dans le schéma de destination, ou configuré via |
Mode de pipeline | Continu | Déclenché uniquement |
Calculer | Classique pour la passerelle, serverless pour le pipeline d'ingestion géré | Classic compute uniquement. Serverless n'est pas pris en charge. |
Refresh complet automatique | Non pris en charge pour les flux existants basés sur une passerelle MySQL | Pris en charge |
Tables maximales | 250 par pipeline | 250 par pipeline |
Pour la configuration de la base de données source, consultez Configurer MySQL pour l'ingestion dans Databricks. La même configuration source s'applique aux deux architectures.
Comment s'exécute un pipeline CDC intégré
Chaque mise à jour de pipeline exécute deux étapes en séquence :
- Extraction. Le pipeline se connecte à la base de données source via la connexion Unity Catalog. Lors de la première exécution ou d'un refresh complet, il capture un instantané initial. Lors des exécutions ultérieures, il capture les modifications incrémentielles (insertions, mises à jour et suppressions) à l'aide du journal binaire (binlog). Le pipeline écrit les données extraites vers un volume de staging Unity Catalog.
- Application. Le pipeline lit le volume de staging et applique les modifications aux tables de streaming de destination dans Unity Catalog. Les opérations Merge utilisent les clés primaires configurées et le type SCD. Le pipeline garantit une sémantique exactement une fois.
Chaque mise à jour du pipeline extrait les modifications, puis s'arrête automatiquement après avoir rattrapé la source, dans la limite d'une durée d'exécution maximale. Pour plus de détails, consultez Fermeture intelligente pour les pipelines CDC intégrés. Pour ingérer des données de manière récurrente, planifiez le pipeline à l'aide d'une tâche Lakeflow Jobs.
Exigences
-
Votre workspace est activé pour Unity Catalog.
-
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. -
Votre Workspace doit avoir la fonctionnalité de connecteur CDC intégré activée. Contactez l'équipe de votre compte Databricks.
-
Vous avez terminé la configuration de la source MySQL. Consultez Configurer MySQL pour l'ingestion dans Databricks.
-
Vous disposez des autorisations suivantes :
CREATE CONNECTIONsur le metastore (si vous créez une nouvelle connexion Unity Catalog), ouUSE CONNECTIONsur une connexion existante.USE CATALOGsur le catalogue de destination.USE SCHEMAetCREATE TABLEsur le schéma de destination.CREATE VOLUMEsur le schéma de destination, ou sur le schéma spécifié dansdata_staging_options. Un volume de staging est requis même sidata_staging_optionsn’est pas défini, car le pipeline en crée un automatiquement dans le schéma de destination.
Exigences en matière de compute
Les pipelines CDC intégrés pour MySQL nécessitent du compute classique. Le compute Serverless n'est pas pris en charge.
- Compute classique : Le plan de compute classique s'exécute dans le Virtual Private Cloud (VPC) ou VNet de votre Workspace Databricks et doit atteindre votre instance MySQL via le réseau. Les chemins réseau pris en charge incluent le peering de Virtual Private Cloud (VPC) ou VNet, les endpoints publics et, pour MySQL on-premise, AWS Direct Connect, Azure ExpressRoute ou VPN.
Pour le compute classique, utilisez des autorisations de création de cluster illimitées ou une politique de cluster personnalisée avec cluster_type défini sur dlt et runtime_engine défini sur STANDARD. Databricks recommande au moins 8 cœurs pour une extraction efficace.
Créez une connexion Unity Catalog à MySQL
Créez une connexion Unity Catalog à MySQL avant de créer un pipeline. Voir Créer une connexion MySQL.
Créer un pipeline CDC intégré
Créez des pipelines CDC intégrés à l'aide de l'API, de la CLI Databricks, des Notebooks ou des Declarative Automation Bundles. La création d'interface utilisateur n'est pas encore disponible.
Toutes les demandes de création de pipeline doivent inclure "channel": "PREVIEW".
- Declarative Automation Bundles
- Databricks notebook
- Databricks CLI
- REST API
Définissez la ressource de pipeline dans un fichier de bundle (par exemple, resources/integrated_cdc_pipeline.yml) :
variables:
pipeline_name:
description: 'Name for the integrated CDC pipeline'
connection_name:
description: 'Unity Catalog connection name'
dest_catalog:
description: 'Destination catalog for ingested data'
dest_schema:
description: 'Destination schema for ingested data'
resources:
pipelines:
integrated_cdc_pipeline:
name: ${var.pipeline_name}
channel: PREVIEW
catalog: ${var.dest_catalog}
schema: ${var.dest_schema}
ingestion_definition:
connection_name: ${var.connection_name}
connector_type: CDC
objects:
- table:
source_schema: 'my_database'
source_table: 'customers'
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
destination_table: 'customers'
table_configuration:
scd_type: 'SCD_TYPE_1'
Pour exécuter le pipeline selon un calendrier, définissez un Job (par exemple, resources/integrated_cdc_job.yml) qui déclenche le pipeline. Étant donné que chaque étape d'extraction dure au moins 10 minutes, un intervalle de 60 minutes ou plus est un bon point de départ :
resources:
jobs:
integrated_cdc_job:
name: '${var.pipeline_name}-job'
tasks:
- task_key: 'cdc_ingestion'
pipeline_task:
pipeline_id: ${resources.pipelines.integrated_cdc_pipeline.id}
schedule:
quartz_cron_expression: '0 0 * * * ?'
timezone_id: 'UTC'
Déployez le bundle avec la CLI Databricks :
databricks bundle deploy
databricks bundle run integrated_cdc_job
Pour plus d'information, consultez Que sont les bundles d'automatisation déclaratifs ?.
from databricks.sdk import WorkspaceClient
from databricks.sdk.service.pipelines import (
ConnectorType,
IngestionConfig,
IngestionPipelineDefinition,
TableSpec,
)
w = WorkspaceClient()
pipeline = w.pipelines.create(
name="<pipeline-name>",
channel="PREVIEW",
catalog="<destination-catalog>",
schema="<destination-schema>",
ingestion_definition=IngestionPipelineDefinition(
connection_name="<unity-catalog-connection-name>",
connector_type=ConnectorType.CDC,
objects=[
IngestionConfig(
table=TableSpec(
source_schema="<source-database>",
source_table="<source-table>",
destination_catalog="<destination-catalog>",
destination_schema="<destination-schema>",
)
)
],
),
)
print(f"Pipeline created: {pipeline.pipeline_id}")
databricks pipelines create --json '{
"name": "<pipeline-name>",
"channel": "PREVIEW",
"catalog": "<destination-catalog>",
"schema": "<destination-schema>",
"ingestion_definition": {
"connection_name": "<unity-catalog-connection-name>",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_schema": "<source-database>",
"source_table": "<source-table>"
}
}
]
}
}'
L'exemple suivant réplique deux tables d'une base de données MySQL. Les deux héritent de la destination de niveau supérieur main.ingestion. Vous pouvez omettre serverless car il est default sur false, et les pipelines CDC intégrés MySQL ne s'exécutent que sur du compute classique. Le compute Serverless n'est pas pris en charge.
POST /api/2.0/pipelines
{
"name": "my-integrated-cdc-pipeline",
"channel": "PREVIEW",
"catalog": "main",
"schema": "ingestion",
"ingestion_definition": {
"connection_name": "my-mysql-connection",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_schema": "my_database",
"source_table": "customers",
"table_configuration": {
"scd_type": "SCD_TYPE_1"
}
}
},
{
"table": {
"source_schema": "my_database",
"source_table": "orders",
"table_configuration": {
"scd_type": "SCD_TYPE_1"
}
}
}
],
"data_staging_options": {
"catalog_name": "main",
"schema_name": "ingestion_staging"
}
}
}
Pour répliquer chaque table dans une base de données source, utilisez un objet schema au lieu d'objets table individuels :
POST /api/2.0/pipelines
{
"name": "my-integrated-cdc-schema-pipeline",
"channel": "PREVIEW",
"catalog": "main",
"schema": "ingestion",
"ingestion_definition": {
"connection_name": "my-mysql-connection",
"connector_type": "CDC",
"objects": [
{
"schema": {
"source_schema": "my_database",
"destination_catalog": "main",
"destination_schema": "ingestion"
}
}
]
}
}
Pour start une mise à jour de pipeline :
POST /api/2.0/pipelines/<pipeline-id>/updates
{
"full_refresh": false
}
Planifier les mises à jour récurrentes
Les pipelines CDC intégrés s'exécutent uniquement en mode Trigger. Pour ingérer des données selon un calendrier récurrent, créez une tâche Lakeflow Jobs qui exécute le pipeline. La durée de la mise à jour varie en fonction de la quantité de données modifiées dont dispose la source, et un important arriéré pourrait ne pas se terminer en une seule mise à jour (voir Clôture intelligente pour les pipelines CDC intégrés). Programmez les pipelines suffisamment souvent pour que les mises à jour suivantes puissent suivre le rythme. Un point de départ de 60 minutes convient bien à la plupart des charges de travail. Si un Trigger se déclenche alors qu'une mise à jour précédente est toujours en cours d'exécution, la nouvelle mise à jour est mise en file d'attente.
Référence de configuration
Paramètres du pipeline
parameter | Type | Description |
|---|---|---|
| chaîne | Un nom pour le pipeline. |
| chaîne | Doit être |
| Booléen | Facultatif. La valeur par défaut est |
| chaîne | Le catalogue de destination default. Utilisé lorsqu'un |
| chaîne | Le schéma de destination default. Utilisé lorsqu'un |
| chaîne | La connexion Unity Catalog à la base de données source. |
| chaîne | Doit être |
| tableau | La liste des tables ou schémas à ingérer. |
| objet | Facultatif. Le catalogue et le schéma où le pipeline crée le volume de préparation. Default au schéma de destination du pipeline. |
Spécification de table
parameter | Obligatoire | Description |
|---|---|---|
| Oui | Le nom de la base de données source MySQL. |
| Oui | Le nom de la table source. |
| Non | Le catalogue de destination. Valeur par défaut : |
| Non | Le schéma de destination. Valeur par défaut : |
| Non | Le nom de la table de destination. La valeur default est |
Configuration de la table
parameter | Par défaut | Description |
|---|---|---|
|
|
|
| Détection automatique | Les colonnes utilisées pour trier les événements CDC. Détecté automatiquement en fonction du mécanisme CDC source s'il n'est pas spécifié. |
| Désactivé | Configure le refresh complet automatique lorsque des Opérations DDL non prises en charge sont détectées. Voir Politique de refresh complet automatique. |
Pour les mappages de types de données MySQL, consultez la référence du connecteur MySQL. Les pipelines CDC intégrés prennent en charge l'élargissement automatique des types : lorsqu'un type de colonne source est élargi (par exemple, INT à BIGINT), la table de destination s'adapte automatiquement.
Surveiller le pipeline
Après avoir créé et start un pipeline CDC intégré, surveillez son statut à l'aide des éléments suivants :
-
Interface utilisateur Databricks. Ouvrez le pipeline dans la section Pipelines pour afficher l'état de la mise à jour, les métriques d'ingestion par table et la lignée.
-
API REST.
TextGET /api/2.0/pipelines/<pipeline-id> -
API Événements.
TextGET /api/2.0/pipelines/<pipeline-id>/events
La première mise à jour du pipeline effectue un instantané complet de toutes les tables sélectionnées. Contrairement aux mises à jour incrémentielles, l'instantané initial se termine en une seule mise à jour. L'achèvement de l'instantané peut prendre plus de temps que les mises à jour incrémentielles ultérieures.
Pour vérifier l'ingestion :
-- Check row counts in the destination table
SELECT COUNT(*) FROM <destination_catalog>.<destination_schema>.<destination_table>;
-- View recent changes (SCD Type 1 tables)
SELECT * FROM <destination_catalog>.<destination_schema>.<destination_table>
ORDER BY __START_AT DESC
LIMIT 10;
Pour le full refresh et le comportement d'auto full refresh, voir refresh des tables cibles.
Les pipelines CDC intégrés activent le dimensionnement automatique vertical par default. Si une mise à jour de pipeline échoue en raison d'un manque de mémoire, la prochaine mise à jour provisionne automatiquement un Driver plus grand. Pour remplacer ce comportement, utilisez une règle de cluster personnalisée.
Limitations
- **Bêta.** Le connecteur CDC intégré nécessite une activation au niveau du workspace. Veuillez contacter l'équipe de votre compte Databricks.
- Mode Trigger uniquement. Les pipelines CDC intégrés ne prennent pas en charge l'exécution continue (toujours active). Planifiez les pipelines à l'aide d'une tâche Lakeflow Jobs.
- Création via API uniquement. La création de pipeline est disponible via l'API REST, le CLI Databricks, les Notebook et les Declarative Automation Bundles. La création d'interface utilisateur n'est pas encore prise en charge.
- Le Canal de distribution doit être
PREVIEW. Les spécifications de pipeline doivent inclure"channel": "PREVIEW". - Le type de connexion et de connecteur sont immuables.
connection_nameetconnector_typene peuvent pas être modifiés après la création du pipeline. Pour modifier la source, créez un nouveau pipeline. - Maximum de 250 tables par pipeline.
- **Tables sans clés primaires.** Le pipeline traite toutes les colonnes non LOB comme une clé composite. Les lignes en double peuvent être regroupées en une seule ligne.
- L'instantané initial se termine en une seule mise à jour. Le connecteur CDC intégré effectue l'instantané initial en une seule mise à jour du pipeline, même pour les grandes tables.
- La mise à jour du runtime est gérée automatiquement. La fermeture intelligente détermine le moment où chaque mise à jour s'arrête. Une mise à jour se termine après avoir rattrapé la source, limitée par un runtime maximum. Voir Fermeture intelligente pour les pipelines CDC intégrés. Vous ne pouvez pas configurer le runtime minimum ou maximum. Un important backlog de modifications peut s'étendre sur plusieurs mises à jour. Les mises à jour planifiées ultérieures reprennent là où la mise à jour précédente s'est arrêtée.
- La purge du journal binaire nécessite un full refresh. Si le log binaire MySQL est purgé avant que le pipeline ne traite les modifications, effectuez une full refresh sur les tables affectées. Le pipeline détecte cette condition et affiche une erreur dans le log des événements.
- Le compute Serverless n'est pas pris en charge. Les pipelines de CDC intégrés à MySQL nécessitent un compute classique.
Dépannage
Certains codes d'erreur utilisent le préfixe INGESTION_GATEWAY_. Ceci est une convention de nommage héritée et n'indique pas qu'une passerelle d'ingestion distincte est requise.
Erreur | Cause | Résolution |
|---|---|---|
| Le pipeline n'est pas en Mode de publication directe. | Le Mode de publication directe est défini automatiquement pour les pipelines CDC intégrés. Si vous voyez cette erreur, recréez le pipeline. |
| La journalisation binaire n’est pas activée ou | Activez la journalisation binaire avec |
| La table source spécifiée n'existe pas ou a été supprimée. | Vérifiez que la table existe et que l'utilisateur de la connexion a accès. |
| Le schéma source n'existe pas. | Vérifiez que le schéma existe dans la base de données source. |
| Le type de base de données source n'est pas pris en charge. | Le connecteur CDC intégré prend en charge MySQL, SQL Server et Oracle. |
| Il manque | Ajoutez |
| L’indicateur de fonctionnalité Workspace n’est pas activé. | Contactez l'équipe de votre compte Databricks pour activer le connecteur CDC intégré sur votre Workspace. |
Si vous rencontrez un problème non abordé ici :
- Examinez le journal des événements du pipeline dans l'interface utilisateur de Databricks ou via
GET /api/2.0/pipelines/<pipeline-id>/events. - Testez la connexion Unity Catalog à partir de l’Explorateur de catalogues pour confirmer que la source est accessible.
- Confirmez que la journalisation binaire est activée sur la base de données source avec
binlog_format=ROWetbinlog_row_image=FULL. - Vérifiez que l'utilisateur de la base de données dispose des permissions MySQL listées dans Accorder les privilèges d'utilisateur MySQL.
- Vérifiez que votre spécification de pipeline inclut
"channel": "PREVIEW".