Ingerir dados do SQL Server

Aprenda como ingerir dados do SQL Server para Azure Databricks usando o Lakeflow Connect.

O conector do SQL Server oferece suporte ao Banco de Dados SQL do Azure, à Instância Gerenciada SQL do Azure e aos bancos de dados SQL do Amazon RDS. Isso inclui o SQL Server em execução em máquinas virtuais (VMs) do Azure e o Amazon EC2. O conector também oferece suporte ao SQL Server local usando a rede Azure ExpressRoute e AWS Direct Connect.

Requisitos

  • Para criar um gateway de ingestão e um pipeline de ingestão, deve primeiro cumprir os seguintes requisitos:

    • Seu espaço de trabalho está habilitado para o Catálogo Unity.

    • A computação sem servidor está habilitada para seu espaço de trabalho. Consulte Requisitos de computação sem servidor.

    • Se tencionas criar uma conexão: tens CREATE CONNECTION privilégios no metastore. Consulte Gerenciar privilégios no Catálogo Unity.

      Se o seu conector suportar a criação de pipeline baseada na interface do utilizador, poderá criar a conexão e o pipeline simultaneamente ao concluir as etapas nesta página. No entanto, se você usar a criação de pipeline baseada em API, deverá criar a conexão no Catalog Explorer antes de concluir as etapas nesta página. Consulte Conectar-se a fontes de ingestão gerenciadas.

    • Se você planeja usar uma conexão existente: você tem USE CONNECTION privilégios ou ALL PRIVILEGES na conexão.

    • No catálogo de destino, você tem privilégios de USE CATALOG.

    • Você tem USE SCHEMA, CREATE TABLE e CREATE VOLUME privilégios num esquema existente ou CREATE SCHEMA privilégios no catálogo de destino.

    • Você tem acesso a uma instância primária do SQL Server. Os recursos de registo de alterações e captura de dados de alterações não são suportados em réplicas de leitura ou instâncias secundárias.

    • Permissões irrestritas para criar clusters ou uma política personalizada (somente API). Uma política personalizada para o gateway deve atender aos seguintes requisitos:

      • Família: Computação de Trabalhos

      • A família de políticas substitui:

        {
          "cluster_type": {
            "type": "fixed",
            "value": "dlt"
          },
          "num_workers": {
            "type": "unlimited",
            "defaultValue": 1,
            "isOptional": true
          },
          "runtime_engine": {
            "type": "fixed",
            "value": "STANDARD",
            "hidden": true
          }
        }
        
      • O Databricks recomenda especificar os menores nós de trabalho possíveis para gateways de ingestão porque eles não afetam o desempenho do gateway. A política de computação a seguir permite que o Azure Databricks dimensione o gateway de ingestão para atender às necessidades de sua carga de trabalho. O requisito mínimo para o nó driver é 8 núcleos para permitir uma extração eficiente e eficaz de dados da sua base de dados de origem.

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

      Para obter mais informações sobre políticas de cluster, consulte Selecionar uma política de computação.

  • Para ingerir a partir do SQL Server, deve primeiro completar os passos em Configurar o Microsoft SQL Server para a ingestão no Azure Databricks.

Criar um gateway e um pipeline de ingestão

Warning

Não pares manualmente o gateway de ingestão. O gateway deve funcionar continuamente para captar alterações antes de os registos de alterações serem truncados na base de dados de origem. Se o gateway ficar parado, as alterações podem perder-se devido à retenção de logs, exigindo uma atualização completa de todas as tabelas afetadas. Parar e reiniciar o gateway também reabastece a VM, o que aumenta o tempo de arranque. Se precisar de resolver problemas com gateway, consulte Troubleshoot SQL Server ingestion ou contacte o Suporte da Databricks.

