Pular para o conteúdo principal

Referência do conector Kafka

Esta página documenta as opções de conector, opções de configuração de tabela e configurações de transformador JSON para o conector gerenciado do Apache Kafka no Lakeflow Connect.

Opções de conector​

As seguintes opções configuram a fonte do Kafka para cada tabela de destino no pipeline de ingestão. Especifique estas opções em connector_options.kafka_options na definição do pipeline. Consulte Exemplos para obter exemplos completos de pipeline.

Opção

Tipo

Padrão

Descrição

topics

Lista de strings

—

Lista de nomes de tópicos para assinar. Mutuamente exclusivo com topic_pattern. É necessário topics ou topic_pattern.

topic_pattern

String

—

Expressões regulares em Java para correspondência de nomes de tópicos aos quais se inscrever. Mutuamente exclusivo com topics.

starting_offset

String

latest

Por onde começar a leitura quando não há ponto de verificação (somente na primeira execução)? Valores válidos: latest, earliest.

key_transformer

Transformer

—

Configuração do desserializador para chaves de mensagem. Se não for definida, a coluna key é mantida como BINARY. Consulte as opções do Transformer.

value_transformer

Transformer

—

Configuração do desserializador para valores de mensagem. Se não for definida, a coluna de valor será mantida como BINARY. Consulte Opções do Transformer.

Opção

Tipo

Padrão

Descrição

topics

Lista de strings

—

Lista de nomes de tópicos para assinar. Mutuamente exclusivo com topic_pattern. É necessário topics ou topic_pattern.

topic_pattern

String

—

Expressões regulares em Java para correspondência de nomes de tópicos aos quais se inscrever. Mutuamente exclusivo com topics.

starting_offset

String

latest

Por onde começar a leitura quando não há ponto de verificação (somente na primeira execução)? Valores válidos: latest, earliest.

key_transformer

Transformer

—

Configuração do desserializador para chaves de mensagem. Se não for definida, a coluna key é mantida como BINARY. Consulte as opções do Transformer.

value_transformer

Transformer

—

Configuração do desserializador para valores de mensagem. Se não for definida, a coluna de valor será mantida como BINARY. Consulte Opções do Transformer.

Opções do Transformer​

Os transformadores definem como as keys e valores de mensagens binárias do Kafka são desserializados em colunas estruturadas. Especifique o formato de serialização e as opções correspondentes específicas do formato em key_transformer ou value_transformer. É possível configurar um transformador para a key, o valor, ou ambos de forma independente. Se nenhum transformador for definido, a coluna será retida como BINARY.

Para JSON, você pode fornecer um esquema explícito, usar inferência de esquema com evolução ou omitir json_options inteiramente para armazenar o valor como uma coluna VARIANT.

Opção

Aplica-se a

Tipo

Padrão

Descrição

format

Todos

String

—

Formato de serialização dos dados. Valores válidos: STRING, JSON, AVRO, PROTOBUF. STRING não requer opções adicionais. Se nenhum json_options for especificado no transformador, o valor será analisado como VARIANT por default. Consulte Formato de dados Variant para obter mais informações. Para AVRO e PROTOBUF, consulte Opções do Avro e Opções do Protobuf.

json_options.schema

JSON

String

—

Esquema em linha no formato DDL do Spark (por exemplo, "id BIGINT, name STRING"). Mutuamente exclusivo com schema_file_path.

json_options.schema_file_path

JSON

String

—

Caminho para um arquivo de esquema .ddl. Mutuamente exclusivo com schema. Oferece suporte a caminhos de Volumes do Unity Catalog (/Volumes/...).

json_options.schema_evolution_mode

JSON

String

—

Modo de evolução do esquema para inferência automática do esquema. Consulte Modos de Evolução do Esquema.

json_options.schema_hints

JSON

String

—

Pares "column_name type" separados por vírgula para influenciar a inferência de esquema (por exemplo, "id BIGINT, ts TIMESTAMP"). É necessário definir schema_evolution_mode. Consulte Substituir a inferência de esquema usando dicas de esquema.

Opção

Aplica-se a

Tipo

Padrão

Descrição

format

Todos

String

—

Formato de serialização dos dados. Valores válidos: STRING, JSON, AVRO, PROTOBUF. STRING não requer opções adicionais. Se nenhum json_options for especificado no transformador, o valor será analisado como VARIANT por default. Consulte Formato de dados Variant para obter mais informações. Para AVRO e PROTOBUF, consulte Opções do Avro e Opções do Protobuf.

