Ingérer des données depuis SQL Server
Découvrez comment ingérer des données de SQL Server dans Databricks à l’aide de Lakeflow Connect.
Le connecteur SQL Server prend en charge Azure SQL Database, Azure SQL Managed Instance et les bases de données SQL Amazon RDS. Cela inclut SQL Server s'exécutant sur des machines virtuelles Azure (VM) et Amazon EC2. Le connecteur prend également en charge SQL Server on-premise en utilisant la mise en réseau Azure ExpressRoute et AWS Direct Connect.
Exigences
-
Pour créer une passerelle d'ingestion et un pipeline d'ingestion, vous devez d'abord satisfaire aux exigences 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. -
Vous avez accès à une instance SQL Server principale. Les fonctionnalités de suivi des modifications et de capture de changement de données ne sont pas prises en charge sur les réplicas en lecture ou les instances secondaires.
-
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 nœuds worker les plus petits possibles pour les passerelles d’ingestion, car ils n’ont pas d’impact sur 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 SQL Server, vous devez d'abord suivre les étapes dans Configurer Microsoft SQL Server pour l'ingestion dans Databricks.
Créer une passerelle et un pipeline d'ingestion
Ne pas arrêter manuellement la passerelle d'ingestion. La passerelle doit fonctionner en continu pour capturer les modifications avant que les logs de modification ne soient tronqués dans la base de données source. Si la passerelle est arrêtée, les modifications peuvent être perdues en raison de la conservation des logs, nécessitant un full refresh de toutes les tables affectées. L'arrêt et le redémarrage de la passerelle re-provisionnent également la VM, ce qui augmente le temps de Startup. Si vous avez besoin de résoudre les problèmes de passerelle, consultez Résoudre les problèmes d'ingestion SQL Server ou contactez le support Databricks.
- Databricks UI
- Declarative Automation Bundles
- Databricks notebook
- Terraform
-
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 SQL Server .
-
Sur la page **Connexion** de l'assistant d'ingestion, sélectionnez la connexion qui stocke vos identifiants d'accès SQL Server. Si vous disposez du privilège
CREATE CONNECTIONsur le metastore, vous pouvez cliquer surCréer une connexion pour créer une nouvelle connexion avec les détails d'authentification dans Créer une connexion SQL Server.
-
Cliquez sur Suivant .
-
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.
-
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 CATALOGetCREATE SCHEMAsur le catalogue, vous pouvez cliquer surCréer un schéma dans le menu déroulant pour créer un nouveau schéma.
-
(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.
-
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.
-
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 CATALOGetCREATE SCHEMAsur le catalogue, vous pouvez cliquer surCréer un schéma dans le menu déroulant pour créer un nouveau schéma.
-
Cliquez sur Créer un pipeline et continuer .
-
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).
-
Cliquez sur Suivant , puis sur Enregistrer et continuer .
-
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 CATALOGetCREATE SCHEMAsur le catalogue, vous pouvez cliquer surCréer un schéma dans le menu déroulant pour créer un nouveau schéma.
-
Cliquez sur Enregistrer et continuer .
-
Sur la page Configuration de la base de données , cliquez sur Valider pour confirmer que votre source est correctement configurée pour l'ingestion Databricks. Toutes les configurations manquantes sont renvoyées. Pour les étapes de résolution, cliquez sur Terminer la configuration . Cliquez ensuite sur Suivant . Alternativement, cliquez sur Ignorer la validation .
-
(Facultatif) Sur la page Schedules and notifications , cliquez sur
Créer un calendrier . Définissez la fréquence pour refresh les tables de destination.
-
(Facultatif) Cliquez sur
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 à l'aide des Declarative Automation Bundles, vous devez avoir accès à une connexion existante. Pour obtenir des instructions, consultez Créer une connexion SQL Server.
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. Spécifiez l'emplacement de staging dans la section gateway_definition de votre fichier YAML de pipeline de bundle.
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. Cela aide à prendre en compte toutes les politiques de rétention des logs de modification que vous avez 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.
Les bundles peuvent contenir des définitions YAML de jobs et de tâches, sont gérés à l'aide de la CLI Databricks, et peuvent être partagés et exécutés dans différents workspaces cibles (tels que le développement, la pré-production et la production). Pour plus d'information, consultez Que sont les Declarative Automation Bundles ?.
-
Créez un bundle à l'aide de la CLI Databricks :
Bashdatabricks bundle init -
Ajoutez votre configuration de pipeline et de job au bundle. Consultez Exemples pour un exemple complet avec toutes les options disponibles.
-
Déployez le pipeline à l'aide de la CLI Databricks :
Bashdatabricks bundle deploy
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.
Vous pouvez utiliser Terraform pour déployer et gérer des pipelines d'ingestion SQL Server. 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.
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.

