Pular para o conteúdo principal

Crie um conector personalizado para o Lakeflow Connect

info

Beta

Este recurso está em Beta. Os administradores do Workspace podem controlar o acesso a este recurso na página Pré-visualizações . Consulte Gerenciar prévias do Databricks.

Os conectores personalizados permitem ingerir dados de uma fonte que o Lakeflow Connect não suporta com um conector gerenciado. Você constrói e testa seu conector, depois o implanta e coloca em execução em seu próprio Workspace do Databricks. Você não precisa registrá-lo na comunidade ou contribuir com qualquer repository compartilhado para usá-lo.

Desenvolva seu conector usando as ferramentas e os modelos no **repository** Lakeflow Comunidade Connectors no GitHub. O repository inclui ferramentas de desenvolvimento com tecnologia de IA para auxiliar em cada fase, incluindo pesquisa de origem, configuração de autenticação, implementação e teste.

Se posteriormente você quiser compartilhar seu conector com outros usuários, poderá contribuí-lo para a comunidade. Para usar um conector da comunidade existente, consulte Conectores da comunidade no Lakeflow Connect.

Requisitos

Antes de começar, certifique-se de ter:

  • Python 3.10 ou acima
  • Um espaço de trabalho do Databricks com o Unity Catalog habilitado
  • Credenciais de API para a fonte à qual você deseja se conectar
  • Git instalado localmente

Configurar o repository

Clone o repository Lakeflow comunidade Connectors e instale as dependências de desenvolvimento.

  1. Clone o repository:

    Bash
    git clone https://github.com/databrickslabs/lakeflow-community-connectors.git
    cd lakeflow-community-connectors
  2. Crie um ambiente virtual e instale as dependências:

    Bash
    python -m venv .venv
    source .venv/bin/activate
    pip install -e ".[dev]"
  3. Revise as implementações de conector existentes em src/databricks/labs/community_connector/sources/, então comece a desenvolver seu conector em um novo diretório sob esse caminho. Siga os comandos e habilidades de desenvolvimento assistido por AI do repository. Para o fluxo de trabalho recomendado, use:

    Text
    /develop-connector <your-source>
    /validate-connector <your-source>

Implementar a interface LakeflowConnect

Cada conector implementa a interface LakeflowConnect, que define como seu conector autentica, descobre tabelas, retorna esquemas e lê dados.

Python
class LakeflowConnect:
def __init__(self, options: dict[str, str]) -> None:
"""Initialize with connection parameters"""

def list_tables(self) -> list[str]:
"""Return names of all tables supported by this connector."""

def get_table_schema(self, table_name: str, table_options: dict[str, str]) -> StructType:
"""Return the Spark schema for a table."""

def read_table_metadata(self, table_name: str, table_options: dict[str, str]) -> dict:
"""Return metadata: primary_keys, cursor_field, ingestion_type
(snapshot|cdc|cdc_with_deletes|append)."""

def read_table(self, table_name: str, start_offset: dict,
table_options: dict[str, str]) -> (Iterator[dict], dict):
"""Yield records as JSON dicts and return the next offset
for incremental reads."""

def read_table_deletes(self, table_name: str, start_offset: dict,
table_options: dict[str, str]) -> (Iterator[dict], dict):
"""Optional: Only required if ingestion_type is 'cdc_with_deletes'."""

Descrições de métodos

A tabela a seguir descreve cada método na interface LakeflowConnect:

Método

Descrição

__init__

Recebe os parâmetros de conexão como um dicionário e inicializa o cliente da API para sua origem.

list_tables

Retorna os nomes de todas as tabelas (ou Endpoint de API) que seu conector expõe. O Databricks usa esta lista para preencher a IU de seleção de tabelas.

get_table_schema

Retorna um StructType do Spark que descreve o esquema da tabela fornecida. Chamado antes da primeira execução do pipeline e em cada execução quando a evolução do esquema está habilitada.

read_table_metadata

Retorna um dicionário com primary_keys, cursor_field e ingestion_type. O ingestion_type deve ser um dos snapshot, cdc, cdc_with_deletes ou append.

read_table

Produz registros como dicionários Python e retorna o próximo deslocamento para leituras incrementais. Na primeira execução, start_offset está vazio. Em execuções subsequentes, ele contém o deslocamento retornado pela execução anterior.

read_table_deletes

Opcional. Implemente este método apenas se ingestion_type for cdc_with_deletes. Produz chaves de registro excluídas e retorna o próximo offset.

Método

Descrição

__init__

Recebe os parâmetros de conexão como um dicionário e inicializa o cliente da API para sua origem.

list_tables

Retorna os nomes de todas as tabelas (ou Endpoint de API) que seu conector expõe. O Databricks usa esta lista para preencher a IU de seleção de tabelas.

get_table_schema

