Ingerir dados do Apache Kafka
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 CONNECTIONno 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 CONNECTIONouALL PRIVILEGESno objeto de conexão. -
Você deve ter privilégios
USE CATALOGno catálogo de destino. -
É necessário ter os privilégios
USE SCHEMAeCREATE TABLEem um esquema existente ou os privilégiosCREATE SCHEMAno 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.
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.
- Declarative Automation Bundles
- Databricks notebook
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?.
-
Crie um novo pacote utilizando a CLI do Databricks:
Bashdatabricks bundle init -
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. -
Implante o pacote usando a CLI do Databricks:
Bashdatabricks bundle deploy
- Importe o seguinte notebook em seu workspace do Databricks:
-
Manter a célula um como está. Não modifique o campo
channel— ele deve permanecerPREVIEW. -
Modifique a célula com os detalhes de configuração do seu pipeline, incluindo seu formato de desserialização desejado. Veja pipeline.ingestion_definition e Exemplos.
-
Clique em Executar tudo .
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:
- Declarative Automation Bundles
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:
- Declarative Automation Bundles
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)
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:
- Declarative Automation Bundles
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.