Conector Apache Kafka (DataStage)

Use o conector Apache Kafka em DataStage® para gravar e ler fluxos de eventos de e para tópicos.

Pré-requisito

Crie a conexão Para obter instruções, consulte Conectando a uma origem de dados no DataStage e Conexão Apache Kafka.

Lendo dados dos tópicos no Kafka

É possível configurar o conector Apache Kafka para ler mensagens de uma origem de dados Kafka .

Figura 1. Exemplo de leitura de dados do conector Apache Kafka
Exemplo de leitura de dados do conector Apache Kafka
Configurando o conector Apache Kafka como uma origem
  1. Na tela de design da tarefa, dê um clique duplo no conector Apache Kafka
  2. Na guia Estágio , clique em Propriedades.
  3. Selecione Usar DataStage propriedades.
  4. Insira os valores para as propriedades necessárias:
    • Nome do tópico - Kafka tópico do qual ler mensagens.
    • Serializador de chave -O serializador dos tipos de dados de chave em Apache Kafka. O valor do serializador de Chave deve ser compatível com a coluna de chave (se definida) no link Saída . Se você selecionar um serializador de chaves Avro ou Avro as JSON , o serializador de valores deverá ser configurado como Avro ou Avro as JSON
    • serializador de valor -O serializador para os tipos de dados de valor no Apache Kafka. O valor do serializador de Valor deve ser compatível com a coluna de chave (se definida) no link Saída .
    • Grupo de consumidores -uma sequência que identifica exclusivamente o grupo de processos do consumidor ao qual este consumidor pertence.. Ao configurar o mesmo ID do grupo para diversos processos, você indica que todos eles fazem parte do mesmo grupo de consumidores
    • Máximo de registros de pesquisa -O número máximo de registros que são retornados em uma única chamada para a pesquisa
    • Máximo de mensagens -O número máximo de mensagens para consumir ou produzir de ou para o tópico em uma base de processo do reprodutor. Esse valor é um múltiplo do Máximo de registros de pesquisa. Essa configuração não estará ativa se você selecionar Modo contínuo.
    • Reconfigurar política -O valor que indica a política se não houver deslocamento inicial no servidor Kafka ou se o deslocamento atual não existir mais no servidor. Por exemplo, se os dados forem excluídos, será possível usar os valores earliest e latest para reconfigurar automaticamente o deslocamento para o mais antigo do mais recente.
  5. Configurar propriedades opcionais:
    • Modo contínuo -O conector buscará mensagens em um modo contínuo aguardando novas mensagens recebidas indefinidamente.
      • Mensagem de parada -A mensagem que para o processamento de modo contínuo.
    • Transação -A transação que envia as configurações do marcador "Fim de onda":
      • Contagem de registros -O número de registros por transação. O valor 0 significa todos os registos disponíveis.
      • Intervalo de tempo -O intervalo de tempo para uma transação Use esta propriedade se você configurar a Contagem de registros como 0..
    • Fim do Grupo -Envie um marcador "Fim do Grupo" após cada transação.
    • Fim dos dados -Insira um marcador "Fim do grupo" para o conjunto final de registros quando seu número for menor que o valor especificado para a contagem de registros de transação. Se o valor da transação Contagem de registros for 0 (todos os registros disponíveis), haverá apenas um grupo de transações. Nesse caso, configure o valor End of data para Yes para que o marcador "End of wave" seja inserido para esse grupo de transações.
    • Nível de criação de log -configura a criação de log do cliente Kafka . Consulte Controlando a saída do log
    • Propriedades de configuração do cliente -Define propriedades adicionais do cliente Kafka . Consulte Kafka propriedades de configuração do cliente
  6. Clique na guia Saída para definir as colunas do conector. As colunas devem ser compatíveis com o tipo de dados definido nas propriedades serializador de chave e serializador de valor . Se você usar um serializador Avro ou Avro como JSON , os nomes de colunas deverão corresponder aos campos de registro Avro. Consulte Registro de esquema para obter detalhes.
    • value -A coluna que contém os valores de mensagem do Kafka Ele deve ser compatível com a propriedade serializador de valor
    • key -A coluna que contém as chaves de mensagem do Kafka Ele deve ser compatível com a propriedade Key serializer
    • partition -A coluna que contém o número de partição do qual a mensagem Kafka vem. Ele deve ser compatível com o tipo de dados Integer.
    • timestamp -a coluna que contém o registro de data e hora para a mensagem Kafka . Ele deve ser compatível com o tipo de dados Timestamp.
    • offset -A coluna que contém o offset para a mensagem Kafka . Deve ser compatível com o tipo de dado BigInt .