json_options.schema

JSON

String

—

Esquema em linha no formato DDL do Spark (por exemplo, "id BIGINT, name STRING"). Mutuamente exclusivo com schema_file_path.

json_options.schema_file_path

JSON

String

—

Caminho para um arquivo de esquema .ddl. Mutuamente exclusivo com schema. Oferece suporte a caminhos de Volumes do Unity Catalog (/Volumes/...).

json_options.schema_evolution_mode

JSON

String

—

Modo de evolução do esquema para inferência automática do esquema. Consulte Modos de Evolução do Esquema.

json_options.schema_hints

JSON

String

—

Pares "column_name type" separados por vírgula para influenciar a inferência de esquema (por exemplo, "id BIGINT, ts TIMESTAMP"). É necessário definir schema_evolution_mode. Consulte Substituir a inferência de esquema usando dicas de esquema.

Opções do Avro​

Defina essas opções em avro_options quando format: AVRO. Forneça o esquema embutido, a partir de um arquivo ou de um registro de esquemas.

Opção

Tipo

Padrão

Descrição

avro_options.schema

String

—

Esquema Avro embutido no formato JSON. Mutuamente exclusivo com schema_file_path e schema_registry.

avro_options.schema_file_path

String

—

Caminho para um arquivo de esquema .avsc. Oferece suporte a caminhos de Unity Catalog Volumes (/Volumes/...). Mutuamente exclusivo com schema e schema_registry.

avro_options.schema_registry

Objeto

—

Resolver o esquema de um registro de esquema no runtime em vez de schema ou schema_file_path. Consulte Opções de registro de esquema.

avro_options.parse_mode

String

PERMISSIVE

How to handle records that fail to deserialize. Valores válidos: PERMISSIVE (default), que grava os bytes brutos de cada registro que falha ao desserializar em uma coluna _corrupt_record (tipo BINARY), define as outras colunas do registro como null e continua o processamento; e FAILFAST, que faz com que o pipeline falhe no primeiro registro que falha ao desserializar.

Opção

Tipo

Padrão

Descrição

avro_options.schema

String

—

Esquema Avro embutido no formato JSON. Mutuamente exclusivo com schema_file_path e schema_registry.

avro_options.schema_file_path

String

—

Caminho para um arquivo de esquema .avsc. Oferece suporte a caminhos de Unity Catalog Volumes (/Volumes/...). Mutuamente exclusivo com schema e schema_registry.

avro_options.schema_registry

Objeto

—

Resolver o esquema de um registro de esquema no runtime em vez de schema ou schema_file_path. Consulte Opções de registro de esquema.

avro_options.parse_mode

String

PERMISSIVE

How to handle records that fail to deserialize. Valores válidos: PERMISSIVE (default), que grava os bytes brutos de cada registro que falha ao desserializar em uma coluna _corrupt_record (tipo BINARY), define as outras colunas do registro como null e continua o processamento; e FAILFAST, que faz com que o pipeline falhe no primeiro registro que falha ao desserializar.

Opções do Protobuf​

Defina estas opções em protobuf_options quando format: PROTOBUF. Forneça um arquivo de conjunto de descritores compilado (.desc) e o nome da mensagem ou resolva o esquema de um registro de esquema.

Opção

Tipo

Padrão

Descrição

protobuf_options.desc_file_path

String

—

Caminho para um arquivo de conjunto de descritores Protobuf compilado (.desc). Dá suporte a caminhos de Unity Catalog Volumes (/Volumes/...). Obrigatório, a menos que schema_registry esteja definido.

protobuf_options.message_name

String

—

Nome do tipo de mensagem Protobuf totalmente qualificado (por exemplo, com.example.events.UserEvent). Exige com desc_file_path.

protobuf_options.schema_registry

Objeto

—

Resolver o esquema de um registro de esquema em Runtime em vez de desc_file_path. Consulte Opções de registro de esquema.

protobuf_options.recursive_fields_max_depth

Integer

—

Profundidade máxima de expansão para campos Protobuf recursivos, aos quais o Spark SQL não oferece suporte nativo. Valores válidos: -1 (desativar campos recursivos), 0 (descartá-los), 1–10.

protobuf_options.parse_mode

String

PERMISSIVE

How to handle records that fail to deserialize. Valores válidos: PERMISSIVE (default), que grava os bytes brutos de cada registro que falha ao desserializar em uma coluna _corrupt_record (tipo BINARY), define as outras colunas do registro como null e continua o processamento; e FAILFAST, que faz com que o pipeline falhe no primeiro registro que falha ao desserializar.