Les colonnes Upserted records et Deleted records ne sont pas affichées par default. Vous pouvez les activer en cliquant sur le bouton de configuration des colonnes et en les sélectionnant.
Exemples
Utilisez ces exemples pour configurer votre pipeline.
Configuration de pipeline
- Declarative Automation Bundles
- Databricks notebook
Le bundle suivant définit un pipeline de passerelle, un pipeline d'ingestion et un Job planifié. Les options commentées affichent toutes les configurations disponibles. Mettez à jour les sections variables et targets avec les détails de votre source et de votre destination.
bundle:
name: lakeflow-connect-sqlserver
# Variables parameterize the bundle for different environments and sources.
# Set values here, override per-target, or pass with: databricks bundle deploy -var="key=value"
variables:
# The name of the Unity Catalog connection to your SQL Server instance.
# This connection must already exist and be of type SQLSERVER.
connection_name:
description: 'Unity Catalog connection name for the SQL Server source'
# The SQL Server database name to ingest from.
# In Lakeflow Connect, this maps to source_catalog in the table/schema spec.
source_database:
description: 'SQL Server database name (maps to source_catalog in table specs)'
# The SQL Server schema to ingest from (for example, "dbo", "sales").
source_schema:
description: 'SQL Server schema name to ingest from'
# The Unity Catalog catalog where ingested Delta tables are created.
dest_catalog:
description: 'Destination Unity Catalog catalog for ingested tables'
# The Unity Catalog schema where ingested Delta tables are created.
dest_schema:
description: 'Destination Unity Catalog schema for ingested tables'
# The Unity Catalog catalog for the gateway's internal staging volume.
# Can be the same as dest_catalog. Must not be a foreign catalog.
staging_catalog:
description: 'Catalog for gateway staging volume'
# The Unity Catalog schema for the gateway's internal staging volume.
staging_schema:
description: 'Schema for gateway staging volume'
resources:
pipelines:
# --- Gateway pipeline ---
# Extracts change data from SQL Server and stages it in a Unity Catalog
# volume. Must run continuously to capture changes before change logs are
# truncated in the source database.
gw_pipeline:
name: 'lfc-sqlserver-gateway-${bundle.target}'
# Gateway pipelines must be continuous.
continuous: true
# "CURRENT" (stable) or "PREVIEW" (early access).
channel: 'CURRENT'
# (Optional) Associate with a budget policy for cost tracking.
# budget_policy_id: "<policy-uuid>"
# The gateway runs on classic compute. Cluster settings are managed
# automatically. You can optionally customize the cluster:
# clusters:
# - label: "default"
# autoscale:
# min_workers: 1
# max_workers: 4
# # node_type_id: "i3.xlarge"
# # Restrict the cluster to an approved cluster policy.
# # policy_id: "<cluster-policy-id>"
catalog: ${var.staging_catalog}
schema: ${var.staging_schema}
gateway_definition:
# (Required) Unity Catalog connection name (type SQLSERVER).
connection_name: ${var.connection_name}
# (Required) Catalog and schema for the staging volume.
gateway_storage_catalog: ${var.staging_catalog}
gateway_storage_schema: ${var.staging_schema}
# (Optional) Custom staging volume name. If not set, the system
# auto-generates: __databricks_ingestion_gateway_staging_data-<pipeline_id>
# gateway_storage_name: "my_custom_staging_volume"
# --- Ingestion pipeline ---
# Reads staged data from the gateway and applies it to Delta tables.
mi_pipeline:
name: 'lfc-sqlserver-ingestion-${bundle.target}'
# Continuous mode is not supported for the ingestion pipeline.
# Use a scheduled job to trigger runs.
continuous: false
channel: 'CURRENT'
# (Optional) Associate with a budget policy for cost tracking.
# budget_policy_id: "<policy-uuid>"
# The ingestion pipeline runs on serverless compute only.
serverless: true
# (Optional) Development mode for faster iteration (no retries).
# development: true
catalog: ${var.dest_catalog}
schema: ${var.dest_schema}
# (Optional) Email notifications for pipeline events.
# notifications:
# - email_recipients:
# - "team@example.com"
# alerts:
# - "on-update-failure"
# - "on-update-fatal-failure"
# - "on-flow-failure"
# (Optional) Run as a service principal for production.
# run_as:
# service_principal_name: "my-service-principal"
ingestion_definition:
# (Required) References the gateway pipeline. The connection is
# inherited from the gateway. Do not specify connection_name here.
ingestion_gateway_id: ${resources.pipelines.gw_pipeline.id}
# Pipeline-level table configuration defaults. These apply to all
# tables unless overridden at the schema or table level.
table_configuration:
# SCD Type: How changes are applied to destination tables.
# SCD_TYPE_1: Overwrites rows with latest values (default).
# SCD_TYPE_2: Preserves history with __START_AT/__END_AT columns.
# Requires CDC on source. CT does not support SCD_TYPE_2.
# APPEND_ONLY: Inserts only. Updates and deletes are ignored.
scd_type: 'SCD_TYPE_1'
# (Optional) Auto full refresh policy. Triggers a snapshot when the
# pipeline detects issues resolvable by re-reading all source data
# (for example, CT/CDC retention window expired).
# auto_full_refresh_policy:
# enabled: true
# min_interval_hours: 24
# (Optional) Schedule automatic full refreshes.
# full_refresh_window:
# start_hour: 2
# days_of_week:
# - "SUNDAY"
# time_zone_id: "America/Los_Angeles"
objects:
# Option 1: Schema-level ingestion. Ingests all tables from a source
# schema. New tables added to the schema are picked up automatically.
- schema:
source_catalog: ${var.source_database}
source_schema: ${var.source_schema}
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
# (Optional) Override table_configuration for this schema.
# table_configuration:
# scd_type: "SCD_TYPE_2"
# Option 2: Table-level ingestion. Provides granular control.
# Replace or combine with the schema-level spec.
# - table:
# source_catalog: ${var.source_database}
# source_schema: ${var.source_schema}
# source_table: "customers"
# destination_catalog: ${var.dest_catalog}
# destination_schema: ${var.dest_schema}
# # (Optional) Rename the table at the destination.
# # destination_table: "customers_v2"
# table_configuration:
# scd_type: "SCD_TYPE_1"
# # Include only specific columns (mutually exclusive with exclude_columns).
# # include_columns:
# # - "customer_id"
# # - "first_name"
# # - "email"
# # Exclude specific columns. All other columns are included.
# # exclude_columns:
# # - "internal_notes"
# # Override the primary key used for change detection.
# # primary_keys:
# # - "customer_id"
# # Logical ordering columns for change resolution.
# # sequence_by:
# # - "updated_at"
# # Auto full refresh for this table.
# # auto_full_refresh_policy:
# # enabled: true
# # min_interval_hours: 48
# (Optional) Grant additional users or groups access.
# permissions:
# - user_name: "analyst@example.com"
# level: "CAN_VIEW"
# - group_name: "data-engineers"
# level: "CAN_RUN"
# --- Scheduled job ---
# Triggers the ingestion pipeline on a schedule.
jobs:
mi_schedule:
name: 'lfc-sqlserver-ingestion-schedule-${bundle.target}'
# Quartz cron syntax: "seconds minutes hours day month day-of-week"
# Examples: "0 0 * * * ?" (hourly), "0 0 */4 * * ?" (every 4 hours)
schedule:
quartz_cron_expression: '0 */30 * * * ?'
timezone_id: 'UTC'
tasks:
- task_key: 'run_ingestion'
pipeline_task:
pipeline_id: ${resources.pipelines.mi_pipeline.id}
# email_notifications:
# on_failure:
# - "team@example.com"
# Deploy to different workspaces with: databricks bundle deploy -t <target>
targets:
dev:
default: true
workspace:
host: https://<workspace-url>.cloud.databricks.com
variables:
connection_name: '<sqlserver-connection>'
source_database: '<database-name>'
source_schema: 'dbo'
dest_catalog: '<dest-catalog>'
dest_schema: '<dest-schema>'
staging_catalog: '<staging-catalog>'
staging_schema: '<staging-schema>'
Voici une section Configuration d'exemple d'une spécification de pipeline :
# The name of the UC connection with the credentials to access the source database
connection_name = "my_connection"
# The name of the UC catalog and schema to store the replicated tables
target_catalog_name = "main"
target_schema_name = "lakeflow_sqlserver_connector_cdc"
# The name of the UC catalog and schema to store the staging volume with intermediate
# CDC and snapshot data. Use the destination catalog/schema by default.
stg_catalog_name = target_catalog_name
stg_schema_name = target_schema_name
# The name of the Gateway pipeline to create
gateway_pipeline_name = "cdc_gateway"
# The name of the Ingestion pipeline to create
ingestion_pipeline_name = "cdc_ingestion"
# Construct the full list of tables to replicate.
# IMPORTANT: The letter case of catalog, schema, and table names must match exactly
# the case used in the source database system tables.
tables_to_replicate = replicate_full_db_schema("MY_DB", ["MY_DB_SCHEMA"])
# Append tables from additional schemas as needed:
# + replicate_tables_from_db_schema("MY_DB", "MY_SCHEMA_2", ["table3", "table4"])
Modèles courants
Pour les configurations de pipeline avancées, consultez Modèles courants pour les pipelines d'ingestion gérés.
Étapes suivantes
start, planifiez et configurez des alertes sur votre pipeline. Consultez les Tâches de maintenance courantes du pipeline.