Gravando dados em tópicos no Kafka ..

É possível configurar o conector Apache Kafka para gravar dados em uma origem de dados Kafka .

Figura 2 Exemplo de gravação de dados para a origem de dados Kafka
Exemplo de gravação de dados para a origem de dados Kafka
Configurando o conector Apache Kafka como um destino
  1. Na tela de design da tarefa, dê um clique duplo no conector Apache Kafka
  2. Na guia Estágio , clique em Propriedades.
  3. Selecione Usar DataStage propriedades.
  4. Insira os valores para as propriedades necessárias:
    • Nome do tópico -O nome do tópico no qual as mensagens devem ser gravadas no estágio de envio de dados.
    • Serializador de chave -O serializador dos tipos de dados de chave em Apache Kafka. O valor do serializador de chave deve ser compatível com a coluna-chave (se definida) no link Entrada . Se você selecionar um serializador de chaves Avro ou um Avro as JSON , o serializador de valores deverá ser configurado como Avro ou como Avro as JSON
    • serializador de valor -O serializador para os tipos de dados de valor no Apache Kafka. O valor do serializador Value deve ser compatível com a coluna-chave (se definida) no link Entrada .
  5. Configurar propriedades opcionais:
    • Nível de criação de log -configura a criação de log do cliente Kafka . Consulte Controlando a saída do log
    • Propriedades de configuração do cliente -Define propriedades adicionais do cliente Kafka . Consulte as propriedades de configuração do cliente Kafka
  6. Clique na guia Entrada para definir as colunas do conector. As colunas devem ser compatíveis com o tipo de dado definido nas propriedades serializador de chave e serializador de valor . Se você usar um serializador Avro ou um Avro as JSON , os nomes de colunas deverão corresponder aos campos de registro Avro Consulte Registro de esquema para obter detalhes.
    • value -A coluna que contém os valores de mensagem do Kafka Ele deve ser compatível com a propriedade serializador de valor
    • key -A coluna que contém as chaves de mensagem do Kafka Ele deve ser compatível com a propriedade Key serializer

Registro de esquema

O registro de esquema fornece a capacidade de usar mensagens mais complexas do Kafka . É possível usar o registro de esquema para definições gerais da coluna DataStage em designs de tarefa.

Por padrão, as mensagens são retornadas como um valor dentro de uma única coluna definida no conector Apache Kafka . Ao usar os formatos Avro e Avro as JSON e o registro de esquema, você ativa a decomposição de mensagens Kafka complexas nas colunas DataStage .

Quando você grava em um tópico, o conector cria um esquema baseado nas definições de coluna.

Quando o conector lê a partir de um tópico Kafka , o conector obtém um esquema do registro de esquema. Com base nesse esquema, os campos de registro Avro são extraídos em colunas que são definidas como saída do conector..

Por exemplo, quando o esquema Avro é definido como o seguinte:
{
  "fields": [
    {
      "name": "id",
      "type": "int"
    },
    {
      "name": "name",
      "type": "string"
    },
    {
      "name": "salary",
      "type": "double"
    }
  ],
  "name": "KafkaConnectorSchema",
  "type": "record"
}
No conector Apache Kafka , defina as colunas conforme a seguir:
Tabela 1. Nome da coluna e tipo de dados
Nome da coluna Tipo de dados
ID Número Inteiro
nome VarChar
salário Duplo

Para usar o registro de esquema, use a seguinte configuração:

  • Você deve configurar a propriedade de configuração serializador de chave ou serializador de valor como Avro ou como Avro as JSON.
  • Depois de configurar o serializador de chave para Avro ou Avro as JSON, não é possível configurar o serializador de valor para um valor diferente de Avro ou Avro as JSON.
  • Ao configurar a propriedade de configuração do Serializador de Chaves como Avro, a definição de coluna para o link deve estar em conformidade com a definição do esquema Avro
  • As colunas de chave devem ser selecionadas com a caixa de seleção Chave ou seus nomes devem iniciar com um prefixo key_ . Use apenas um método Não é possível usar uma caixa de seleção Chave para uma coluna e um prefixo key_ para outra coluna.