Opção

Tipo

Padrão

Descrição

protobuf_options.desc_file_path

String

—

Caminho para um arquivo de conjunto de descritores Protobuf compilado (.desc). Dá suporte a caminhos de Unity Catalog Volumes (/Volumes/...). Obrigatório, a menos que schema_registry esteja definido.

protobuf_options.message_name

String

—

Nome do tipo de mensagem Protobuf totalmente qualificado (por exemplo, com.example.events.UserEvent). Exige com desc_file_path.

protobuf_options.schema_registry

Objeto

—

Resolver o esquema de um registro de esquema em Runtime em vez de desc_file_path. Consulte Opções de registro de esquema.

protobuf_options.recursive_fields_max_depth

Integer

—

Profundidade máxima de expansão para campos Protobuf recursivos, aos quais o Spark SQL não oferece suporte nativo. Valores válidos: -1 (desativar campos recursivos), 0 (descartá-los), 1–10.

protobuf_options.parse_mode

String

PERMISSIVE

How to handle records that fail to deserialize. Valores válidos: PERMISSIVE (default), que grava os bytes brutos de cada registro que falha ao desserializar em uma coluna _corrupt_record (tipo BINARY), define as outras colunas do registro como null e continua o processamento; e FAILFAST, que faz com que o pipeline falhe no primeiro registro que falha ao desserializar.

Opções de registro de esquema​

Defina schema_registry em avro_options ou protobuf_options para resolver o esquema em runtime a partir de um registro de esquema compatível com o Confluent. Por default, o pipeline se autentica no registro usando a conexão de origem Kafka do pipeline, que armazena a URL do registro e a API key. Para se autenticar com uma conexão do Unity Catalog diferente, defina connection_name. Consulte Propriedades de conexão.

Opção

Tipo

Padrão

Descrição

schema_registry.confluent_options.subject

String

—

Obrigatório. O assunto a ser resolvido no registro de esquema compatível com o Confluent.

schema_registry.connection_name

String

—

Uma conexão do Unity Catalog usada para autenticar no registro. Usa como default a conexão de origem do Kafka do pipeline. Defina isso quando o registro usar credenciais diferentes.

schema_registry.protobuf_message_name

String

—

Somente Protobuf. Seleciona uma mensagem quando o assunto define mais de uma mensagem Protobuf. Simples (Location) ou totalmente qualificado (com.example.protos.Location). Usa como default a primeira mensagem no esquema.

Opção

Tipo

Padrão

Descrição

schema_registry.confluent_options.subject

String

—

Obrigatório. O assunto a ser resolvido no registro de esquema compatível com o Confluent.

schema_registry.connection_name

String

—

Uma conexão do Unity Catalog usada para autenticar no registro. Usa como default a conexão de origem do Kafka do pipeline. Defina isso quando o registro usar credenciais diferentes.

schema_registry.protobuf_message_name

String

—

Somente Protobuf. Seleciona uma mensagem quando o assunto define mais de uma mensagem Protobuf. Simples (Location) ou totalmente qualificado (com.example.protos.Location). Usa como default a primeira mensagem no esquema.

Opções de configuração de tabela​

As seguintes opções são definidas em table_configuration em um objeto de tabela, um irmão de connector_options. Consulte Exemplos para obter exemplos completos de pipeline.

Opção

Tipo

Padrão

Descrição

source_metadata_column

String

—

Nome de uma coluna de estrutura adicionada à tabela de destino que contém os metadados de origem do Kafka para cada registro. Consulte Coluna de metadados de origem. O nome não deve ser key ou value.

Opção

Tipo

Padrão

Descrição

source_metadata_column

String

—

Nome de uma coluna de estrutura adicionada à tabela de destino que contém os metadados de origem do Kafka para cada registro. Consulte Coluna de metadados de origem. O nome não deve ser key ou value.

Opções de fanout​

info

Visualização

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

As opções de fanout direcionam cada registro de uma única origem Kafka para uma de muitas tabelas de destino. Especifique estas opções em fanout_options em um objeto de esquema (não um objeto de tabela) na sua definição de pipeline. Consulte Direcionar registros para várias tabelas (fanout) para obter um exemplo completo de pipeline e Limitações de fanout para as restrições.

Opção

Tipo

Padrão

Descrição

fanout_by

String

—