Retorna um StructType do Spark que descreve o esquema da tabela fornecida. Chamado antes da primeira execução do pipeline e em cada execução quando a evolução do esquema está habilitada.

read_table_metadata

Retorna um dicionário com primary_keys, cursor_field e ingestion_type. O ingestion_type deve ser um dos snapshot, cdc, cdc_with_deletes ou append.

read_table

Produz registros como dicionários Python e retorna o próximo deslocamento para leituras incrementais. Na primeira execução, start_offset está vazio. Em execuções subsequentes, ele contém o deslocamento retornado pela execução anterior.

read_table_deletes

Opcional. Implemente este método apenas se ingestion_type for cdc_with_deletes. Produz chaves de registro excluídas e retorna o próximo offset.

Desenvolva seu conector

Siga os passos para construir e validar um novo conector:

  1. Pesquisar a API da fonte : estude as especificações da API da fonte, os mecanismos de autenticação, os limites de taxa e os esquemas de dados disponíveis. Identifique quais tabelas ou Endpoint expor.

  2. Configurar a autenticação : gere a especificação da conexão, configure as credenciais para a fonte e verifique a conectividade a partir do seu ambiente de desenvolvimento.

  3. Implemente o conector : codifique todos os métodos de interface LakeflowConnect necessários para conectar-se à API de origem e retornar dados no formato esperado.

  4. Testar e iterar : faça a execução dos conjuntos de testes padrão em um sistema de fonte real e corrija quaisquer problemas. Consulte Testar seu conector para obter detalhes.

  5. Documente o conector : escreva um README.md voltado para o usuário e gere o arquivo YAML de especificação do conector que descreve os parâmetros configuráveis do conector.

  6. Construir o artefato de implantação : faça a execução do script de build para produzir o artefato de arquivo único que pode ser implantado em um Workspace.

Testar seu conector

O repository fornece várias abordagens de teste:

Conjunto de testes genérico (obrigatório)

Esta suíte se conecta a uma origem real usando suas credenciais fornecidas para verificar a funcionalidade de ponta a ponta, incluindo autenticação, descoberta de esquema e leituras de dados.

Bash
python -m pytest tests/generic/ --connector <your-source> --credentials credentials.json

Teste de write-back (recomendado)

O teste de write-back executa ciclos de gravação-leitura-verificação para validar leituras e exclusões incrementais. Isso confirma que seu acompanhamento de offset e a lógica de CDC funcionam corretamente.

Bash
python -m pytest tests/writeback/ --connector <your-source> --credentials credentials.json

Testes de unidade

Escreva testes unitários para qualquer lógica personalizada complexa em seu conector, como tratamento de paginação, coerção de tipo ou recuperação de erros.

Construir o artefato de implantação

Depois que seu conector passar pelos conjuntos de testes, empacote-o para que um pipeline possa executá-lo. Um conector é implantado em duas partes:

  • Um artefato de origem de arquivo único. Execute o script de merge para simplificar seu conector em um único arquivo Python independente. O pipeline usa este arquivo durante o Runtime em vez do repository completo.

    Bash
    python tools/scripts/merge_python_source.py --connector <your-source>

    O script grava este arquivo em dist/<your-source>/.

  • Python wheels (.whl) para as dependências do conector. O framework do conector e quaisquer bibliotecas de terceiros que seu conector importe devem estar disponíveis para o pipeline como wheels armazenados em um volume do Unity Catalog. Se esses wheels estiverem ausentes, o pipeline poderá falhar durante a descoberta de origem. Você pode fazer o upload deles por conta própria e referenciá-los na IU, ou permitir que a CLI do Community Connector os compile e faça o upload para você. Consulte Deploy with the Community Connector CLI.

Implante seu conector de uma das duas maneiras:

  • A interface do usuário do Databricks é o caminho de apontar e clicar. Use-o para uma implantação única quando você mesmo fornecer os wheels do conector no campo Dependências da biblioteca . Consulte Implantar na interface do usuário do Databricks.
  • A CLI community-connector é o caminho programável. Use-o quando estiver desenvolvendo localmente e quiser que as wheels do conector sejam criadas e upload para você, ou quando quiser uma implantação repetível que possa automatizar. Consulte Implantar com a CLI do Community Connector.

Implantar na interface do usuário do Databricks

Implante seu conector na interface do usuário do Databricks em duas fases: adicione o conector e, em seguida, crie o pipeline.

Adicionar o conector personalizado

