Hinweis
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, sich anzumelden oder das Verzeichnis zu wechseln.
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, das Verzeichnis zu wechseln.
Erfahren Sie, wie Sie Daten aus SQL Server mithilfe von Lakeflow Connect in Azure Databricks aufnehmen.
Der SQL Server-Connector unterstützt Azure SQL-Datenbank, azure SQL Managed Instance und Amazon RDS SQL-Datenbanken. Dazu gehört SQL Server, der auf virtuellen Azure-Computern (VMs) und Amazon EC2 ausgeführt wird. Der Connector unterstützt auch sql Server lokal mit Azure ExpressRoute und AWS Direct Connect-Netzwerk.
Anforderungen
Um ein Datenaufnahme-Gateway und eine Datenaufnahme-Pipeline zu erstellen, müssen Sie zuerst die folgenden Anforderungen erfüllen:
Ihr Arbeitsbereich ist für Unity Catalog aktiviert.
Serverlose Berechnung ist für Ihren Arbeitsbereich aktiviert. Siehe Serverlose Computeanforderungen.
Wenn Sie eine Verbindung erstellen möchten: Sie verfügen über
CREATE CONNECTIONBerechtigungen für den Metastore. Weitere Informationen finden Sie unter Verwalten von Berechtigungen in Unity Catalog.Wenn Ihr Connector die benutzeroberflächenbasierte Pipelineerstellung unterstützt, können Sie die Verbindung und die Pipeline gleichzeitig erstellen, indem Sie die Schritte auf dieser Seite ausführen. Wenn Sie jedoch apibasierte Pipelineerstellung verwenden, müssen Sie die Verbindung im Katalog-Explorer erstellen, bevor Sie die Schritte auf dieser Seite ausführen. Siehe Herstellen einer Verbindung mit verwalteten Aufnahmequellen.
Wenn Sie beabsichtigen, eine vorhandene Verbindung zu verwenden: Sie verfügen über
USE CONNECTIONBerechtigungen oderALL PRIVILEGESfür die Verbindung.Sie verfügen über
USE CATALOG-Berechtigungen für den Zielkatalog.Sie verfügen über
USE SCHEMA,CREATE TABLEundCREATE VOLUMEBerechtigungen für ein vorhandenes Schema oder überCREATE SCHEMABerechtigungen für den Zielkatalog.
Sie haben Zugriff auf eine primäre SQL Server-Instanz. Die Funktionen zur Änderungsverfolgung und zur Erfassung von Änderungsdaten werden auf Leserepliken oder sekundären Instanzen nicht unterstützt.
Uneingeschränkte Berechtigungen zum Erstellen von Clustern oder einer benutzerdefinierten Richtlinie (nur API). Eine benutzerdefinierte Richtlinie für das Gateway muss die folgenden Anforderungen erfüllen:
Familie: Job-Compute
Außerkraftsetzungen der Richtlinienfamilie:
{ "cluster_type": { "type": "fixed", "value": "dlt" }, "num_workers": { "type": "unlimited", "defaultValue": 1, "isOptional": true }, "runtime_engine": { "type": "fixed", "value": "STANDARD", "hidden": true } }Databricks empfiehlt, die kleinstmöglichen Workerknoten für Erfassungs-Gateways anzugeben, da sie sich nicht auf die Leistung der Gateways auswirken. Mit der folgenden Berechnungsrichtlinie können Azure Databricks das Aufnahmegateway skalieren, um die Anforderungen Ihrer Workload zu erfüllen. Die Mindestanforderung ist 8 Kerne, um eine effiziente und leistungsfähige Datenextraktion aus Der Quelldatenbank zu ermöglichen.
{ "driver_node_type_id": { "type": "fixed", "value": "Standard_E64d_v4" }, "node_type_id": { "type": "fixed", "value": "Standard_F4s" } }
Weitere Informationen zu Clusterrichtlinien finden Sie unter Auswählen einer Richtlinie für Rechenressourcen.
Zum Aufnehmen von SQL Server müssen Sie zunächst die Schritte unter Konfigurieren von Microsoft SQL Server für die Aufnahme in Azure Databricks ausführen.
Ein Gateway und eine Aufnahmepipeline erstellen
Warning
Stoppen Sie das Erfassungs-Gateway nicht manuell. Das Gateway muss kontinuierlich laufen, um Änderungen zu erfassen, bevor Änderungsprotokolle in der Quelldatenbank gekürzt werden. Wenn das Gateway beendet wird, können Änderungen aufgrund der Protokollaufbewahrung gelöscht werden, sodass eine vollständige Aktualisierung aller betroffenen Tabellen erforderlich ist. Das Beenden und Neustarten des Gateways stellt auch die VM erneut bereit, was die Startzeit verlängert. Wenn Sie Probleme mit dem Gateway beheben müssen, lesen Sie Problembehandlung bei der SQL Server-Datenerfassung oder wenden Sie sich an den Databricks-Support.
Databricks UI
Klicken Sie in der Randleiste des Azure Databricks-Arbeitsbereichs auf "Datenaufnahme".
Klicken Sie auf der Seite " Daten hinzufügen " unter "Databricks-Connectors" auf SQL Server.
Wählen Sie auf der Seite Connection des Erfassungsassistenten die Verbindung aus, in der Ihre SQL Server-Zugangsdaten gespeichert sind. Wenn Sie über das
CREATE CONNECTION-Privileg für den Metastore verfügen, Sie können aufVerbindung erstellen klicken, um eine neue Verbindung mit den Authentifizierungsdetails in Create a SQL Server connection zu erstellen.
Klicke auf Weiter.
Geben Sie auf der Seite Einrichtung der Erfassung einen eindeutigen Namen für die Erfassungspipeline ein. Diese Pipeline verschiebt Daten vom Zwischenspeicherort zum Ziel.
Wählen Sie einen Katalog und ein Schema aus, in das Ereignisprotokolle geschrieben werden sollen. Das Ereignisprotokoll enthält Überwachungsprotokolle, Datenqualitätsprüfungen, Pipelinefortschritte und Fehler. Wenn Sie
und Berechtigungen im Katalog haben, können Sie im Dropdownmenü auf das klicken, um ein neues Schema zu erstellen.Plussymbol (Optional) Festlegen der automatischen vollständigen Aktualisierung für alle Tabellen auf "Ein". Wenn die automatische Aktualisierung aktiviert ist, versucht die Pipeline automatisch, Probleme wie Protokollbereinigungsereignisse und bestimmte Arten der Schemaentwicklung zu beheben, indem die betroffene Tabelle vollständig aktualisiert wird. Wenn die Verlaufsverfolgung aktiviert ist, wird dieser Verlauf durch eine vollständige Aktualisierung gelöscht.
Geben Sie einen eindeutigen Namen für das Dateneingangs-Gateway ein. Das Gateway ist eine Pipeline, die Änderungen aus der Quelle extrahiert und für das Laden in die Erfassungspipeline bereitstellt.
Wählen Sie einen Katalog und ein Schema für den Zwischenspeicherort aus. An diesem Speicherort wird für das Staging von extrahierten Daten ein Volume erstellt. Wenn Sie
und Berechtigungen im Katalog haben, können Sie im Dropdownmenü auf das klicken, um ein neues Schema zu erstellen.Plussymbol Klicken Sie auf "Pipeline erstellen", und fahren Sie fort.
Wählen Sie auf der Seite "Quelle" die aufzunehmenden Tabellen aus. Wenn Sie bestimmte Tabellen auswählen, können Sie Tabelleneinstellungen konfigurieren:
a) (Optional) Geben Sie auf der Registerkarte "Einstellungen " einen Zielnamen für jede aufgenommene Tabelle an. Dies ist nützlich, um zwischen Zieltabellen zu unterscheiden, wenn Sie ein Objekt mehrmals in dasselbe Schema aufnehmen. Siehe Name einer Zieltabelle.
a) (Optional) Ändern Sie die Standardeinstellung für die Verlaufsverfolgung . Siehe "Verlaufsverfolgung aktivieren" (SCD-Typ 2).
Klicken Sie auf "Weiter", und klicken Sie dann auf " Speichern", und fahren Sie fort.
Wählen Sie auf der Seite "Ziel " einen Katalog und ein Schema aus, in das Daten geladen werden sollen. Wenn Sie
und Berechtigungen im Katalog haben, können Sie im Dropdownmenü auf das klicken, um ein neues Schema zu erstellen.Plussymbol Klicken Sie auf Speichern und fortfahren.
Klicken Sie auf der Seite "Datenbankeinrichtung " auf " Überprüfen ", um zu bestätigen, dass Ihre Quelle für die Aufnahme von Azure Databricks ordnungsgemäß konfiguriert ist. Fehlende Konfigurationen werden zurückgegeben. Um Schritte zur Behebung auszuführen, klicken Sie auf "Konfiguration abschließen". Klicken Sie dann auf Weiter. Alternativ können Sie auf "Überprüfung überspringen" klicken.
(Optional) Klicken Sie auf der Seite "Zeitpläne und Benachrichtigungen " auf das
Zeitplan erstellen. Legen Sie die Häufigkeit fest, mit der die Zieltabellen aktualisiert werden.
(Optional) Klicken Sie auf
Fügen Sie eine Benachrichtigung hinzu, um E-Mail-Benachrichtigungen für Erfolg oder Fehler des Pipelinevorgangs festzulegen, und klicken Sie dann auf " Speichern und Ausführen der Pipeline".
Deklarative Automatisierungspakete
Bevor Sie Daten mithilfe deklarativer Automatisierungspakete erfassen, müssen Sie Zugriff auf eine bestehende Verbindung haben. Anweisungen finden Sie unter Create a SQL Server connection.
Der Stagingkatalog und das Stagingschema können mit dem Zielkatalog und -schema identisch sein. Bei dem Staging-Katalog kann es sich nicht um einen externen Katalog handeln. Geben Sie den Stagingspeicherort im gateway_definition-Abschnitt der YAML-Datei der Paketpipeline an.
Das Gateway zum Einbinden extrahiert Snapshot- und Änderungsdaten aus der Quelldatenbank und speichert sie im Unity Catalog Staging Volume. Sie müssen das Gateway als fortlaufende Pipeline ausführen. Dadurch können Sie alle Aufbewahrungsrichtlinien für Änderungsprotokolle berücksichtigen, die Sie in der Quelldatenbank haben.
Die Erfassungspipeline wendet die Momentaufnahme- und Änderungsdaten aus dem Stagingvolume in Zielstreamingtabellen an.
Bundles können YAML-Definitionen von Aufträgen und Aufgaben enthalten, mithilfe der Databricks CLI verwaltet und in verschiedenen Zielarbeitsbereichen (z. B. Entwicklung, Staging und Produktion) freigegeben und ausgeführt werden. Weitere Informationen finden Sie unter Was sind deklarative Automatisierungs-Bundles?.
Erstellen Eines Bündels mithilfe der Databricks CLI:
databricks bundle initFügen Sie die Pipeline- und Jobkonfiguration dem Bundle hinzu. Siehe Beispiele für ein vollständiges Beispiel mit allen verfügbaren Optionen.
Stellen Sie die Pipeline mithilfe der Databricks CLI bereit:
databricks bundle deploy
Databricks-Notizbuch
Aktualisieren Sie die Configuration-Zelle im folgenden Notebook mit der Quellverbindung, dem Zielkatalog, dem Zielschema und den Tabellen, um Daten aus der Quelle zu erfassen.
Terraform
Sie können Terraform verwenden, um SQL Server-Aufnahmepipelines bereitzustellen und zu verwalten. Ein vollständiges Beispielframework, einschließlich Terraform-Konfigurationen zum Erstellen von Gateways und Aufnahmepipelinen, finden Sie im Repository für Lakeflow Connect Terraform-Beispiele auf GitHub.
Überprüfen der erfolgreichen Datenerfassung
Die Listenansicht auf der Detailseite der Pipeline zeigt die Anzahl der Datensätze an, die beim Einbinden der Daten verarbeitet werden. Diese Zahlen werden automatisch aktualisiert.
Die Spalten Upserted records und Deleted records werden standardmäßig nicht angezeigt. Sie können sie aktivieren, indem Sie auf die Schaltfläche für die Spaltenkonfiguration (
) klicken und die Spalten auswählen.
Beispiele
Verwenden Sie diese Beispiele, um Ihre Pipeline zu konfigurieren.
Pipelinekonfiguration
Deklarative Automatisierungspakete
Das folgende Paket definiert eine Gateway-Pipeline, eine Ingestion-Pipeline und einen geplanten Job. Auskommentierte Optionen zeigen alle verfügbaren Konfigurationen an. Aktualisieren Sie die Abschnitte variables und targets mit Ihren Quell- und Zieldetails.
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>'
Databricks-Notizbuch
Im Folgenden sehen Sie einen Beispielabschnitt Configuration einer Pipelinespezifikation:
# 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"])
Allgemeine Muster
Weitere Informationen zu erweiterten Pipelinekonfigurationen finden Sie unter "Allgemeine Muster für verwaltete Aufnahmepipelinen".
Nächste Schritte
Starten Sie, planen Sie und legen Sie Benachrichtigungen für Ihre Pipeline fest. Siehe allgemeine Pipelinewartungsaufgaben.