Aller au contenu principal

Créez un pipeline CDC intégré pour SQL Server

info

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 de CDC intégré ingère les données modifiées de SQL Server vers Databricks à l'aide d'un seul pipeline. Contrairement à l'architecture standard basée sur des passerelles, qui nécessite une passerelle d'ingestion et un pipeline d'ingestion distincts, un pipeline CDC intégré exécute les étapes d'extraction et d'application dans une seule mise à jour de pipeline.

Quand utiliser le connecteur CDC intégré

Le tableau suivant compare les pipelines CDC intégrés avec l'architecture standard basée sur les passerelles :

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

Le pipeline intègre l'extraction dans chaque mise à jour.

Référence de connexion

ingestion_gateway_id

connection_name (une connexion Unity Catalog)

Type de connecteur

Implicite

Explicite : connector_type: CDC

Volume intermédiaire

La passerelle gère le volume intermédiaire en interne.

Vous configurez le volume de staging via data_staging_options. Le pipeline en crée un automatiquement s'il n'est pas spécifié.

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

Le pipeline intègre l'extraction dans chaque mise à jour.

Référence de connexion

ingestion_gateway_id

connection_name (une connexion Unity Catalog)

Type de connecteur

Implicite

Explicite : connector_type: CDC

Volume intermédiaire

La passerelle gère le volume intermédiaire en interne.

Vous configurez le volume de staging via data_staging_options. Le pipeline en crée un automatiquement s'il n'est pas spécifié.

Pour la configuration de la base de données source, consultez Configuration de Microsoft SQL Server 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 :

  1. 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 suivantes, il capture les modifications incrémentielles (insertions, mises à jour et suppressions) à l’aide du mécanisme de suivi des modifications intégré de la base de données. Le pipeline écrit les données extraites vers un volume de staging Unity Catalog.
  2. 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 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.

  • Votre Workspace doit avoir la fonctionnalité de connecteur CDC intégré activée. Contactez l'équipe de votre compte Databricks.

  • Vous avez accès à l’instance principale de SQL Server. Le connecteur CDC intégré ne prend pas en charge les réplicas en lecture, les instances de secours ou les instances secondaires.

  • Vous avez terminé la configuration de la source SQL Server. Consultez Configurer Microsoft SQL Server pour l'ingestion dans Databricks.

  • Vous disposez des autorisations suivantes :

    • CREATE CONNECTION sur le metastore (si vous créez une nouvelle connexion Unity Catalog), ou USE CONNECTION sur une connexion existante.
    • USE CATALOG sur le catalogue de destination.
    • USE SCHEMA et CREATE TABLE sur le schéma de destination.
    • CREATE VOLUME sur le schéma de destination, ou sur le schéma spécifié dans data_staging_options. Un volume de staging est requis même si data_staging_options n’est pas défini, car le pipeline en crée un automatiquement dans le schéma de destination.

Exigences en matière de compute

Un pipeline CDC intégré s'exécute sur un compute classique ou serverless :

  • Compute classique : Le plan de compute classique s'exécute dans votre Virtual Private Cloud (VPC) ou VNet du Workspace Databricks et doit pouvoir atteindre votre instance SQL Server sur le réseau. Tout chemin réseau qui permet au plan de compute d'atteindre la base de données est pris en charge, y compris le peering VPC ou VNet, les endpoints publics et, pour SQL Server on-premise, AWS Direct Connect, Azure ExpressRoute ou VPN.
  • Compute serverless : configurez la connectivité réseau serverless entre le compute serverless de Databricks et votre base de données source. Les sources on-premise nécessitent un chemin réseau via la sortie Serverless configurée (par exemple, une passerelle de transit ou un VNet appairé avec ExpressRoute ou VPN).

Pour le compute classique, vous pouvez utiliser des autorisations de création de cluster illimitées ou une stratégie de cluster personnalisée avec cluster_type défini sur dlt, runtime_engine défini sur STANDARD, et au moins 8 cœurs recommandés pour une extraction efficace.

Créez une connexion Unity Catalog à SQL Server

Créez une connexion Unity Catalog à SQL Server avant de créer un pipeline. Voir Créer une connexion SQL Server.

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.

important

Toutes les demandes de création de pipeline doivent inclure "channel": "PREVIEW".

Définissez la ressource de pipeline dans un fichier de bundle (par exemple, resources/integrated_cdc_pipeline.yml) :

YAML
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_catalog: 'my_database'
source_schema: 'dbo'
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 :

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

