Criar um pipeline CDC integrado para Oracle
Beta
Este recurso está em Beta. Os administradores do workspace podem controlar o acesso a esse recurso na página Pré-visualizações . Consulte Gerenciar prévias do Databricks.
Um pipeline de CDC integrado ingere dados de alteração do Oracle para o Databricks usando um único pipeline. O conector CDC integrado combina extração e aplicação em uma única atualização de pipeline.
O conector CDC integrado do Oracle usa LogMiner no modo de transação não confirmada para ler as alterações dos redo logs online e archive logs.
Requisitos
-
O Workspace está habilitado para o Unity Catalog.
-
Para criar uma conexão: É preciso ter
CREATE CONNECTIONprivilégios no metastore. Consulte Gerenciar privilégios no Unity Catalog.Se o seu conector suporta a criação de pipeline baseada na IU, você pode criar a conexão e o pipeline ao mesmo tempo concluindo os passos desta página. No entanto, se você usa a criação de pipeline baseada em API, você deve criar a conexão no Catalog Explorer antes de concluir os passos nesta página. Consulte Conectar-se a fontes de ingestão gerenciadas.
-
Se você planeja usar uma conexão existente: você tem os privilégios
USE CONNECTIONouALL PRIVILEGESna conexão. -
O usuário possui
USE CATALOGprivilégios no catálogo de destino. -
Ter os privilégios
USE SCHEMA,CREATE TABLEeCREATE VOLUMEem um esquema existente ou os privilégiosCREATE SCHEMAno catálogo de destino. -
Seu workspace deve ter o recurso de conector CDC integrado habilitado. Entre em contato com sua equipe de conta da Databricks.
-
Você concluiu a configuração do banco de dados de origem Oracle. Consulte Configurar o Oracle para ingestão no Databricks.
-
Você tem as seguintes permissões:
CREATE CONNECTIONno metastore (se estiver criando uma nova conexão do Unity Catalog), ouUSE CONNECTIONem uma conexão existente.USE CATALOGno catálogo de destinos.USE SCHEMAeCREATE TABLEno esquema de destino.CREATE VOLUMEno esquema de destino ou no esquema especificado emdata_staging_options.
Para bancos de dados multi-tenant Oracle, o usuário de conexão deve ser um usuário comum no CDB$ROOT. Para obter detalhes, consulte Criar o usuário de replicação.
Compute requirements
Um pipeline CDC integrado é executado em compute clássico ou serverless:
- Compute clássico : O plano de compute clássico é executado em sua VPC ou VNet do workspace do Databricks e deve ser capaz de alcançar sua instância Oracle pela rede. Qualquer caminho de rede que permite que o plano de compute alcance o banco de dados tem suporte, incluindo emparelhamento de VPC ou VNet, endpoints públicos e, para Oracle on-premises, AWS Direct Connect, Azure ExpressRoute ou VPN.
- Compute serverless : Configure a conectividade de rede serverless entre o compute serverless do Databricks e seu banco de dados de origem. Fontes on-premises exigem um caminho de rede através da saída serverless configurada (por exemplo, um gateway de trânsito ou VNet emparelhada com ExpressRoute ou VPN).
Para o compute clássico, é possível usar permissões irrestritas de criação de cluster ou uma política de cluster personalizada com cluster_type fixado em dlt, runtime_engine fixado em STANDARD, e pelo menos 8 núcleos recomendados para uma extração eficiente.
Criar conexão do Unity Catalog para Oracle
Crie uma conexão do Unity Catalog com o Oracle antes de criar um pipeline. Consulte Criar uma conexão Oracle.
Criar um pipeline CDC integrado
Crie um pipeline de CDC integrado usando a interface do usuário de ingestão de dados, a API REST, a CLI do Databricks, notebooks ou Declarative Automation Bundles.
Toda solicitação de criação de pipeline programático deve incluir "channel": "PREVIEW". Ao usar a interface do usuário, o Databricks define o canal para você.
Para pipelines CDC integrados do Oracle, source_catalog mapeia para o nome do serviço Oracle. Para bancos de dados multi-tenant, este deve ser o nome do serviço CDB$ROOT.
- Databricks UI
- Declarative Automation Bundles
- Databricks notebook
- Databricks CLI
- REST API
-
Na barra lateral, clique em ingestão de dados , depois selecione Oracle como o tipo de origem.

