Importare dati da SQL Server

Informazioni su come inserire dati da SQL Server in Azure Databricks usando Lakeflow Connect.

Il connettore SQL Server supporta il database SQL di Azure, l'istanza gestita di SQL di Azure e i database SQL di Amazon RDS. Sono inclusi SQL Server in esecuzione in macchine virtuali (VM) di Azure e Amazon EC2. Il connettore supporta anche SQL Server in locale usando la rete di Azure ExpressRoute e AWS Direct Connect.

Requisiti

  • Per creare un gateway di inserimento e una pipeline di inserimento, è prima necessario soddisfare i requisiti seguenti:

    • L'area di lavoro è configurata per Unity Catalog.

    • Il calcolo serverless è abilitato per il tuo spazio di lavoro. Consulta Requisiti di calcolo serverless.

    • Se prevedi di creare una connessione: hai privilegi CREATE CONNECTION nel metastore. Consulta Gestione dei privilegi in Unity Catalog.

      Se il connettore supporta la creazione di pipeline basate sull'interfaccia utente, è possibile creare la connessione e la pipeline contemporaneamente completando i passaggi in questa pagina. Tuttavia, se si usa la creazione di pipeline basate su API, è necessario creare la connessione in Esplora cataloghi prima di completare i passaggi in questa pagina. Vedere Connettersi alle origini di inserimento gestite.

    • Se si prevede di usare una connessione esistente: si dispone di USE CONNECTION privilegi o ALL PRIVILEGES sulla connessione.

    • Hai i privilegi USE CATALOG sul catalogo di destinazione.

    • Si dispone di USE SCHEMA privilegi, CREATE TABLE, e CREATE VOLUME per uno schema esistente o CREATE SCHEMA privilegi sul catalogo di destinazione.

    • Hai accesso a un'istanza primaria di SQL Server. Le funzionalità di Rilevamento delle modifiche e Change Data Capture non sono supportate nelle repliche di lettura o nelle istanze secondarie.

    • Autorizzazioni senza restrizioni per creare cluster o criteri personalizzati (solo API). Un criterio personalizzato per il gateway deve soddisfare i requisiti seguenti:

      • Famiglia: Calcolo del Lavoro

      • Esclusioni della famiglia di criteri:

        {
          "cluster_type": {
            "type": "fixed",
            "value": "dlt"
          },
          "num_workers": {
            "type": "unlimited",
            "defaultValue": 1,
            "isOptional": true
          },
          "runtime_engine": {
            "type": "fixed",
            "value": "STANDARD",
            "hidden": true
          }
        }
        
      • Databricks consiglia di specificare i nodi di lavoro più piccoli possibili per i gateway di inserimento perché non influiscono sulle prestazioni del gateway. I criteri di calcolo seguenti consentono ad Azure Databricks di ridimensionare il gateway di inserimento per soddisfare le esigenze del carico di lavoro. Il requisito minimo è di 8 core per consentire l'estrazione efficiente ed efficiente dei dati dal database di origine.

        {
          "driver_node_type_id": {
            "type": "fixed",
            "value": "Standard_E64d_v4"
          },
          "node_type_id": {
            "type": "fixed",
            "value": "Standard_F4s"
          }
        }
        

      Per altre informazioni sui criteri del cluster, vedere Selezionare un criterio di calcolo.

  • Per inserire da SQL Server, è prima necessario completare i passaggi descritti in Configurare Microsoft SQL Server per l'inserimento in Azure Databricks.

Creare un gateway e una pipeline di inserimento

Avvertimento

Non arrestare manualmente il gateway di inserimento. Il gateway deve essere in esecuzione continua per rilevare le modifiche prima che i registri delle modifiche vengano troncati nel database di origine. Se il gateway viene arrestato, le modifiche possono essere eliminate a causa della conservazione dei log, che richiede un aggiornamento completo di tutte le tabelle interessate. Anche l'arresto e il riavvio del gateway comportano un nuovo provisioning della macchina virtuale, aumentando così il tempo di avvio. Se devi risolvere i problemi del gateway, consulta Risolvere i problemi di acquisizione di SQL Server o contatta l'assistenza Databricks.