Shell
databricks bundle deploy
databricks bundle run integrated_cdc_job

Pour plus d'information, consultez Que sont les bundles d'automatisation déclaratifs ?.

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

name

chaîne

Un nom pour le pipeline.

channel

chaîne

Doit être PREVIEW.

serverless

Booléen

Facultatif. La valeur par défaut est false. Définissez sur true pour le compute serverless ou sur false pour le compute classique. Le compute Serverless nécessite une mise en réseau Serverless vers votre base de données source.

catalog

chaîne

Le catalogue de destination default. Utilisé lorsqu'un destination_catalog par table n'est pas spécifié.

schema

chaîne

Le schéma de destination default. Utilisé lorsqu'un destination_schema par table n'est pas spécifié.

ingestion_definition.connection_name

chaîne

La connexion Unity Catalog à la base de données source.

ingestion_definition.connector_type

chaîne

Doit être CDC.

ingestion_definition.objects

tableau

La liste des tables ou schémas à ingérer.

ingestion_definition.data_staging_options

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.

parameter

Type

Description

name

chaîne

Un nom pour le pipeline.

channel

chaîne

Doit être PREVIEW.

serverless

Booléen

Facultatif. La valeur par défaut est false. Définissez sur true pour le compute serverless ou sur false pour le compute classique. Le compute Serverless nécessite une mise en réseau Serverless vers votre base de données source.

catalog

chaîne

Le catalogue de destination default. Utilisé lorsqu'un destination_catalog par table n'est pas spécifié.

schema

chaîne

Le schéma de destination default. Utilisé lorsqu'un destination_schema par table n'est pas spécifié.

ingestion_definition.connection_name

chaîne

La connexion Unity Catalog à la base de données source.

ingestion_definition.connector_type

chaîne

Doit être CDC.

ingestion_definition.objects

tableau

La liste des tables ou schémas à ingérer.

ingestion_definition.data_staging_options

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

source_catalog

Oui

Le nom de la base de données source.

source_schema

Oui

Le nom du schéma source.

source_table

Oui

Le nom de la table source.

destination_catalog

Non

Le catalogue de destination. Valeur par défaut : catalog du pipeline.

destination_schema

Non

Le schéma de destination. Valeur par défaut : schema du pipeline.

destination_table

Non

Le nom de la table de destination. La valeur default est source_table.

parameter

Obligatoire

Description

source_catalog

Oui

Le nom de la base de données source.

source_schema

Oui

Le nom du schéma source.

source_table

Oui

Le nom de la table source.

destination_catalog

Non

Le catalogue de destination. Valeur par défaut : catalog du pipeline.

destination_schema

Non

Le schéma de destination. Valeur par défaut : schema du pipeline.

destination_table

Non

Le nom de la table de destination. La valeur default est source_table.

Configuration de la table

parameter

Par défaut

Description

primary_keys

Détection automatique

Les colonnes qui identifient chaque ligne. Détection automatique à partir de la clé primaire source si non spécifiée.

scd_type

SCD_TYPE_1

SCD_TYPE_1 conserve uniquement la dernière version. SCD_TYPE_2 conserve l'historique complet et requiert la CDC de SQL Server sur la source. Le SCD de type 2 n'est pas pris en charge avec le suivi des modifications.

sequence_by

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

parameter

Par défaut

Description

primary_keys

Détection automatique

Les colonnes qui identifient chaque ligne. Détection automatique à partir de la clé primaire source si non spécifiée.

scd_type

SCD_TYPE_1

SCD_TYPE_1 conserve uniquement la dernière version. SCD_TYPE_2 conserve l'historique complet et requiert la CDC de SQL Server sur la source. Le SCD de type 2 n'est pas pris en charge avec le suivi des modifications.

sequence_by

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

Pour les mappages de types de données SQL Server, consultez la référence du connecteur SQL Server. 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.

    Text
    GET /api/2.0/pipelines/<pipeline-id>
  • API Événements.

    Text
    GET /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, ce qui peut prendre plus de temps que les mises à jour incrémentielles. Pour les grandes tables, l’instantané initial peut nécessiter plusieurs mises à jour planifiées pour être terminé. Chaque mise à jour suivante reprend là où la précédente s’était arrêtée.

Pour vérifier l'ingestion :

SQL
-- Check row counts in the destination table
SELECT COUNT(*) FROM <destination_catalog>.<destination_schema>.<destination_table>;

