Pular para o conteúdo principal

Ingerir dados do Apache Kafka

info

Beta

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

Esta página mostra como criar um pipeline de ingestão gerenciado do Kafka usando o Databricks LakeFlow Connect.

Requisitos

  • Para criar um pipeline de ingestão, você deve primeiro atender aos seguintes requisitos:

    • Seu workspace deve estar habilitado para o Unity Catalog.

    • O compute serverless deve ser habilitado para seu workspace. Consulte Requisitos de computação serverless.

    • Para criar uma nova conexão, você deve ter privilégios CREATE CONNECTION no metastore. Consulte Gerenciar privilégios no Unity Catalog.

      Se o conector suportar a criação de pipelines baseada na IU, um administrador poderá criar a conexão e o pipeline simultaneamente, ao concluir os passos nesta página. No entanto, se os usuários que criam pipelines usam a autoria de pipeline baseada em API ou são usuários não administradores, um administrador deve primeiro criar a conexão no Catalog Explorer. Consulte Conectar-se a fontes de ingestão gerenciadas.

    • Para usar uma conexão existente, é preciso ter os privilégios USE CONNECTION ou ALL PRIVILEGES no objeto de conexão.

    • Você deve ter privilégios USE CATALOG no catálogo de destino.

    • É necessário ter os privilégios USE SCHEMA e CREATE TABLE em um esquema existente ou os privilégios CREATE SCHEMA no catálogo de destino.

  • Para ingerir do Kafka, primeiro, conclua os passos em Conectar ao Apache Kafka para ingestão gerenciada.

Criar pipeline de ingestão

Cada tópico Kafka é ingerido em uma tabela de transmissão. Para obter uma lista de dados compatíveis e limitações, consulte Dados compatíveis.

nota

A criação de pipelines baseada em interface de usuário não é compatível com o conector Kafka na versão Beta. Use Pacotes de Automação Declarativa ou um Notebook Databricks para criar seu pipeline.

Use Pacotes de Automação Declarativa para gerenciar pipelines Kafka como código. Os pacotes podem conter definições YAML de Jobs e tarefas, são gerenciados usando a CLI do Databricks e podem ser compartilhados e executados em diferentes workspaces de destino (como desenvolvimento, preparo e produção). Para obter mais informações, consulte O que são Pacotes de Automação Declarativa?.

  1. Crie um novo pacote utilizando a CLI do Databricks:

    Bash
    databricks bundle init
  2. Adicione um arquivo de definição de pipeline ao pacote (por exemplo, resources/kafka_pipeline.yml). Veja pipeline.definição_de_ingestão e Exemplos.

  3. Implante o pacote usando a CLI do Databricks:

    Bash
    databricks bundle deploy

Exemplos

Use estes exemplos para configurar seu pipeline.

Pipeline mínima — key e valor binários brutos

Este exemplo ingere um ou mais tópicos Kafka com colunas de key e de valor retidas como BINARY:

YAML
variables:
connection_name:
default: my-kafka-connection
dest_catalog:
default: main
dest_schema:
default: kafka_ingest
resources:
pipelines:
kafka_pipeline:
name: kafka-ingestion-pipeline
serverless: true
continuous: true
channel: PREVIEW
catalog: ${var.dest_catalog}
target: ${var.dest_schema}
ingestion_definition:
connection_name: ${var.connection_name}
objects:
- table:
source_table: N/A
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
destination_table: user_events
connector_options:
kafka_options:
topics: [user-events, power-user-events]

Pipeline com transformadores — valor JSON e key de string

Este exemplo desserializa as chaves de mensagem como STRING e os valores como JSON com a evolução do esquema habilitada:

YAML
variables:
connection_name:
default: my-kafka-connection
dest_catalog:
default: main
dest_schema:
default: kafka_ingest
resources:
pipelines:
kafka_pipeline:
name: kafka-ingestion-pipeline
serverless: true
continuous: true
channel: PREVIEW
catalog: ${var.dest_catalog}
target: ${var.dest_schema}
ingestion_definition:
connection_name: ${var.connection_name}
objects:
- table:
source_table: N/A
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
destination_table: user_events
connector_options:
kafka_options:
topics:
- user-events
starting_offset: latest
key_transformer:
format: STRING
value_transformer:
format: JSON
json_options:
schema_evolution_mode: rescue

Encaminhar registros para várias tabelas (fanout)

info

Visualização

Este recurso está em Prévia Privada. Para experimentar, fale com o seu contato no Databricks.

O Fanout roteia cada registro de uma única origem Kafka para uma das muitas tabelas de destino. Você configura a expansão em um objeto de esquema em vez de um objeto de tabela. Uma chave de roteamento derivada de cada registro determina o nome da tabela de destino: {destination_catalog}.{destination_schema}.{key_value}.

Defina fanout_options no objeto de esquema. O campo fanout_by é uma expressão SQL avaliada em relação ao registro de origem bruto (com as colunas Kafka key e value), e seu resultado se torna o segmento final (o nome da tabela) da tabela de destino {destination_catalog}.{destination_schema}.{result}. Como fanout_by é executado no registro bruto, converta o binário value em uma string antes de extrair um campo dele, como no exemplo a seguir. Seu valor resolvido é usado como o segmento do nome da tabela literalmente, sem aspas ou sanitização, portanto, deve ser um identificador de tabela não citado válido. Valores com espaços, pontos ou outros caracteres que não são válidos em um identificador sem aspas falham na gravação, assim como um valor que é inteiramente composto por dígitos (por exemplo, 123). Um dígito inicial é permitido quando o valor também contém uma letra ou underscore, como 2024_events.

Você pode aplicar, opcionalmente, uma única transformação JSON a cada rota. Essa transformação é executada em cada registro roteado após o roteamento, portanto, não afeta o valor que fanout_by vê. Para a lista completa de opções e as restrições v1, consulte Opções de Fanout e Limitações de Fanout.

O exemplo a seguir lê registros de tópicos que correspondem a um padrão e roteia cada registro para uma tabela nomeada a partir de seu campo event_type. A transformação JSON opcional analisa a coluna de mensagem value diretamente em cada tabela de destino:

YAML
variables:
connection_name:
default: my-kafka-connection
dest_catalog:
default: main
dest_schema:
default: kafka_ingest
resources:
pipelines:
kafka_pipeline:
name: kafka-ingestion-pipeline
serverless: true
continuous: true
channel: PREVIEW
catalog: ${var.dest_catalog}
target: ${var.dest_schema}
ingestion_definition:
connection_name: ${var.connection_name}
objects:
- schema:
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
connector_options:
kafka_options:
topic_pattern: 'events-.*'
starting_offset: earliest
fanout_options:
fanout_by: 'cast(value as string):event_type::string'
transforms:
- format: JSON
input_column: value

Padrões comuns

Para configurações avançadas de pipeline, consulte Padrões comuns para pipelines de ingestão gerenciados.

Passos seguintes

Comece e defina alertas no seu pipeline. Veja Tarefas comuns de manutenção de pipelines.

Recursos adicionais