Interfaccia utente di Databricks

  1. Nella barra laterale dell'area di lavoro di Azure Databricks fare clic su Inserimento dati.

  2. Nella pagina Aggiungi dati , in Connettori Databricks, fare clic su SQL Server.

  3. Nella pagina Connection della procedura guidata di inserimento selezionare la connessione in cui sono archiviate le credenziali di accesso SQL Server. Se si dispone del privilegio CREATE CONNECTION nel metastore, È possibile fare clic su Plus icon. Crea connessione per creare una nuova connessione con i dettagli di autenticazione in Creare una connessione SQL Server.

  4. Fare clic su Avanti.

  5. Nella pagina Configurazione dell'acquisizione immettere un nome univoco per la pipeline di acquisizione. Questa pipeline sposta i dati dall'area di staging alla destinazione.

  6. Selezionare un catalogo e uno schema in cui scrivere i registri eventi. Il registro eventi contiene log di controllo, controlli di qualità dei dati, stato della pipeline ed errori. Se dispongono di privilegi e sul catalogo, è possibile fare clic sull'icona Più nel menu a discesa per selezionare l'opzione "Crea schema" e creare un nuovo schema.

  7. (Facoltativo) Impostare Aggiornamento automatico completo per tutte le tabellesu Sì. Quando l'aggiornamento automatico è attivo, la pipeline tenta automaticamente di risolvere problemi come gli eventi di pulizia dei log e alcuni tipi di evoluzione dello schema aggiornando completamente la tabella interessata. Se il rilevamento della cronologia è abilitato, un aggiornamento completo cancella tale cronologia.

  8. Immettere un nome univoco per il gateway di inserimento. Il gateway è una pipeline che estrae le modifiche dall'origine e le prepara per il caricamento nella pipeline di inserimento.

  9. Selezionare un catalogo e uno schema per il percorso di gestione temporanea. Un volume viene creato in questa posizione per preparare i dati estratti. Se dispongono di privilegi e sul catalogo, è possibile fare clic sull'icona Più nel menu a discesa per selezionare l'opzione "Crea schema" e creare un nuovo schema.

  10. Fare clic su Crea pipeline e continuare.

  11. Nella pagina Origine selezionare le tabelle da inserire. Se si selezionano tabelle specifiche, è possibile configurare le impostazioni della tabella:

    a. (Facoltativo) Nella scheda Impostazioni specificare un nome di destinazione per ogni tabella inserita. Ciò è utile per distinguere le tabelle di destinazione quando si inserisce un oggetto nello stesso schema più volte. Vedere Assegnare un nome a una tabella di destinazione.

    a. (Facoltativo) Modificare l'impostazione di rilevamento cronologia predefinita. Consulta Abilitare il rilevamento della cronologia (SCD tipo 2).

  12. Fare clic su Avanti, quindi su Salva e continua.

  13. Nella pagina Destinazione selezionare un catalogo e uno schema in cui caricare i dati. Se dispongono di privilegi e sul catalogo, è possibile fare clic sull'icona Più nel menu a discesa per selezionare l'opzione "Crea schema" e creare un nuovo schema.

  14. Fare clic su Salva e continua.

  15. Nella pagina Configurazione database fare clic su Convalida per verificare che l'origine sia configurata correttamente per l'inserimento di Azure Databricks. Vengono restituite eventuali configurazioni mancanti. Per i passaggi da risolvere, fare clic su Completa configurazione. Fare quindi clic su Avanti. In alternativa, fare clic su Ignora convalida.

  16. (Facoltativo) Nella pagina Pianificazioni e notifiche fare clic sull'icona Più. Creare una pianificazione. Impostare la frequenza per aggiornare le tabelle di destinazione.

  17. (Facoltativo) Fare clic sull'icona Con il segno più. Aggiungere una notifica per impostare le notifiche di posta elettronica per l'esito positivo o negativo dell'operazione della pipeline, quindi fare clic su Salva ed esegui pipeline.

Pacchetti di automazione dichiarativa

Prima di importare utilizzando i bundle di automazione dichiarativa, è necessario avere accesso a una connessione esistente. Per istruzioni, vedere Creare una connessione SQL Server.

Il catalogo e lo schema di staging possono essere gli stessi del catalogo e dello schema di destinazione. Il catalogo di staging non può essere un catalogo straniero. Specificare il percorso di gestione temporanea nella sezione gateway_definition del file YAML della pipeline di pacchetti.

Il gateway di inserimento estrae gli snapshot e i dati delle modifiche dal database di origine e li archivia nel volume di staging di Unity Catalog. È necessario eseguire il gateway come una pipeline continua. Ciò consente di gestire eventuali criteri di conservazione dei log delle modifiche presenti nel database di origine.

La pipeline di inserimento applica lo snapshot e modifica i dati dal volume di staging alle tabelle di streaming di destinazione.

I bundle possono contenere definizioni YAML di processi e attività, vengono gestiti tramite l'interfaccia della riga di comando di Databricks e possono essere condivisi ed eseguiti in aree di lavoro di destinazione diverse, ad esempio sviluppo, gestione temporanea e produzione. Per altre informazioni, vedere Che cosa sono i bundle di automazione dichiarativa?.

  1. Creare un bundle usando l'interfaccia della riga di comando di Databricks:

    databricks bundle init
    
  2. Aggiungi la pipeline e la configurazione del job al bundle. Vedere Esempi per un esempio completo con tutte le opzioni disponibili.

  3. Distribuire la pipeline usando la CLI di Databricks.

    databricks bundle deploy
    

Notebook di Databricks

Aggiornare la Configuration cella nel notebook seguente con la connessione di origine, il catalogo di destinazione, lo schema di destinazione e le tabelle da acquisire dalla sorgente.

Prendi il computer portatile

Terraform

È possibile usare Terraform per distribuire e gestire le pipeline di inserimento di SQL Server. Per un framework di esempio completo, incluse le configurazioni di Terraform per la creazione di gateway e pipeline di ingestione, vedere il repository degli esempi di Lakeflow Connect Terraform su GitHub.

Verificare il successo dell'inserimento dati

La visualizzazione elenco nella pagina dei dettagli della pipeline mostra il numero di record elaborati durante l'inserimento dei dati. Questi numeri vengono aggiornati automaticamente.

Verificare la replica

Le Upserted records colonne e Deleted records non vengono visualizzate per impostazione predefinita. È possibile abilitarli facendo clic sul pulsante Icona di configurazione colonne e selezionandoli.

Examples

Usare questi esempi per configurare la pipeline.

Configurazione della pipeline

Pacchetti di automazione dichiarativa

Il bundle seguente definisce una pipeline del gateway, una pipeline di inserimento e un processo pianificato. Le opzioni impostate come commento mostrano tutte le configurazioni disponibili. Aggiorna le sezioni variables e targets con i dettagli di origine e destinazione.

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>'

Notebook di Databricks

Di seguito è riportata una sezione di esempio Configuration di una specifica della 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"])

Modelli comuni

Per le configurazioni avanzate della pipeline, vedere Modelli comuni per le pipeline di inserimento gestite.

Passaggi successivi

Avvia, pianifica e imposta avvisi sulla tua pipeline. Vedere Attività comuni di manutenzione della pipeline.

Risorse aggiuntive