-- View recent changes (SCD Type 2 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 ont le dimensionnement automatique vertical activé 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_name et connector_type ne peuvent pas être modifiés après la création du pipeline. Pour modifier la source, créez un nouveau pipeline.
  • Maximum recommandé de 300 tables par pipeline.
  • **Instances principales uniquement.** Le connecteur CDC intégré ne prend pas en charge les réplicas en lecture, les instances de secours ou les instances secondaires.
  • **Tables sans clés primaires.** Le pipeline traite toutes les colonnes non LOB comme une clé composite. Les lignes en double peuvent être réduites en une seule ligne, sauf si vous activez le SCD de type 2.
  • L'instantané initial pourrait s'étendre sur plusieurs mises à jour. Pour les grandes tables, l'instantané initial pourrait ne pas se terminer en une seule mise à jour. Les mises à jour planifiées ultérieures reprennent là où la mise à jour précédente s'est arrêtée.
  • L'exécution de la mise à jour 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 une durée d'exécution maximale. Voir Fermeture intelligente pour les pipelines CDC intégrés. Vous ne pouvez pas configurer le runtime minimum ou maximum. Un backlog de modifications important peut s'étendre sur plusieurs mises à jour. Les mises à jour planifiées suivantes reprennent là où la mise à jour précédente s'est arrêtée.
  • La purge des logs nécessite un refresh complet. Si SQL Server purge les logs de suivi des modifications ou les logs CDC avant que le pipeline ne les traite, effectuez un refresh complet des tables affectées. Le pipeline détecte cette condition et affiche une erreur dans le log des événements.

Dépannage

remarque

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

NOT_IN_DEFAULT_PUBLISHING_MODE

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.

INGESTION_GATEWAY_CDC_NOT_ENABLED

La CDC ou le suivi des modifications n'est pas activé sur une ou plusieurs tables source.

Activez la CDC ou le suivi des modifications sur les tables affectées. Consultez Configurer Microsoft SQL Server pour l'ingestion dans Databricks.

INGESTION_GATEWAY_MISSING_TABLE_IN_SOURCE

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.

INGESTION_GATEWAY_SOURCE_SCHEMA_MISSING_ENTITY

Le schéma source n'existe pas.

Vérifiez que le schéma existe dans la base de données source.

UNSUPPORTED_SOURCE_TYPE_FOR_CDC_CONNECTOR

Le type de base de données source n'est pas pris en charge.

Le connecteur CDC intégré prend en charge SQL Server et Oracle.

SOURCE_TABLE_REQUIRED

Il manque source_table à la spécification de la table.

Ajoutez source_table à chaque spécification de table dans le tableau objects.

Integrated CDC connector is disabled

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.

Erreur

Cause

Résolution

NOT_IN_DEFAULT_PUBLISHING_MODE

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.

INGESTION_GATEWAY_CDC_NOT_ENABLED

La CDC ou le suivi des modifications n'est pas activé sur une ou plusieurs tables source.

Activez la CDC ou le suivi des modifications sur les tables affectées. Consultez Configurer Microsoft SQL Server pour l'ingestion dans Databricks.

INGESTION_GATEWAY_MISSING_TABLE_IN_SOURCE

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.

INGESTION_GATEWAY_SOURCE_SCHEMA_MISSING_ENTITY

Le schéma source n'existe pas.

Vérifiez que le schéma existe dans la base de données source.

UNSUPPORTED_SOURCE_TYPE_FOR_CDC_CONNECTOR

Le type de base de données source n'est pas pris en charge.

Le connecteur CDC intégré prend en charge SQL Server et Oracle.

SOURCE_TABLE_REQUIRED

Il manque source_table à la spécification de la table.

Ajoutez source_table à chaque spécification de table dans le tableau objects.

Integrated CDC connector is disabled

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 :

  1. Examinez le journal des événements du pipeline dans l'interface utilisateur de Databricks ou via GET /api/2.0/pipelines/<pipeline-id>/events.
  2. Testez la connexion Unity Catalog à partir de l’Explorateur de catalogues pour confirmer que la source est accessible.
  3. Confirmez que le suivi des modifications ou le CDC est activé sur la base de données source et les tables.
  4. Vérifiez que l'utilisateur de base de données dispose des autorisations SQL Server répertoriées dans Exigences relatives aux utilisateurs de base de données Microsoft SQL Server.
  5. Vérifiez que votre spécification de pipeline inclut "channel": "PREVIEW".

Ressources supplémentaires