Propriedades de configuração do cliente Kafka

Use as Propriedades de Configuração do Cliente (em Configuração Avançada do Kafka na guia Estágio ) quando precisar configurar uma propriedade de configuração do cliente Kafka , mas o conector Apache Kafka não tiver um diálogo de configuração para ele.

Insira a propriedade com o formato key=value com entradas definidas em cada linha.. Use o caractere # para obter comentários Para obter uma descrição das regras para as propriedades de configuração do cliente, consulte a definição de Propriedades Java

Resolução de problemas

Controlando saída de log

Criação de log do cliente Kafka

Use as propriedades na Configuração avançada do Kafka para controlar a criação de log do cliente do Kafka As entradas de log do cliente Kafka serão gravadas no log da tarefa DataStage juntamente com as entradas de log para a saída dos estágios.

  • Nível de criação de log: configure o nível mínimo de mensagens produzidas por cliente Kafka para gravar no log da tarefa.
  • Entradas do log de aviso e de erro: mude o valor para Log como informativo ou para Log como Aviso para controlar se alguma mensagem ERROR ou FATAL será gravada no log da tarefa com uma severidade que possa causar uma falha da tarefa.

Para visualizar mensagens Debug ou Trace , configure o nível de criação de log associado na variável de ambiente CC_MSG_LEVEL .

Mensagens de erro

Tabela 2. Resolução de problemas de erros
Tipo do erro Detalhes do erro Solução
KafkaProducer -Falha ao atualizar metadados após 20000 ms O cliente Kafka não relata o problema de conexão quando há uma configuração incorreta entre a autenticação configurada no lado do broker do Kafka e no conector Apache Kafka . Por exemplo, o conector Apache Kafka é configurado para conectar-se sem nenhuma configuração de segurança, mas a conexão é feita na porta que é protegida com SSL Quando o conector Apache Kafka é configurado para ler mensagens do Kafka, a tarefa para de responder na etapa de conexão e o cliente Kafka não relata um erro. Deve-se parar a tarefa manualmente Assegure-se de que as configurações de conexão (nome do host, porta, tipo de conexão segura e propriedades) correspondam à configuração de conexão do broker Kafka .
Confirmação não pode ser concluída Quando o conector Apache Kafka é configurado para executar no modo Contínuo, a tarefa pode falhar com esta mensagem:

Kafka Connector was unable to commit messages due to following error: Commit cannot be completed since the group has already rebalanced and assigned the partitions to another member. This means that the time between subsequent calls to poll() was longer than the configured max.poll.interval.ms, which typically implies that the poll loop is spending too much time message processing. You can address this either by increasing the session timeout or by reducing the maximum size of batches returned in poll() with max.poll.records.
Na seção Avançado do estágio, mude o modo de execução para Sequencial
Chave Nula O arquivo keytab está errado ou as credenciais do usuário principal não existem no arquivo keytab. Assegure que o arquivo keytab esteja correto.
Exceção do analisador SAX lançada: A entrada terminou antes que todas as tags iniciadas fossem terminadas. A última tag iniciada foi 'AdvancedKafkaConfigOptions' O valor de Propriedades de configuração do cliente pode conter várias linhas de entradas de valor da chave Como cada propriedade pode ser parametrizada com o formato #name# , essa propriedade também é verificada se um parâmetro está presente. Ao mesmo tempo, o valor das propriedades do cliente Kafka deve estar em conformidade com os requisitos de classe Propriedades Java . O caractere # pode ser interpretado como um comentário ou como o início de um parâmetro de tarefa Quando o caractere # aparecer na última linha dessa propriedade, o pré-processador assumirá que o caractere corresponde a um segundo caractere # e falhará se não for localizado. Em Propriedades de configuração do cliente; certifique-se de que pelo menos uma nova linha vazia esteja presente no final de cada propriedade