Crie um pipeline de ingestão baseado em consultas.
Esta página mostra como criar um pipeline de ingestão baseado em consultas no LakeFlow Connect.
Requisitos
Antes de criar um pipeline de ingestão baseado em consulta, você precisa atender aos seguintes requisitos:
- Unity Catalog está habilitado para seu workspace Databricks .
- Seu ambiente compute serverless permite conectividade de rede com o banco de dados de origem. Consulte a seção Redes e as recomendações de rede para a Lakehouse Federation.
- Para ingestão de conexão externa : Você tem uma conexão existente com o banco de dados de origem ou privilégios
CREATE CONNECTIONno metastore. Consulte Conectar para gerenciar fontes de ingestão. - Para ingestão de catálogos externos : Você precisa ter um catálogo externo já registrado na Federação Lakehouse ou ter privilégios para criar um.
- Você tem privilégios
CREATEeUSE SCHEMAno catálogo e esquema de destino.
Opção 1: Ingestão de conexão estrangeira
Utilize essa abordagem quando você tiver uma conexão que armazena credenciais de autenticação para o banco de dados de origem. As fontes de dados suportadas incluem Oracle, Teradata, SQL Server, MySQL, MariaDB e PostgreSQL.
- Databricks UI
- Declarative Automation Bundles
A IU do Databricks implanta pipelines baseados em query para compute serverless. Para implantar em compute clássico, em vez disso, consulte a **tab Pacotes de Automação Declarativa**.
-
Na barra lateral workspace Databricks , clique em inserção de dados .
-
Na página Adicionar dados , em Conectores do Databricks , clique na sua fonte (por exemplo, Oracle ou SQL Server ). O assistente de ingestão é aberto.
-
Na página **Pipeline de ingestão**, insira um nome para o pipeline.
-
Em Catálogo de destino , selecione um catálogo Unity Catalog para armazenar os dados recebidos.
-
Selecione a conexão do Unity Catalog que armazena as credenciais necessárias para acessar o banco de dados de origem.
Se não houver nenhuma conexão existente, clique em Criar conexão e insira os detalhes da conexão. Você deve ter privilégios
CREATE CONNECTIONno metastore. -
Clique em Criar pipeline e continue .
-
Na página Origem , selecione os esquemas e tabelas a serem importados.
-
Para cada tabela, especifique a coluna do cursor . Deve ser uma única coluna com valores que aumentam monotonicamente (por exemplo,
updated_atourow_id). Se você não selecionar uma coluna de cursor que aumente monotonicamente, o conector realizará uma carga completa em cada execução. -
Opcionalmente, altere a configuração default da história acompanhamento. Para mais informações, consulte Habilitar história acompanhamento (SCD tipo 2).
-
Clique em Avançar .
-
Na página Destino , selecione o catálogo e o esquema do Unity Catalog nos quais deseja gravar.
Se não quiser usar um esquema existente, clique em Criar esquema . Você deve ter privilégios
USE CATALOGeCREATE SCHEMAno catálogo pai. -
Clique em Salvar e continuar .
-
(Opcional) Na página Configurações , clique em Criar programa e defina a frequência refresh .
-
(Opcional) Configure notificações email para sucesso ou falha pipeline .
-
Clique em Salvar e pipelinede execução .
Implante um pipeline de ingestão baseado em query usando Pacotes de Automação Declarativa com o mecanismo de implantação direta. Os Pacotes contêm definições YAML de pipelines e jobs, são gerenciados com a CLI do Databricks e podem ser implantados em vários workspaces de destino. O mecanismo de implantação direta implanta pacotes sem Terraform e respeita as configurações de compute, como serverless e photon. Para mais informações, consulte O que são os Pacotes de Automação Declarativa? e Migrar para o mecanismo de implantação direta.
Este exemplo implanta o pipeline para compute serverless (default). Para implantar no compute clássico em vez disso, consulte o exemplo de compute clássico.
-
Criar um pacote:
Bashdatabricks bundle init -
Habilite o motor de implantação direta. Novos pacotes criados com a CLI do Databricks versão 1.3.0 e acima usam o motor de implantação direta por default. Se o seu pacote foi criado com uma versão anterior, defina o campo
bundle.enginede nível superior emdatabricks.yml(não é uma configuração por destino):YAMLbundle:
engine: directSe o pacote foi previamente implantado com o mecanismo de implantação Terraform, migre o estado da implantação para cada destino antes de implantar. Caso contrário, a implantação volta silenciosamente para o mecanismo Terraform. Execute
databricks bundle deployment migrate -t <target>para cada destino. Para obter mais informações, consulte Migrar um pacote existente. -
Adicione um arquivo de definição de pipeline ao pacote (por exemplo,
resources/query_based_pipeline.yml):YAMLvariables:
dest_catalog:
default: main
dest_schema:
default: ingest_destination_schema
resources:
pipelines:
pipeline_query_based:
name: query-based-ingestion-pipeline
ingestion_definition:
connection_name: <your-uc-connection-name>
objects:
- table:
source_catalog: <source-catalog>
source_schema: <source-schema>
source_table: <source-table>
table_configuration:
query_based_connector_config:
cursor_columns:
- updated_at
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
target: ${var.dest_schema}
catalog: ${var.dest_catalog} -
Adicione um arquivo de definição de tarefa que controle o programa de ingestão (por exemplo,
resources/query_based_job.yml):YAMLresources:
jobs:
query_based_job:
name: query_based_job
trigger:
periodic:
interval: 1
unit: HOURS
email_notifications:
on_failure:
- <email-address>
tasks:
- task_key: refresh_pipeline
pipeline_task:
pipeline_id: ${resources.pipelines.pipeline_query_based.id} -
implantar o pacote:
Bashdatabricks bundle deploy
Exemplo de compute clássico (Beta)
Beta
Classic compute for query-based ingestion pipelines is in Beta. Databricks recommends serverless compute for most workloads.
Para implantar em compute clássico, defina serverless: false e adicione um bloco clusters à definição do pipeline. O mecanismo de implantação direta respeita o campo serverless e outras configurações de compute, como photon. Após a implantação, serverless: false não aparece mais na especificação do pipeline, mas o bloco clusters é mantido e o pipeline é executado no compute clássico. Para o conjunto completo de campos de cluster compatíveis, consulte Configurar o compute clássico para pipelines.
-
Criar um pacote:
Bashdatabricks bundle init -
Habilite o motor de implantação direta. Novos pacotes criados com a CLI do Databricks versão 1.3.0 e acima usam o motor de implantação direta por default. Se o seu pacote foi criado com uma versão anterior, defina o campo
bundle.enginede nível superior emdatabricks.yml(não é uma configuração por destino):YAMLbundle:
engine: directSe o pacote foi previamente implantado com o mecanismo de implantação Terraform, migre o estado da implantação para cada destino antes de implantar. Caso contrário, a implantação volta silenciosamente para o mecanismo Terraform. Execute
databricks bundle deployment migrate -t <target>para cada destino. Para obter mais informações, consulte Migrar um pacote existente. -
Adicione um arquivo de definição de pipeline ao pacote (por exemplo,
resources/query_based_pipeline.yml):YAMLvariables:
dest_catalog:
default: main
dest_schema:
default: ingest_destination_schema
resources:
pipelines:
pipeline_query_based:
name: query-based-ingestion-pipeline
serverless: false
clusters:
- label: default
node_type_id: r6i.xlarge
driver_node_type_id: i3.large
autoscale:
min_workers: 1
max_workers: 5
ingestion_definition:
connection_name: <your-uc-connection-name>
objects:
- table:
source_catalog: <source-catalog>
source_schema: <source-schema>
source_table: <source-table>
table_configuration:
query_based_connector_config:
cursor_columns:
- updated_at
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
target: ${var.dest_schema}
catalog: ${var.dest_catalog} -
Adicione um arquivo de definição de tarefa que controle o programa de ingestão (por exemplo,
resources/query_based_job.yml):YAMLresources:
jobs:
query_based_job:
name: query_based_job
trigger:
periodic:
interval: 1
unit: HOURS
email_notifications:
on_failure:
- <email-address>
tasks:
- task_key: refresh_pipeline
pipeline_task:
pipeline_id: ${resources.pipelines.pipeline_query_based.id} -
implantar o pacote:
Bashdatabricks bundle deploy
Opção 2: Ingestão de catálogo estrangeiro
Use esta abordagem quando você deseja ingerir de um catálogo externo registrado na Lakehouse Federation. A ingestão de catálogo externo oferece suporte a todas as fontes de dados da Lakehouse Federation e ao acompanhamento de exclusão.
- Databricks UI
- Direct Bundles
A IU do Databricks implanta pipelines baseados em query para compute serverless. Para implantar em compute clássico, consulte a tab **Direct Bundles**.
-
Na barra lateral workspace Databricks , clique em inserção de dados .
-
Na página Adicionar dados , em Conectores do Databricks , clique na sua fonte. O assistente de ingestão é aberto.
-
Na página **Pipeline de ingestão**, insira um nome para o pipeline.
-
Em Catálogo de destino , selecione um catálogo Unity Catalog para armazenar os dados recebidos.
-
Para **Tipo de conexão**, selecione **Catálogo externo**, e então selecione o catálogo externo registrado na Lakehouse Federation.
-
Clique em Criar pipeline e continue .
-
Na página Origem , selecione os esquemas e tabelas a serem importados.
-
Para cada tabela, especifique a coluna do cursor . Deve ser uma única coluna com valores que aumentam monotonicamente (por exemplo,
updated_atourow_id). -
Opcionalmente, altere a configuração default da história acompanhamento. Para mais informações, consulte Habilitar história acompanhamento (SCD tipo 2).
-
Clique em Avançar .
-
Na página Destino , selecione o catálogo e o esquema do Unity Catalog nos quais deseja gravar.
Se não quiser usar um esquema existente, clique em Criar esquema . Você deve ter privilégios
USE CATALOGeCREATE SCHEMAno catálogo pai. -
Clique em Salvar e continuar .
-
(Opcional) Na página Configurações , clique em Criar programa e defina a frequência refresh .
-
(Opcional) Configure notificações email para sucesso ou falha pipeline .
-
Clique em Salvar e pipelinede execução .
Implantar um pipeline de ingestão de catálogo estrangeiro usando Pacotes de Automação Declarativa com o motor de implantação direta. Os Pacotes contêm definições YAML de pipelines e jobs, são gerenciados com a CLI do Databricks e podem ser implantados em vários workspaces de destino. O mecanismo de implantação direta implanta pacotes sem Terraform e respeita as configurações de compute, como serverless e photon. Para mais informações, consulte O que são os Pacotes de Automação Declarativa? e Migrar para o mecanismo de implantação direta.
Este exemplo implanta o pipeline para compute serverless (default). Para implantar no compute clássico em vez disso, consulte o exemplo de compute clássico.
-
Criar um pacote:
Bashdatabricks bundle init -
Habilite o motor de implantação direta. Novos pacotes criados com a CLI do Databricks versão 1.3.0 e acima usam o motor de implantação direta por default. Se o seu pacote foi criado com uma versão anterior, defina o campo
bundle.enginede nível superior emdatabricks.yml(não é uma configuração por destino):YAMLbundle:
engine: directSe o pacote foi previamente implantado com o mecanismo de implantação Terraform, migre o estado da implantação para cada destino antes de implantar. Caso contrário, a implantação volta silenciosamente para o mecanismo Terraform. Execute
databricks bundle deployment migrate -t <target>para cada destino. Para obter mais informações, consulte Migrar um pacote existente. -
Adicione um arquivo de definição de pipeline ao pacote (por exemplo,
resources/foreign_catalog_pipeline.yml):YAMLvariables:
dest_catalog:
default: main
dest_schema:
default: ingest_destination_schema
resources:
pipelines:
pipeline_foreign_catalog:
name: foreign-catalog-ingestion-pipeline
ingestion_definition:
ingest_from_uc_foreign_catalog: true
objects:
- table:
source_catalog: <foreign-catalog-name>
source_schema: <source-schema>
source_table: <source-table>
table_configuration:
primary_keys:
- id
query_based_connector_config:
cursor_columns:
- updated_at
deletion_condition: 'deleted_at IS NOT NULL'
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
target: ${var.dest_schema}
catalog: ${var.dest_catalog} -
Adicione um arquivo de definição de tarefa (por exemplo,
resources/foreign_catalog_job.yml):YAMLresources:
jobs:
foreign_catalog_job:
name: foreign_catalog_job
trigger:
periodic:
interval: 1
unit: HOURS
email_notifications:
on_failure:
- <email-address>
tasks:
- task_key: refresh_pipeline
pipeline_task:
pipeline_id: ${resources.pipelines.pipeline_foreign_catalog.id} -
implantar o pacote:
Bashdatabricks bundle deploy
Exemplo de compute clássico (Beta)
Beta
Classic compute for query-based ingestion pipelines is in Beta. Databricks recommends serverless compute for most workloads.
Para implantar em compute clássico, defina serverless: false e adicione um bloco clusters à definição do pipeline. O mecanismo de implantação direta respeita o campo serverless e outras configurações de compute, como photon. Após a implantação, serverless: false não aparece mais na especificação do pipeline, mas o bloco clusters é mantido e o pipeline é executado no compute clássico. Para o conjunto completo de campos de cluster compatíveis, consulte Configurar o compute clássico para pipelines.
-
Criar um pacote:
Bashdatabricks bundle init -
Habilite o motor de implantação direta. Novos pacotes criados com a CLI do Databricks versão 1.3.0 e acima usam o motor de implantação direta por default. Se o seu pacote foi criado com uma versão anterior, defina o campo
bundle.enginede nível superior emdatabricks.yml(não é uma configuração por destino):YAMLbundle:
engine: directSe o pacote foi previamente implantado com o mecanismo de implantação Terraform, migre o estado da implantação para cada destino antes de implantar. Caso contrário, a implantação volta silenciosamente para o mecanismo Terraform. Execute
databricks bundle deployment migrate -t <target>para cada destino. Para obter mais informações, consulte Migrar um pacote existente. -
Adicione um arquivo de definição de pipeline ao pacote (por exemplo,
resources/foreign_catalog_pipeline.yml):YAMLvariables:
dest_catalog:
default: main
dest_schema:
default: ingest_destination_schema
resources:
pipelines:
pipeline_foreign_catalog:
name: foreign-catalog-ingestion-pipeline
serverless: false
clusters:
- label: default
node_type_id: r6i.xlarge
driver_node_type_id: i3.large
autoscale:
min_workers: 1
max_workers: 5
ingestion_definition:
ingest_from_uc_foreign_catalog: true
objects:
- table:
source_catalog: <foreign-catalog-name>
source_schema: <source-schema>
source_table: <source-table>
table_configuration:
primary_keys:
- id
query_based_connector_config:
cursor_columns:
- updated_at
deletion_condition: 'deleted_at IS NOT NULL'
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
target: ${var.dest_schema}
catalog: ${var.dest_catalog} -
Adicione um arquivo de definição de tarefa (por exemplo,
resources/foreign_catalog_job.yml):YAMLresources:
jobs:
foreign_catalog_job:
name: foreign_catalog_job
trigger:
periodic:
interval: 1
unit: HOURS
email_notifications:
on_failure:
- <email-address>
tasks:
- task_key: refresh_pipeline
pipeline_task:
pipeline_id: ${resources.pipelines.pipeline_foreign_catalog.id} -
implantar o pacote:
Bashdatabricks bundle deploy
Configurar acompanhamento incremental
Conectores baseados em query usam uma coluna de cursor para determinar quais linhas são novas ou atualizadas após a última execução da pipeline. Sua escolha da coluna de cursor é fundamental para uma ingestão incremental eficaz.
Ao selecionar uma coluna do cursor, considere o seguinte:
- Utilize uma coluna de registro de data e hora, se possível. Colunas como
updated_atoulast_modifiedsão ideais porque refletem diretamente quando uma linha foi alterada pela última vez. - IDs inteiros funcionam para fontes somente de anexo. Caso as linhas nunca sejam atualizadas, é possível usar uma coluna de ID de incremento automático (como
idourow_id) como o cursor. Evite usar um ID de inteiro como um cursor se as linhas puderem ser atualizadas sem alterar o ID. - A coluna deve aumentar monotonicamente. Os valores nunca devem diminuir. Se um processo como um preenchimento retroativo define a coluna para um valor passado, o conector não reingere as linhas gravadas antes do marcador de alta anterior.
- Você só pode especificar uma única coluna de cursor. Não é possível especificar várias colunas como um cursor composto.
Depois que o conector armazena a marca d'água alta do cursor, ele usa a marca d'água alta como filtro de limite inferior (cursor_column > last_value) na próxima execução. O conector não ingere linhas com um valor de cursor NULO.
Configurar história envio (SCD)
Para acompanhar todo o histórico de alterações de linhas nas tabelas de destino, configure SCD tipo 2. Consulte Ativar acompanhamento de história (SCD tipo 2).
Padrões comuns
Para configurações avançadas pipeline , consulte Padrões comuns para gerenciar pipeline de ingestão.