Crie um pipeline CDC integrado para o SQL Server

Importante

Este recurso está em versão Beta. Os administradores do espaço de trabalho podem controlar o acesso a esse recurso na página Visualizações . Ver Gerir as pré-visualizações de Azure Databricks.

Um pipeline CDC integrado ingere dados de alteração do SQL Server para o Azure Databricks usando um único pipeline. Ao contrário da arquitetura padrão baseada em gateways, que requer um gateway de ingestão separado e um pipeline de ingestão separado, um pipeline integrado de CDC executa as etapas de extração e aplicação numa única atualização do pipeline.

Quando usar o conector CDC integrado

A tabela seguinte compara pipelines de CDC integrados com a arquitetura padrão baseada em gateway:

Feature CDC padrão (baseado em gateway) CDC integrado
Número de oleodutos Dois (porta de entrada para ingestão e fluxo de ingestão) One (canalização unificada)
Configuração Crie um gateway, depois crie um pipeline de ingestão que faça referência ao ID do gateway Crie um único pipeline que faça referência a uma ligação ao Unity Catalog
Modo Gateway A porta de entrada funciona continuamente O pipeline integra a extração em cada atualização
Referência de conexão ingestion_gateway_id connection_name (uma conexão ao Unity Catalog)
Tipo de conector Implícito Explícito: connector_type: CDC
Volume de encenação O gateway gere o volume de preparação internamente Configuras o volume de staging através de data_staging_options. O pipeline cria automaticamente um se não for especificado.

Para configuração da base de dados de origem, veja Configure Microsoft SQL Server para ingestão em Azure Databricks. A mesma configuração de origem aplica-se a ambas as arquiteturas.

Como funciona um pipeline CDC integrado

Cada atualização do pipeline executa duas etapas em sequência:

  1. Extração. O pipeline liga-se à base de dados de origem através da ligação Unity Catalog. Na primeira execução ou numa atualização completa, captura uma captura inicial. Nas execuções subsequentes, capta alterações incrementais (inserções, atualizações e eliminações) utilizando o mecanismo de acompanhamento de alterações integrado na base de dados. O pipeline escreve os dados extraídos num volume de staging do Unity Catalog.
  2. Aplicação. O pipeline lê do volume de preparação e aplica alterações às tabelas de transmissão de destino no Unity Catalog. As operações de fusão utilizam as chaves primárias configuradas e o tipo SCD. O pipeline garante uma semântica exata uma vez.

Cada atualização do pipeline extrai alterações e depois para automaticamente depois de ter alcançado a fonte, limitada por um tempo máximo de execução. Para mais detalhes, consulte Fecho inteligente para oleodutos CDC integrados. Para ingerir dados de forma recorrente, agende o pipeline usando uma tarefa Lakeflow Jobs .

Requirements

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

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

  • O seu espaço de trabalho deve ter ativada a funcionalidade integrada do conector CDC. Contacte a sua equipa de contas Azure Databricks.
  • Tens acesso à instância primária SQL Server. O conector CDC integrado não suporta réplicas de leitura, instâncias em espera ou instâncias secundárias.
  • Concluiu a configuração do código-fonte do SQL Server. Consulte Configurar o Microsoft SQL Server para ingestão no Azure Databricks.
  • Tem as seguintes permissões:
    • CREATE CONNECTION na metastore (se criares uma nova ligação ao Unity Catalog), ou USE CONNECTION numa ligação já existente.
    • USE CATALOG no catálogo de destino.
    • USE SCHEMA e CREATE TABLE no esquema de destino.
    • CREATE VOLUME no esquema de destino, ou no esquema especificado em data_staging_options. É necessário um volume de preparação, mesmo que data_staging_options não esteja definido, porque o pipeline cria automaticamente um no esquema de destino.

Requisitos de computação

Um pipeline integrado de CDC é executado em computação clássica ou serverless:

  • Computação clássica: O plano de computação clássico corre no seu VPC ou VNet do seu espaço de trabalho Azure Databricks e deve conseguir aceder à sua instância do SQL Server através da rede. Qualquer caminho de rede que permita ao plano de cálculo aceder à base de dados é suportado, incluindo peering VPC ou VNet, endpoints públicos e, para SQL Server local, AWS Direct Connect, Azure ExpressRoute ou VPN.
  • Computação sem servidor: Configure a conectividade de rede sem servidor entre a computação sem servidor do Azure Databricks e a sua base de dados de origem. As fontes on-premises requerem um caminho de rede através da saída serverless configurada (por exemplo, um transit gateway ou VNet peered com ExpressRoute ou VPN).