Obrigatório. Expressão SQL avaliada em relação ao registro de origem bruto (com as colunas Kafka key e value) cujo resultado determina o nome da tabela de destino. O valor torna-se o segmento final do nome da tabela: {destination_catalog}.{destination_schema}.{value}. A coluna value é binária, então converta-a em uma string antes de extrair um campo, por exemplo cast(value as string):event_type::string. A expressão deve ser resolvida para uma STRING não nula, e essa string é usada como o segmento do nome da tabela literalmente, sem aspas ou higienização, portanto, deve ser um identificador de tabela sem aspas 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 dígitos (por exemplo, 123). Um dígito inicial é permitido quando o valor também contém uma letra ou underscore, como 2024_events. As tabelas de destino são criadas automaticamente se ainda não existirem.

transforms

Lista de Transformer

—

Uma transformação aplicada a cada registro roteado após o roteamento, antes de ser gravada em sua tabela de destino. Como é executado após o roteamento, não afeta o valor que fanout_by vê. É permitida no máximo uma transformação, e ela deve usar format: JSON. Consulte opções de transformação fanout.

Opção

Tipo

Padrão

Descrição

fanout_by

String

—

Obrigatório. Expressão SQL avaliada em relação ao registro de origem bruto (com as colunas Kafka key e value) cujo resultado determina o nome da tabela de destino. O valor torna-se o segmento final do nome da tabela: {destination_catalog}.{destination_schema}.{value}. A coluna value é binária, então converta-a em uma string antes de extrair um campo, por exemplo cast(value as string):event_type::string. A expressão deve ser resolvida para uma STRING não nula, e essa string é usada como o segmento do nome da tabela literalmente, sem aspas ou higienização, portanto, deve ser um identificador de tabela sem aspas 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 dígitos (por exemplo, 123). Um dígito inicial é permitido quando o valor também contém uma letra ou underscore, como 2024_events. As tabelas de destino são criadas automaticamente se ainda não existirem.

transforms

Lista de Transformer

—

Uma transformação aplicada a cada registro roteado após o roteamento, antes de ser gravada em sua tabela de destino. Como é executado após o roteamento, não afeta o valor que fanout_by vê. É permitida no máximo uma transformação, e ela deve usar format: JSON. Consulte opções de transformação fanout.

Opções de transformação fanout​

Cada entrada em transforms usa as seguintes opções. Apenas o formato JSON é compatível com transformações fanout.

Opção

Tipo

Padrão

Descrição

format

String

JSON

Opcional. Formato de serialização da transformação. Somente JSON é compatível para transformações fanout, e é o default quando omitido. STRING, AVRO e PROTOBUF não são suportados.

input_column

String

—

A coluna da qual a transformação lê e para a qual grava (por exemplo, value). A transformação JSON fanout analisa esta coluna no local.

Opção

Tipo

Padrão

Descrição

format

String

JSON

Opcional. Formato de serialização da transformação. Somente JSON é compatível para transformações fanout, e é o default quando omitido. STRING, AVRO e PROTOBUF não são suportados.

input_column

String

—

A coluna da qual a transformação lê e para a qual grava (por exemplo, value). A transformação JSON fanout analisa esta coluna no local.

nota

A opção de transformação output_column não é aplicada às transformações fanout. A transformação fanout JSON sempre grava seu resultado de volta em input_column no local; qualquer valor output_column é ignorado.

Propriedades de conexão​

Ao criar a conexão do Unity Catalog com o Kafka no Catalog Explorer, você deve especificar as seguintes propriedades, dependendo do método de autenticação. Consulte Criar uma conexão Kafka para obter instruções sobre como criar uma conexão.

Nome de usuário e senha​

Propriedade

Descrição

Nome da conexão

Um nome exclusivo para a conexão do Unity Catalog.

Tipo de conexão

Select Kafka .

Tipo de autenticação

Selecione **Nome de usuário e senha**.

Nome de usuário

O nome de usuário usado para autenticar com o cluster Kafka.

Senha

A senha usada para autenticar com o cluster Kafka.

Servidores bootstrap

O endereço do servidor de inicialização do cluster Kafka (por exemplo, broker1:9092,broker2:9092).

URL do registro de esquema (opcional)

O URL do seu registro de esquema.

API key do registro de esquema (opcional)

A API key para o seu registro de esquema.

Segredo da API de registro de esquema (opcional)

O segredo da API para o seu registro de esquema.

Propriedade

Descrição

Nome da conexão

Um nome exclusivo para a conexão do Unity Catalog.

Tipo de conexão

Select Kafka .