Interface do usuário do Databricks

  1. Na barra lateral do espaço de trabalho do Azure Databricks, clique em Ingestão de Dados.

  2. Na página Adicionar dados , em Conectores Databricks, clique em SQL Server.

  3. Na página Connection do assistente de ingestão, selecione a ligação que armazena as suas credenciais de acesso SQL Server. Se tiver o privilégio CREATE CONNECTION na metastore, pode clicar em Plus icon. Criar ligação para criar uma nova ligação com os detalhes de autenticação em Criar uma ligação SQL Server.

  4. Clique em Next.

  5. Na página de Configuração de Ingestão , introduza um nome único para o pipeline de ingestão. Esta canal move dados do local intermediário para o destino.

  6. Selecione um catálogo e um esquema para escrever registos de eventos. O registo de eventos contém registos de auditoria, verificações de qualidade de dados, progresso do pipeline e erros. Se tiver USE CATALOG e CREATE SCHEMA privilégios no catálogo, pode clicar no ícone Mais. Criar esquema no menu suspenso para criar um novo esquema.

  7. (Opcional) Defina a atualização automática completa para todas as tabelas para Ligado. Quando a atualização automática está ativada, o pipeline tenta automaticamente corrigir problemas como eventos de limpeza de registos e certos tipos de evolução de esquemas, atualizando completamente a tabela afetada. Se o rastreamento do histórico estiver ativado, uma atualização completa apaga esse histórico.

  8. Introduza um nome exclusivo para a porta de entrada de ingestão. O gateway é um pipeline que extrai as alterações da fonte de dados e as prepara para que o pipeline de ingestão carregue.

  9. Selecione um catálogo e um esquema para o local de Staging. É criado um volume neste local para escalonar os dados extraídos. Se tiver USE CATALOG e CREATE SCHEMA privilégios no catálogo, pode clicar no ícone Mais. Criar esquema no menu suspenso para criar um novo esquema.

  10. Clique em Criar pipeline e continue.

  11. Na página Origem , selecione as tabelas a serem ingeridas. Se selecionares tabelas específicas, podes configurar as definições das tabelas:

    a. (Opcional) Na aba Definições, especifique um nome de destino para cada tabela ingerida. Isto é útil para diferenciar tabelas de destino quando ingere um objeto no mesmo esquema várias vezes. Veja Nomeie uma tabela de destinos.

    a. (Opcional) Altere a configuração padrão de rastreamento de histórico History tracking. Consulte Ativar rastreamento de histórico (SCD tipo 2).

  12. Clica em Próximo, depois em Guardar e continua.

  13. Na página de Destino , selecione um catálogo e um esquema para carregar dados. Se tiver USE CATALOG e CREATE SCHEMA privilégios no catálogo, pode clicar no ícone Mais. Criar esquema no menu suspenso para criar um novo esquema.

  14. Clique em Salvar e continuar.

  15. Na página de Configuração da Base de Dados, clique em Validar para confirmar se a sua fonte está devidamente configurada para a ingestão do Azure Databricks. Quaisquer configurações em falta são retornadas. Para os passos a resolver, clique em Completar configuração. Em seguida, clique em Avançar. Alternativamente, clique em Saltar validação.

  16. (Opcional) Na página de Horários e notificações , clique no ícone Plus. Crie um horário. Defina a frequência para atualizar as tabelas de destino.

  17. (Opcional) Clique no ícone Mais. Adicione uma notificação para definir notificações por email para o sucesso ou fracasso da operação do pipeline, depois clique em Guardar e executar pipeline.

Pacotes de Automação Declarativa

Antes de ingerir usando Pacotes de Automação Declarativa, deve ter acesso a uma ligação existente. Para instruções, veja Criar uma ligação SQL Server.

O catálogo e o esquema de preparo podem ser os mesmos que o catálogo e o esquema de destino. O catálogo de preparo não pode ser um catálogo estrangeiro. Especifique a localização de staging na gateway_definition secção do ficheiro YAML do seu pipeline de bundle.

O gateway de ingestão extrai dados instantâneos e de alteração do banco de dados de origem e os armazena no volume de preparo do Catálogo Unity. Você deve executar o gateway como um pipeline contínuo. Isso ajuda a acomodar quaisquer políticas de retenção de log de alterações que você tenha no banco de dados de origem.

O pipeline de ingestão aplica o instantâneo e os dados de alteração do volume de staging em tabelas de streaming de destino.

Os pacotes podem conter definições YAML de trabalhos e tarefas, são gerenciados usando a CLI do Databricks e podem ser compartilhados e executados em diferentes espaços de trabalho de destino (como desenvolvimento, preparação e produção). Para mais informações, veja O que são os Pacotes de Automação Declarativa?.

  1. Crie um bundle usando a CLI Databricks:

    databricks bundle init
    
  2. Adicione a sua canalização e a configuração de tarefas ao pacote. Consulte Exemplos para um exemplo completo com todas as opções disponíveis.

  3. Implante o pipeline usando a CLI do Databricks:

    databricks bundle deploy
    

Caderno de Notas do Databricks

Atualize a célula Configuration no bloco de notas seguinte com a conexão de origem, o catálogo de destino, o esquema de destino e as tabelas para carregar a partir da origem.

Obter notebook

Terraform

Pode usar o Terraform para implementar e gerir pipelines de ingestão do SQL Server. Para um quadro de exemplos completo, incluindo configurações do Terraform para criar gateways e canais de ingestão, consulte o repositório de exemplos do Lakeflow Connect Terraform no GitHub.

Verificar a ingestão bem-sucedida de dados

A exibição de lista na página de detalhes do pipeline mostra o número de registros processados à medida que os dados são ingeridos. Esses números são atualizados automaticamente.

Verificar a replicação

As Upserted records colunas e Deleted records não são mostradas por padrão. Você pode ativá-los clicando no botão do ícone de configuração de colunas e selecionando-os.

Exemplos

Use estes exemplos para configurar o seu pipeline.

Configuração do pipeline

Pacotes de Automação Declarativa

O conjunto seguinte define um pipeline de gateway, um pipeline de ingestão e uma tarefa agendada. As opções comentadas mostram todas as configurações disponíveis. Atualize as secções variables e targets com os seus detalhes de origem e destino.

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

Caderno de Notas do Databricks

A seguir está uma secção de exemplo Configuration de uma especificação 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"])

Padrões comuns

Para configurações avançadas de pipelines, consulte os Padrões comuns para pipelines de ingestão geridos.

Passos seguintes

Inicia, agenda e configura alertas no seu pipeline. Ver Tarefas comuns de manutenção de oleodutos.

Recursos adicionais