Primeiro, adicione seu conector para que ele apareça como um bloco na página Adicionar dados :

  1. Na barra lateral do seu workspace do Databricks, clique em +Novo > Adicionar ou fazer upload de dados e, em Conectores da comunidade , adicione um conector personalizado.
  2. Para Nome da fonte , insira o nome do seu conector. Isso deve corresponder ao nome do diretório que contém o código-fonte do seu conector (sources/<source-name>).
  3. Para Nome de exibição , insira um nome amigável para o conector. Se você deixar este campo em branco, ele usará o nome da fonte como default.
  4. Para Dependências de biblioteca , adicione os arquivos Python wheel (.whl) que seu conector precisa a partir de um volume do Unity Catalog. Consulte Criar o artefato de implantação.
  5. Para Especificação de conexão , cole a especificação de conexão do conector em YAML, correspondendo ao seu arquivo connector_spec.yaml.
  6. Clique em Salvar . O conector aparece como um bloco Personalizado em Conectores da comunidade .

Criar o pipeline de ingestão

Em seguida, crie o pipeline que ingere dados da sua origem:

  1. Selecione o bloco do seu conector para abrir o assistente Ingerir dados .
  2. Na etapa Conexão , clique em + Criar conexão ou selecione uma conexão existente, insira os detalhes da conexão para sua fonte e clique em Próximo .
  3. Na etapa Configuração de ingestão , insira um Nome do pipeline , defina o Local do log de eventos (catálogo e esquema), escolha um Tipo de compute e clique em Criar pipeline e continuar .
  4. No passo Source , selecione as tabelas para ingestão.
  5. No passo Destino , escolha o catálogo e o esquema onde as tabelas ingeridas são gravadas.
  6. No **o passo** Schedules and notifications , defina uma programar e notificações opcionais e, em seguida, conclua.
  7. Execute o pipeline manualmente ou conforme sua programação.

Para configurar o pipeline ainda mais, você pode editar ingest.py no editor de pipelines. Consulte Opções de configuração do pipeline.

Opções de configuração do pipeline

Você pode configurar as seguintes opções em ingest.py:

Opção

Descrição

connection_name

Obrigatório. O nome da conexão que armazena as credenciais de autenticação para a fonte.

objects

Obrigatório. Uma lista de tabelas para ingestão. Cada entrada tem o formato {"table": {"source_table": "..."}}. Você também pode especificar um destination_table opcional dentro do objeto table.

destination_catalog

O catálogo onde as tabelas ingeridas são gravadas. Adota como default o catálogo definido durante a criação do pipeline.

destination_schema

O esquema em que as tabelas ingeridas são gravadas. Adota como default o esquema definido durante a criação do pipeline.

scd_type

A estratégia de dimensões que mudam lentamente (SCD): SCD_TYPE_1, SCD_TYPE_2 ou APPEND_ONLY. O default é SCD_TYPE_1.

primary_keys

Substituir as chaves primárias default de uma tabela. Forneça uma lista de nomes de colunas.

Opção

Descrição

connection_name

Obrigatório. O nome da conexão que armazena as credenciais de autenticação para a fonte.

objects

Obrigatório. Uma lista de tabelas para ingestão. Cada entrada tem o formato {"table": {"source_table": "..."}}. Você também pode especificar um destination_table opcional dentro do objeto table.

destination_catalog

O catálogo onde as tabelas ingeridas são gravadas. Adota como default o catálogo definido durante a criação do pipeline.

destination_schema

O esquema em que as tabelas ingeridas são gravadas. Adota como default o esquema definido durante a criação do pipeline.

scd_type

A estratégia de dimensões que mudam lentamente (SCD): SCD_TYPE_1, SCD_TYPE_2 ou APPEND_ONLY. O default é SCD_TYPE_1.

primary_keys

Substituir as chaves primárias default de uma tabela. Forneça uma lista de nomes de colunas.

Implantar com a CLI do conector da comunidade

O comando publish da CLI cria o framework e as wheels do conector a partir da sua fonte local, faz o upload delas para um volume do Unity Catalog, registra seus caminhos no manifesto do conector e publica o conector como um bloco Personalizado na página Adicionar dados . Aponte-o para a especificação do seu conector local para que ele não procure o conector no repository upstream:

Bash
community-connector publish <your-source> \
--spec src/databricks/labs/community_connector/sources/<your-source>/connector_spec.yaml

Para reutilizar wheels que você já compilou e pular o passo de compilação, passe-os com --package. Para substituir um conector que você publicou anteriormente, adicione --overwrite. Para a lista completa de opções, incluindo --package, --volume-path, --catalog e --schema, consulte a referência de comandopublish.

Para executar todo o fluxo de trabalho a partir da linha de comando, incluindo a criação de uma conexão, a criação e atualização do pipeline de ingestão, a publicação e a despublicação, consulte a referência da CLI do Community Connector.

Contribua com seu conector para a comunidade

Seu conector é executado em seu workspace, independentemente de você contribuir com ele ou não. Se você deseja compartilhá-lo para que outros usuários possam descobri-lo e usá-lo, abra um pull request no repository Lakeflow Community Connectors. Conectores contribuídos tornam-se conectores da comunidade, que a comunidade mantém e que não são suportados por SLAs da Databricks.