Executar um pipeline CDC integrado em modo contínuo
Aplica-se a : Conectores SaaS
Conectores de banco de dados
Conectores baseados em query
Beta
Esse 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.
O modo contínuo executa um pipeline CDC integrado como uma transmissão sempre ativa em vez de seguir uma programação. Por default, um pipeline CDC integrado é executado no modo Trigger, onde cada atualização extrai e aplica dados de alteração e, em seguida, para. Use o modo contínuo para:
- Ingestão de baixa latência. Os dados de alteração são aplicados às tabelas de transmissão de destino à medida que chegam, normalmente dentro de alguns minutos, em vez de aguardar a próxima atualização programada.
- Fontes com retenção limitada de logs de alteração. Alguns bancos de dados armazenam alterações em logs de transação que podem crescer muito ou ser limpos entre as atualizações. A execução contínua mantém o pipeline sincronizado com a origem, reduzindo o risco de ficar atrás da janela de logs disponível.
Ativar modo contínuo
Para executar um pipeline CDC integrado no modo contínuo, defina continuous como true nas configurações do pipeline. O pipeline usa o modo otimizado para escala por default. Para os passos completos de criação de pipeline, consulte a página de pipeline integrado para o seu conector: Criar um pipeline CDC integrado para SQL Server, Criar um pipeline CDC integrado para MySQL ou Criar um pipeline CDC integrado para Oracle.
- Declarative Automation Bundles
- Databricks notebook
- Databricks CLI
- REST API
resources:
pipelines:
continuous_cdc_pipeline:
name: my-continuous-cdc-pipeline
channel: PREVIEW
continuous: true
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
from databricks.sdk import WorkspaceClient
from databricks.sdk.service.pipelines import (
ConnectorType,
IngestionConfig,
IngestionPipelineDefinition,
TableSpec,
)
w = WorkspaceClient()
pipeline = w.pipelines.create(
name="my-continuous-cdc-pipeline",
channel="PREVIEW",
continuous=True,
catalog="main",
schema="ingestion",
ingestion_definition=IngestionPipelineDefinition(
connection_name="my-sqlserver-connection",
connector_type=ConnectorType.CDC,
objects=[
IngestionConfig(
table=TableSpec(
source_catalog="my_database",
source_schema="dbo",
source_table="customers",
)
)
],
),
)
print(f"Pipeline created: {pipeline.pipeline_id}")
databricks pipelines create --json '{
"name": "my-continuous-cdc-pipeline",
"channel": "PREVIEW",
"continuous": true,
"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"
}
}
]
}
}'
POST /api/2.0/pipelines
{
"name": "my-continuous-cdc-pipeline",
"channel": "PREVIEW",
"continuous": true,
"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"
}
}
]
}
}
Converter um pipeline Trigger em contínuo
Para alternar um pipeline Trigger existente para o modo contínuo:
- Interrompa a atualização atual, se houver uma em execução. Consulte Parar a atualização atual.
- Atualize o pipeline e defina
continuouscomotrue.
A operação de atualização substitui toda a especificação do pipeline, portanto, inclua a definição completa do pipeline, não apenas o campo alterado.
- Declarative Automation Bundles
- Databricks notebook
- Databricks CLI
- REST API
Defina continuous: true no recurso de pipeline e, em seguida, reimplante o pacote:
databricks bundle deploy
from databricks.sdk import WorkspaceClient
w = WorkspaceClient()
w.pipelines.update(
pipeline_id="<pipeline-id>",
name="my-continuous-cdc-pipeline",
channel="PREVIEW",
continuous=True,
catalog="main",
schema="ingestion",
ingestion_definition=existing_ingestion_definition,
)
databricks pipelines update --json '{
"pipeline_id": "<pipeline-id>",
"name": "my-continuous-cdc-pipeline",
"channel": "PREVIEW",
"continuous": true,
"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"
}
}
]
}
}'
PUT /api/2.0/pipelines/<pipeline-id>
{
"pipeline_id": "<pipeline-id>",
"name": "my-continuous-cdc-pipeline",
"channel": "PREVIEW",
"continuous": true,
"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"
}
}
]
}
}
Modos de execução
Pipelines CDC contínuos suportam dois modos de execução:
Mode | Descrição |
|---|---|
Otimizado para escala (default) | Ingere até 500 tabelas rotacionando internamente a ingestão entre elas, sem autoscale agressivo. Use este modo para ingerir um grande número de tabelas em um único pipeline. |
Otimizado para velocidade (Beta) | Executa a transmissão de ingestão para todas as tabelas continuamente para minimizar a latência de ingestão, normalmente dentro de alguns minutos. O modo otimizado de velocidade oferece suporte a até 50 tabelas e usa autoscale agressivo. Use o modo otimizado de velocidade quando a latência for a prioridade máxima. |
O modo de execução é definido pela configuração pipelines.managedIngestion.continuous.runMode do Spark no pipeline. O modo otimizado para escala é o default. Para habilitar o modo otimizado para velocidade, defina runMode como SPEED.
Ativar o modo otimizado de velocidade
Para habilitar o modo otimizado para velocidade, defina a configuração pipelines.managedIngestion.continuous.runMode do Spark como SPEED ao criar o pipeline, além de continuous: true:
- Declarative Automation Bundles
- Databricks notebook
- Databricks CLI
- REST API
resources:
pipelines:
continuous_cdc_pipeline:
name: my-continuous-cdc-pipeline
channel: PREVIEW
continuous: true
catalog: main
schema: ingestion
configuration:
pipelines.managedIngestion.continuous.runMode: SPEED
ingestion_definition:
connection_name: my-sqlserver-connection
connector_type: CDC
objects:
- table:
source_catalog: my_database
source_schema: dbo
source_table: customers
from databricks.sdk import WorkspaceClient
from databricks.sdk.service.pipelines import (
ConnectorType,
IngestionConfig,
IngestionPipelineDefinition,
TableSpec,
)
w = WorkspaceClient()
pipeline = w.pipelines.create(
name="my-continuous-cdc-pipeline",
channel="PREVIEW",
continuous=True,
catalog="main",
schema="ingestion",
configuration={"pipelines.managedIngestion.continuous.runMode": "SPEED"},
ingestion_definition=IngestionPipelineDefinition(
connection_name="my-sqlserver-connection",
connector_type=ConnectorType.CDC,
objects=[
IngestionConfig(
table=TableSpec(
source_catalog="my_database",
source_schema="dbo",
source_table="customers",
)
)
],
),
)
print(f"Pipeline created: {pipeline.pipeline_id}")
databricks pipelines create --json '{
"name": "my-continuous-cdc-pipeline",
"channel": "PREVIEW",
"continuous": true,
"catalog": "main",
"schema": "ingestion",
"configuration": {
"pipelines.managedIngestion.continuous.runMode": "SPEED"
},
"ingestion_definition": {
"connection_name": "my-sqlserver-connection",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_catalog": "my_database",
"source_schema": "dbo",
"source_table": "customers"
}
}
]
}
}'
POST /api/2.0/pipelines
{
"name": "my-continuous-cdc-pipeline",
"channel": "PREVIEW",
"continuous": true,
"catalog": "main",
"schema": "ingestion",
"configuration": {
"pipelines.managedIngestion.continuous.runMode": "SPEED"
},
"ingestion_definition": {
"connection_name": "my-sqlserver-connection",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_catalog": "my_database",
"source_schema": "dbo",
"source_table": "customers"
}
}
]
}
}
Para usar o modo otimizado de escala, remova a configuração pipelines.managedIngestion.continuous.runMode.
Selecione o modo correto
Use a comparação a seguir para escolher entre o modo Trigger e os dois modos de execução contínua:
Capacidade | Acionado | Contínuo (otimizado para escala) | Contínuo (otimizado para velocidade) |
|---|---|---|---|
Número máximo de tabelas suportado | 300 | 500 | 50 |
Compute necessário | Baixo (executa conforme uma programação) | Médio (sempre ativo) | Alto (sempre ativado com autoscale agressivo) |
Carregamento de origem | Baixo (executa conforme uma programação) | Alta | Alta (queries contínuas) |
Consistência | Baixo (risco de rollover do log de alterações) | Alta | Alta |
Os limites nesta tabela aplicam-se a pipelines CDC contínuos. Conectores individuais podem ter limites menores. Consulte a documentação do seu conector.
Parar a atualização atual
Interrompa a atualização do pipeline em execução antes de converter um pipeline para o modo contínuo ou executar um refresh completo seletivo. Substitua <pipeline-id> pelo ID do seu pipeline, que você pode encontrar na UI de pipelines.
- Databricks notebook
- Databricks CLI
- REST API
from databricks.sdk import WorkspaceClient
w = WorkspaceClient()
w.pipelines.stop(pipeline_id="<pipeline-id>")
databricks pipelines stop <pipeline-id>
POST /api/2.0/pipelines/<pipeline-id>/stop
Fully refresh a subset de tabelas
Em um pipeline contínuo, você pode fazer um refresh completo de um subconjunto de tabelas enquanto todas as outras tabelas continuam a ingerir dentro da mesma atualização. Isso é útil quando uma única tabela precisa de um refresh completo (por exemplo, após uma alteração de esquema incompatível) sem interromper o restante do pipeline.
Para realizar um refresh completo seletivo:
- Parar a atualização atual. Consulte Parar a atualização atual.
- Inicie uma nova atualização que lista as tabelas para refresh completo em
full_refresh_selectione definerefresh_selectioncomo o caractere curinga["*"].
- Databricks notebook
- Databricks CLI
- REST API
from databricks.sdk import WorkspaceClient
w = WorkspaceClient()
w.pipelines.start_update(
pipeline_id="<pipeline-id>",
full_refresh_selection=["customers", "orders"],
refresh_selection=["*"],
)
databricks pipelines start-update <pipeline-id> --json '{
"full_refresh_selection": ["customers", "orders"],
"refresh_selection": ["*"]
}'
POST /api/2.0/pipelines/<pipeline-id>/updates
{
"full_refresh_selection": ["customers", "orders"],
"refresh_selection": ["*"]
}
As tabelas em full_refresh_selection são totalmente refreshed, enquanto todas as outras tabelas continuam sendo refreshed dentro da mesma atualização. Após a conclusão do refresh completo, o pipeline retoma automaticamente a ingestão contínua normal para todas as tabelas. Não há necessidade de interromper a atualização e começar uma nova sem full_refresh_selection.
No modo contínuo, uma das seleções de refresh deve incluir o caractere curinga * para que todas as outras tabelas continuem a ingestão enquanto as tabelas selecionadas são totalmente atualizadas. O refresh de apenas um subconjunto de tabelas sem * (um refresh parcial) não é suportado no modo contínuo.
Limitações
O modo contínuo tem as seguintes limitações:
- As atualizações reiniciam para aplicar mudanças de estado. O pipeline usa um mecanismo de cancelar e reiniciar para recarregar o gráfico do pipeline ou aplicar mudanças de estado, como mudanças de esquema.
- O refresh completo pode exigir várias reinicializações. Um refresh completo pode exigir várias reinicializações do pipeline para ser concluído, porque o Snapshot de origem é preparado de forma assíncrona.