Para computação clássica, pode usar permissões de criação de cluster irrestritas ou uma política personalizada de cluster com cluster_type fixo para dlt, runtime_engine fixo para STANDARD, e pelo menos 8 núcleos recomendados para extração eficiente.

Criar uma ligação ao Unity Catalog para SQL Server

Crie uma ligação Unity Catalog ao SQL Server antes de criar um pipeline. Veja Criar uma ligação SQL Server.

Criar uma canalização CDC integrada

Crie pipelines de CDC integrados usando a API, a CLI da Databricks, notebooks ou bundles de automação declarativa. A criação de interfaces ainda não está disponível.

Importante

Todos os pedidos de criação de pipeline devem incluir "channel": "PREVIEW".

Pacotes de Automação Declarativa

Defina o recurso pipeline num ficheiro bundle (por exemplo, resources/integrated_cdc_pipeline.yml):

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'

Para executar o pipeline de forma agendada, defina uma tarefa (por exemplo, resources/integrated_cdc_job.yml) que acione o pipeline. Como cada etapa de extração dura pelo menos 10 minutos, um intervalo de 60 minutos ou mais é um bom ponto de partida:

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'

Implemente o bundle com a CLI Databricks:

databricks bundle deploy
databricks bundle run integrated_cdc_job

Para mais informações, veja O que são os Pacotes de Automação Declarativa?.

Caderno de Notas do Databricks

from databricks.sdk import WorkspaceClient
from databricks.sdk.service.pipelines import (
    ConnectorType,
    IngestionConfig,
    IngestionPipelineDefinition,
    TableSpec,
)

w = WorkspaceClient()

pipeline = w.pipelines.create(
    name="<pipeline-name>",
    channel="PREVIEW",
    catalog="<destination-catalog>",
    schema="<destination-schema>",
    ingestion_definition=IngestionPipelineDefinition(
        connection_name="<unity-catalog-connection-name>",
        connector_type=ConnectorType.CDC,
        objects=[
            IngestionConfig(
                table=TableSpec(
                    source_catalog="<source-database>",
                    source_schema="<source-schema>",
                    source_table="<source-table>",
                    destination_catalog="<destination-catalog>",
                    destination_schema="<destination-schema>",
                )
            )
        ],
    ),
)

print(f"Pipeline created: {pipeline.pipeline_id}")

CLI do Databricks

databricks pipelines create --json '{
  "name": "<pipeline-name>",
  "channel": "PREVIEW",
  "catalog": "<destination-catalog>",
  "schema": "<destination-schema>",
  "ingestion_definition": {
    "connection_name": "<unity-catalog-connection-name>",
    "connector_type": "CDC",
    "objects": [
      {
        "table": {
          "source_catalog": "<source-database>",
          "source_schema": "<source-schema>",
          "source_table": "<source-table>"
        }
      }
    ]
  }
}'

API REST

O exemplo seguinte replica duas tabelas de uma base de dados do SQL Server. A tabela customers usa SCD Tipo 1, e a tabela orders usa SCD Tipo 2 (que requer SQL Server CDC na fonte). Ambos herdam o destino de nível superior main.ingestion. O exemplo omite serverless, que por defeito é false (computação clássica). Adicione "serverless": true para executar em computação sem servidor em vez disso.

POST /api/2.0/pipelines

{
  "name": "my-integrated-cdc-pipeline",
  "channel": "PREVIEW",
  "catalog": "main",
  "schema": "ingestion",
  "ingestion_definition": {
    "connection_name": "my-sqlserver-connection",
    "connector_type": "CDC",
    "objects": [
      {
        "table": {
          "source_catalog": "my_database",
          "source_schema": "dbo",
          "source_table": "customers",
          "table_configuration": {
            "scd_type": "SCD_TYPE_1"
          }
        }
      },
      {
        "table": {
          "source_catalog": "my_database",
          "source_schema": "dbo",
          "source_table": "orders",
          "table_configuration": {
            "scd_type": "SCD_TYPE_2"
          }
        }
      }
    ],
    "data_staging_options": {
      "catalog_name": "main",
      "schema_name": "ingestion_staging"
    }
  }
}

Para replicar todas as tabelas num esquema de origem, use um schema objeto em vez de objetos individuais table . O pipeline ignora tabelas sem CDC ou controlo de alterações ativados na origem.

POST /api/2.0/pipelines

{
  "name": "my-integrated-cdc-schema-pipeline",
  "channel": "PREVIEW",
  "catalog": "main",
  "schema": "ingestion",
  "ingestion_definition": {
    "connection_name": "my-sqlserver-connection",
    "connector_type": "CDC",
    "objects": [
      {
        "schema": {
          "source_catalog": "my_database",
          "source_schema": "dbo",
          "destination_catalog": "main",
          "destination_schema": "ingestion"
        }
      }
    ]
  }
}