-
Selecione uma conexão para usar. Escolha uma conexão existente do Unity Catalog ou crie uma.


-
Forneça um nome para o pipeline e um local para o log de eventos. O local do log de eventos é onde o Databricks armazena dados de staging e os metadados usados para realizar CDC.

-
Clique em Avançar . O Databricks faz o provisionamento do compute e cria o pipeline. O passo pode levar algum tempo e exibir
Waiting for resources. [[ ## completed ##]] Quando for concluído, selecione as tabelas de origem para ingestão.
-
Selecione o esquema de destino onde o pipeline grava os dados capturados da origem. O pipeline cria automaticamente tabelas com os mesmos nomes da origem no esquema selecionado.

-
Clique em Validar e aguarde a conclusão bem-sucedida da validação.

-
Defina como programar o pipeline. O pipeline é executado enquanto houver dados disponíveis, para após atingir um estado parado e retoma do mesmo ponto no próximo Trigger. [[ ## completed ##]]

-
Revise o pipeline. A view de lista mostra os fluxos e estatísticas sobre os dados replicados.

-
Para verificar o que o pipeline está fazendo, e particularmente para revisar avisos ou mensagens de erro quando uma atualização falha, abra o painel Event logs à direita.

O pipeline agora está configurado e em execução. Você pode fazer query nas tabelas que o pipeline cria no esquema de destino e tratá-las como tabelas Bronze na arquitetura medallion.
Defina o recurso de pipeline em um arquivo de pacote (por exemplo, resources/oracle_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:
oracle_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: 'ORCL'
source_schema: 'HR'
source_table: 'EMPLOYEES'
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
destination_table: 'employees'
table_configuration:
scd_type: 'SCD_TYPE_1'
Para executar o pipeline em um agendamento, defina um Job que acione o pipeline. Como cada etapa de extração executa por pelo menos 10 minutos, um intervalo de 60 minutos ou mais é um bom ponto de partida:
resources:
jobs:
oracle_integrated_cdc_job:
name: '${var.pipeline_name}-job'
tasks:
- task_key: 'cdc_ingestion'
pipeline_task:
pipeline_id: ${resources.pipelines.oracle_integrated_cdc_pipeline.id}
schedule:
quartz_cron_expression: '0 0 * * * ?'
timezone_id: 'UTC'
Implante o pacote com a CLI do Databricks:
databricks bundle deploy
databricks bundle run oracle_integrated_cdc_job
Para obter mais informações, consulte O que são Pacotes de Automação Declarativa?.
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="<oracle-connection-name>",
connector_type=ConnectorType.CDC,
objects=[
IngestionConfig(
table=TableSpec(
source_catalog="<oracle-service-name>",
source_schema="<oracle-schema>",
source_table="<oracle-table>",
destination_catalog="<destination-catalog>",
destination_schema="<destination-schema>",
)
)
],
),
)
print(f"Pipeline created: {pipeline.pipeline_id}")
databricks pipelines create --json '{
"name": "<pipeline-name>",
"channel": "PREVIEW",
"catalog": "<destination-catalog>",
"schema": "<destination-schema>",
"ingestion_definition": {
"connection_name": "<oracle-connection-name>",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_catalog": "<oracle-service-name>",
"source_schema": "<oracle-schema>",
"source_table": "<oracle-table>"
}
}
]
}
}'
POST /api/2.0/pipelines
{
"name": "my-oracle-integrated-pipeline",
"channel": "PREVIEW",
"catalog": "main",
"schema": "ingestion",
"ingestion_definition": {
"connection_name": "my-oracle-connection",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_catalog": "ORCL",
"source_schema": "HR",
"source_table": "EMPLOYEES",
"table_configuration": {
"scd_type": "SCD_TYPE_1"
}
}
}
],
"data_staging_options": {
"catalog_name": "main",
"schema_name": "ingestion_staging"
}
}
}
Para replicar todas as tabelas em um esquema de origem, utilize um objeto schema em vez de objetos table individuais:
POST /api/2.0/pipelines
{
"name": "my-oracle-schema-pipeline",
"channel": "PREVIEW",
"catalog": "main",
"schema": "ingestion",
"ingestion_definition": {
"connection_name": "my-oracle-connection",
"connector_type": "CDC",
"objects": [
{
"schema": {
"source_catalog": "ORCL",
"source_schema": "HR",
"destination_catalog": "main",
"destination_schema": "ingestion"
}
}
]
}
}
Para iniciar uma atualização do pipeline:
POST /api/2.0/pipelines/<pipeline-id>/updates
{
"full_refresh": false
}
Programar atualizações recorrentes
Pipelines de CDC integrados operam apenas no modo triggered. Para ingerir dados em um agendamento recorrente, crie uma tarefa Lakeflow Jobs que executa o pipeline. Cada atualização dura aproximadamente 30 minutos e pode não concluir o processamento de todo o backlog de alterações em uma única atualização. Programe os pipelines com frequência suficiente para que as atualizações subsequentes possam acompanhar. Um ponto de partida de 60 minutos funciona bem para a maioria das cargas de trabalho.
Referência de configuração
Parâmetros do pipeline
Parâmetro | Tipo | Descrição |
|---|---|---|
| string | Um nome para o pipeline. |
| string | Deve ser |
| Booleana | Opcional. O padrão é |
| string | O catálogo de destinos default. |
| string | O esquema de destino default. |
| string | A conexão do Unity Catalog com o Oracle. |
| string | Deve ser |
| matriz | A lista de tabelas ou esquemas para ingestão. |
| objeto | Opcional. O catálogo e o esquema onde o pipeline cria o volume de preparo. Usa por padrão o esquema de destino do pipeline. |
Especificação da tabela
Parâmetro | Obrigatório | Descrição |
|---|---|---|
| Sim | O nome do serviço Oracle. Para bancos de dados multi-tenant, use o nome do serviço |
| Sim | O esquema Oracle (normalmente o proprietário da tabela). |
| Sim | O nome da tabela Oracle. |
| Não | Catálogo de destinos. O padrão é o |
| Não | O esquema de destino. O padrão é o |
| Não | O nome da tabela de destino. O padrão é |
Configuração da tabela
Parâmetro | Padrão | Descrição |
|---|---|---|
| Detectado automaticamente | As colunas que identificam cada linha. Detectado automaticamente da chave primária de origem se não especificada. |
|
|
|
| Detectado automaticamente | As colunas usadas para ordenar eventos de CDC. |
Para mapeamentos de tipos de dados do Oracle, consulte Mapeamentos de tipos de dados.
Sensibilidade a maiúsculas e minúsculas para identificadores Oracle
Oracle armazena identificadores sem aspas em maiúsculas. Ao especificar source_catalog, source_schema, source_table e primary_keys na configuração do pipeline, o uso de maiúsculas e minúsculas deve corresponder a como o Oracle armazena o identificador. Para a maioria dos bancos de dados, isso significa usar letras maiúsculas. Se um identificador foi criado com aspas duplas que preservou um caso diferente, utilize esse caso exato.
Monitorar o pipeline
Depois de criar e iniciar um pipeline CDC integrado, monitore seu status usando o seguinte:
-
Interface do usuário do Databricks. Abra o pipeline na seção Pipelines para exibir o status de atualização, as métricas de ingestão por tabela e a linhagem.
-
API REST.
TextGET /api/2.0/pipelines/<pipeline-id> -
API de Eventos.
TextGET /api/2.0/pipelines/<pipeline-id>/events
A view em lista na página de detalhes do pipeline mostra o número de registros processados à medida que os dados são ingeridos. [[ ## completed ##]] Estes números fazem refresh automaticamente.

A primeira atualização do pipeline realiza um snapshot completo de todas as tabelas selecionadas, o que pode levar mais tempo do que as atualizações incrementais. Para tabelas grandes, o snapshot inicial pode exigir várias atualizações agendadas para ser concluído.
Você pode fazer query dos dados ingeridos no Unity Catalog.

Para comportamento de refresh completo e refresh automático completo, consulte Fazer refresh completo das tabelas de destino.
Pipelines CDC integrados têm o autoscale vertical ativado por default. Se uma atualização de pipeline falhar devido a uma condição de falta de memória, a próxima atualização provisiona automaticamente um driver maior.
Limitações
Limitações gerais
- Beta. O conector de CDC integrado e o conector Oracle exigem habilitação no nível do workspace. Entre em contato com sua equipe de account da Databricks.
- Apenas modo acionado. Pipelines CDC integrados não são compatíveis com execução contínua (sempre ativa). Programar pipelines com uma tarefa de LakeFlow Jobs.
- O canal deve ser
PREVIEW. As especificações programáticas do pipeline devem incluir"channel": "PREVIEW". - Máximo recomendado de aproximadamente 500 tabelas por pipeline de ingestão.
- Pipelines de CDC integrados ainda não oferecem suporte a alterações de esquema (operações DDL).
- Snapshot inicial pode abranger várias atualizações para tabelas grandes.
- Cada atualização tem uma execução de aproximadamente 30 minutos. O pipeline não processa necessariamente todo o backlog de alterações em uma única atualização. As atualizações agendadas subsequentes retomam o processamento de onde a atualização anterior parou. Não é possível configurar este Runtime.
- Conexão e tipo de conector são imutáveis após a criação do pipeline.
Limitações específicas do Oracle
- Implantações do Oracle não suportadas: Oracle RAC, Exadata na configuração RAC, Physical Standby, Oracle Autonomous Databases e instâncias de banco de dados Amazon RDS de multi-tenant.
- Tipos de dados não
XMLsuportados:,JSONe tipos de dados espaciais. - Tabelas Ignoradas pelo LogMiner : o LogMiner ignora qualquer tabela que contenha
BFILE, tabelas aninhadas, colunas de identidade, colunas de validade temporal, colunasPKREFou colunasPKOID. Consulte limitações do LogMiner. - Comprimento do identificador: Nomes de tabelas e colunas não podem exceder 30 caracteres.
- Recursos pós-12.2 : o conector não oferece suporte a tipos de dados e recursos adicionados após o Oracle Database 12c Release 2, incluindo
BOOLEAN,VECTOReJSON.
Solução de problemas
Se uma atualização de pipeline falhar:
- Revise o log de eventos do pipeline na UI do Databricks ou por meio de
GET /api/2.0/pipelines/<pipeline-id>/events. - Teste a conexão do Unity Catalog do Explorador de Catálogos para confirmar se o Oracle está acessível.
- Confirme se o modo de log de arquivo e o registro suplementar estão habilitados. Consulte o Passo 1: Verificar o modo de log de arquivo e a retenção de logs.
- Verifique se o usuário de replicação possui os privilégios concedidos por
DBX_ORACLE_SETUP_UTIL.GRANT_PERMISSIONS. Consulte requisitos de usuário do Oracle database. - Para bancos de dados multi-tenant, confirme se o usuário é um usuário comum em
CDB$ROOTe quesource_catalogé o nome do serviçoCDB$ROOT. - Verifique se a especificação do seu pipeline inclui
"channel": "PREVIEW".
Se o Oracle limpar os logs de arquivo morto antes que o pipeline possa processá-los, execute um refresh completo nas tabelas afetadas.