Tipo de autenticação

Selecione **Nome de usuário e senha**.

Nome de usuário

O nome de usuário usado para autenticar com o cluster Kafka.

Senha

A senha usada para autenticar com o cluster Kafka.

Servidores bootstrap

O endereço do servidor de inicialização do cluster Kafka (por exemplo, broker1:9092,broker2:9092).

URL do registro de esquema (opcional)

O URL do seu registro de esquema.

API key do registro de esquema (opcional)

A API key para o seu registro de esquema.

Segredo da API de registro de esquema (opcional)

O segredo da API para o seu registro de esquema.

Credencial de Serviço​

Propriedade

Descrição

Nome da conexão

Um nome exclusivo para a conexão do Unity Catalog.

Tipo de conexão

Select Kafka .

Tipo de autenticação

Selecione **Credencial de serviço**.

Credencial de serviço

Selecione uma credencial de serviço do Unity Catalog existente ou clique em Criar nova credencial de serviço .

Servidores bootstrap

O endereço do servidor de inicialização do cluster Kafka (por exemplo, broker1:9092,broker2:9092).

URL do registro de esquema (opcional)

O URL do seu registro de esquema.

API key do registro de esquema (opcional)

A API key para o seu registro de esquema.

Segredo da API de registro de esquema (opcional)

O segredo da API para o seu registro de esquema.

Propriedade

Descrição

Nome da conexão

Um nome exclusivo para a conexão do Unity Catalog.

Tipo de conexão

Select Kafka .

Tipo de autenticação

Selecione **Credencial de serviço**.

Credencial de serviço

Selecione uma credencial de serviço do Unity Catalog existente ou clique em Criar nova credencial de serviço .

Servidores bootstrap

O endereço do servidor de inicialização do cluster Kafka (por exemplo, broker1:9092,broker2:9092).

URL do registro de esquema (opcional)

O URL do seu registro de esquema.

API key do registro de esquema (opcional)

A API key para o seu registro de esquema.

Segredo da API de registro de esquema (opcional)

O segredo da API para o seu registro de esquema.

Esquema da tabela de destino​

O conector Kafka grava em tabelas de transmissão (somente acréscimo). As colunas gravadas na tabela de destino dependem se os transformadores estão configurados.

Sem transformadores (binário bruto)​

Quando nenhum key_transformer ou value_transformer está configurado, a tabela de destino contém as seguintes colunas:

Coluna

Tipo

Descrição

key

BINARY

O conteúdo binário bruto da chave da mensagem do Kafka.

value

BINARY

O conteúdo binário bruto do valor da mensagem Kafka.

Coluna

Tipo

Descrição

key

BINARY

O conteúdo binário bruto da chave da mensagem do Kafka.

value

BINARY

O conteúdo binário bruto do valor da mensagem Kafka.

Com um transformador STRING​

Quando format: STRING é definido em um transformador, a coluna correspondente é escrita como STRING em vez de BINARY.

Com um transformador JSON​

Quando format: JSON é definido em um transformador:

  • Se nenhum(a) json_options for/forem especificado(a)s, a coluna será escrita como VARIANT.
  • Se json_options.schema ou json_options.schema_file_path for especificado, o JSON é analisado em colunas tipadas correspondentes ao esquema.
  • Se json_options.schema_evolution_mode estiver definido, a inferência de esquema será usada e o esquema evoluirá automaticamente.

Com um transformador Avro​

Quando format: AVRO for definido em um transformador, a coluna será analisada em colunas tipadas correspondentes ao esquema Avro. No modo PERMISSIVE (o default), uma coluna _corrupt_record do tipo BINARY também é adicionada; ela contém os bytes brutos de qualquer registro que falhe na desserialização e é null para registros analisados com sucesso.

Com um transformador Protobuf​

Quando format: PROTOBUF é definido em um transformador, a coluna é analisada em colunas tipadas correspondentes à definição de mensagem Protobuf. No modo PERMISSIVE (o default), uma coluna _corrupt_record do tipo BINARY também é adicionada; ela contém os bytes brutos de qualquer registro cuja desserialização falhe e é null para registros analisados com sucesso.

Coluna de metadados de origem​

Defina source_metadata_column em table_configuration para adicionar uma coluna struct com esse nome à tabela de destino. O struct contém os seguintes campos de metadados de origem do Kafka para cada registro: topic, partition, offset, timestamp, timestampType e headers. O nome da coluna não deve ser key ou value. Consulte Opções de configuração de tabela.