Para iniciar uma atualização do pipeline:

POST /api/2.0/pipelines/<pipeline-id>/updates

{
  "full_refresh": false
}

Agendar atualizações recorrentes

Por predefinição, os pipelines CDC integrados são executados em modo acionado. Para execução contínua, consulte Executar um pipeline CDC integrado em modo contínuo. Para ingerir dados num cronograma recorrente, crie uma tarefa Lakeflow Jobs que execute o pipeline. A duração da atualização varia consoante a quantidade de dados de alterações existentes na origem, e uma grande acumulação de pendentes pode não ficar concluída numa única atualização (ver Encerramento inteligente para pipelines CDC integrados). Programa os pipelines com frequência suficiente para que as atualizações subsequentes consigam recuperar o atraso. Um ponto de partida de 60 minutos funciona bem para a maioria das cargas de trabalho. Se um gatilho for ativado enquanto uma atualização anterior ainda está a correr, a nova atualização fica em fila.

Referência de configuração

Parâmetros do pipeline

Parâmetro Tipo Description
name cadeia (de caracteres) Um nome para o oleoduto.
channel cadeia (de caracteres) Deve ser PREVIEW.
serverless booleano Optional. O valor padrão é false. Defina como true para computação serverless ou false para computação clássica. A computação sem servidor requer ligação em rede sem servidor à sua base de dados de origem.
catalog cadeia (de caracteres) O catálogo de destino padrão. Usado quando um per-tabela destination_catalog não é especificado.
schema cadeia (de caracteres) O esquema de destino por defeito. Usado quando um per-tabela destination_schema não é especificado.
ingestion_definition.connection_name cadeia (de caracteres) A ligação do Unity Catalog à base de dados de origem.
ingestion_definition.connector_type cadeia (de caracteres) Deve ser CDC.
ingestion_definition.objects matriz A lista de tabelas ou esquemas a ingerir.
ingestion_definition.data_staging_options objecto Optional. O catálogo e o esquema onde o pipeline cria o volume de preparação. Predefinido para o esquema de destino do pipeline.

Especificação da tabela

Parâmetro Required Description
source_catalog Yes O nome da base de dados de origem.
source_schema Yes O nome do esquema de origem.
source_table Yes O nome da tabela de origem.
destination_catalog No O catálogo de destinos. Por predefinição, utiliza o catalog da pipeline.
destination_schema No O esquema de destino. Por predefinição, utiliza o schema da pipeline.
destination_table No O nome da tabela de destino. O valor padrão é source_table.

Configuração da tabela

Parâmetro Predefinido Description
primary_keys Detetação automática As colunas que identificam cada linha. Detetado automaticamente a partir da chave primária de origem, caso não seja especificada.
scd_type SCD_TYPE_1 SCD_TYPE_1 mantém apenas a versão mais recente. SCD_TYPE_2 mantém o histórico completo e exige SQL Server CDC na fonte. O SCD Tipo 2 não é suportado com rastreamento de alterações.
sequence_by Detetação automática As colunas eram usadas para ordenar os eventos do CDC. Detetado automaticamente com base no mecanismo CDC da origem, caso não seja especificado.

Para mapeamentos de tipos de dados do SQL Server, consulte referência do conector do SQL Server. Pipelines CDC integrados suportam alargamento automático de tipos: quando um tipo de coluna de origem é alargado (por exemplo, INT para BIGINT), a tabela de destino adapta-se automaticamente.

Monitorizar o oleoduto

Depois de criar e iniciar um pipeline CDC integrado, monitorize o seu estado usando o seguinte:

  • Azure Databricks UI. Abra o pipeline na secção Pipelines para ver o estado da atualização, as métricas de ingestão de cada tabela e a linhagem de dados.

  • REST API.

    GET /api/2.0/pipelines/<pipeline-id>
    
  • API de eventos.

    GET /api/2.0/pipelines/<pipeline-id>/events
    

A primeira atualização do pipeline realiza um snapshot completo de todas as tabelas selecionadas, o que pode demorar mais do que as atualizações incrementais. Para tabelas grandes, a captura inicial pode exigir várias atualizações agendadas para ficar concluída. Cada atualização subsequente retoma de onde a anterior terminou.

Para verificar a ingestão:

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

Para atualização completa e comportamento automático de atualização completa, consulte Tabelas alvo de atualização completa.

As pipelines CDC integradas têm o dimensionamento automático vertical ativado por predefinição. Se uma atualização do pipeline falhar devido a uma condição de falta de memória, a atualização seguinte provisiona automaticamente um driver maior. Para sobrepor este comportamento, use uma política de cluster personalizada.

Limitações

  • Beta. O conector CDC integrado requer ativação ao nível do espaço de trabalho. Contacte a sua equipa de contas Azure Databricks.
  • Ativado por predefinição. Por defeito, os pipelines de CDC integrados são executados no modo acionado; agende-os através de uma tarefa do Lakeflow Jobs. O modo contínuo está disponível em Beta. Veja : Executar um pipeline CDC integrado em modo contínuo.
  • Criação exclusivamente através da API. A criação de pipelines está disponível por meio da API REST, da CLI da Databricks, de notebooks e de pacotes de automatização declarativa. A criação de interfaces ainda não é suportada.
  • O canal deve ser PREVIEW. As especificações do pipeline devem incluir "channel": "PREVIEW".
  • A ligação e o tipo de conector são imutáveis. connection_name e connector_type não podem ser alteradas após a criação do pipeline. Para alterar a origem, crie uma nova canalização.
  • Máximo recomendado de 300 tabelas por pipeline.
  • Apenas nas instâncias principais. O conector CDC integrado não suporta réplicas de leitura, instâncias em espera ou instâncias secundárias.
  • Tabelas sem chaves primárias. O pipeline trata todas as colunas que não sejam LOB como uma chave composta. As linhas duplicadas podem colapsar para uma única linha, a menos que atives o SCD Tipo 2.
  • O snapshot inicial pode abranger várias atualizações. Para tabelas grandes, o snapshot inicial pode não terminar numa única atualização. As atualizações agendadas subsequentes são retomadas a partir do ponto onde a atualização anterior terminou.
  • O tempo de execução da atualização é gerido automaticamente: O fecho inteligente determina quando cada atualização termina. Uma atualização termina depois de alcançar a fonte, limitada por um tempo máximo de execução. Consulte Encerramento inteligente para pipelines de CDC integrados. Não podes configurar o tempo de execução mínimo ou máximo. Um grande atraso de alterações pode abranger várias atualizações. As atualizações agendadas subsequentes são retomadas a partir do ponto onde a atualização anterior terminou.
  • A purga de registos requer uma atualização completa. Se o SQL Server eliminar os registos do controlo de alterações ou os registos CDC antes de o pipeline os processar, efetue uma atualização completa das tabelas afetadas. O pipeline deteta esta condição e apresenta um erro no registo de eventos.

Resolução de Problemas

Note

Alguns códigos de erro usam o INGESTION_GATEWAY_ prefixo. Esta é uma convenção de nomenclatura legada e não indica que seja necessário um gateway de ingestão separado.

Erro Cause Resolução
NOT_IN_DEFAULT_PUBLISHING_MODE O pipeline não está em Modo de Publicação Direta. O Modo de Publicação Direta é definido automaticamente para pipelines CDC integrados. Se vir este erro, recrie o pipeline.
INGESTION_GATEWAY_CDC_NOT_ENABLED O CDC ou o acompanhamento de alterações não estão ativados numa ou mais tabelas de origem. Ative o CDC ou o acompanhamento de alterações nas tabelas afetadas. Consulte Configurar o Microsoft SQL Server para ingestão no Azure Databricks.
INGESTION_GATEWAY_MISSING_TABLE_IN_SOURCE A tabela de origem especificada não existe ou foi retirada. Verifique se a tabela existe e que o utilizador da ligação tem acesso.
INGESTION_GATEWAY_SOURCE_SCHEMA_MISSING_ENTITY O esquema fonte não existe. Verifique se o esquema existe na base de dados de origem.
UNSUPPORTED_SOURCE_TYPE_FOR_CDC_CONNECTOR O tipo de base de dados de origem não é suportado. O conector CDC integrado suporta SQL Server e Oracle.
SOURCE_TABLE_REQUIRED A especificação da tabela está em falta source_table. Adicione source_table a cada especificação de tabela no objects array.
Integrated CDC connector is disabled O sinalizador de funcionalidade do espaço de trabalho não está ativado. Contacte a sua equipa de contas Azure Databricks para ativar o conector CDC integrado no seu espaço de trabalho.

Se encontrar um problema não abordado aqui:

  1. Revise o registo de eventos do pipeline na interface Azure Databricks ou através de GET /api/2.0/pipelines/<pipeline-id>/events.
  2. Teste a ligação ao Unity Catalog a partir do Explorador de Catálogos para confirmar que a fonte é acessível.
  3. Confirme que o controlo de alterações ou o CDC está ativado na base de dados de origem e nas tabelas.
  4. Verifique se o utilizador da base de dados tem as permissões do SQL Server listadas em Requisitos do utilizador da base de dados do Microsoft SQL Server.
  5. Verifica se a especificação do teu pipeline inclui "channel": "PREVIEW".

Recursos adicionais