Acessar dados do Kafka no Cloud Storage

Se você precisar carregar dados de um tópico do serviço gerenciado do Google Cloud para Apache Kafka em um bucket do Cloud Storage, faça isso com um modelo do Dataflow. É possível usar o Google Cloud console, a API REST ou a Google Cloud CLI.

Este documento ajuda a configurar o modelo Kafka para Cloud Storage Dataflow usando o Google Cloud console.

Google Cloud Produtos usados

O modelo Kafka para Cloud Storage Dataflow usa os seguintes produtos faturáveis Google Cloud . Use a Calculadora de preços para gerar uma estimativa de custo com base no uso previsto.

  • Dataflow: Dataflow é um serviço de processamento de dados totalmente gerenciado. O modelo Kafka para Cloud Storage Dataflow usa o Dataflow para criar um pipeline que lê dados do tópico do Kafka, realiza as transformações necessárias e os grava no Cloud Storage. Os recursos de escalonamento automático e auto-recuperação do Dataflow garantem que o pipeline seja executado de maneira confiável e eficiente.
  • Cloud Storage: serve como destino dos dados do Kafka. Você precisa de um bucket do Cloud Storage para armazenar os dados transferidos pelo pipeline do Dataflow.

Além disso, a solução também usa o serviço gerenciado do Google Cloud para Apache Kafka.

  • Serviço gerenciado do Google Cloud para Apache Kafka: um Google Cloud serviço que ajuda você a executar o Apache Kafka. Fornece os dados de origem para o pipeline. Você precisa de um cluster e um tópico do serviço gerenciado para Apache Kafka com dados que você quer transferir para o Cloud Storage. Para mais informações sobre os preços do serviço gerenciado do Google Cloud para Apache Kafka, consulte o guia de preços.

Antes de começar

Antes de iniciar o modelo Kafka para Cloud Storage Dataflow, verifique se você concluiu as seguintes etapas:

  1. Crie um cluster e um tópico do serviço gerenciado para Apache Kafka.

    Uma maneira de criar um cluster e um tópico é seguir o guia de início rápido do Serviço Gerenciado para Apache Kafka.

    Se o tópico contiver registros Avro, consulte Especificar o formato da mensagem para mais requisitos de recursos.

  2. Ative as seguintes Google Cloud APIs:

    • Dataflow

    • Cloud Storage

    gcloud services enable dataflow.googleapis.com storage-api.googleapis.com \
    
  3. Criar um bucket do Cloud Storage.

    Para mais informações sobre como criar um bucket do Cloud Storage, consulte Criar um bucket.

Conceder o papel de cliente do Kafka gerenciado à conta de serviço do worker do Dataflow

Para conectar o job do Dataflow ao Serviço Gerenciado para Apache Kafka, é necessário conceder permissões específicas à conta de serviço do worker do Dataflow. Essa conta de serviço é a identidade usada para todas as VMs de worker no job do Dataflow, e todas as solicitações feitas dessas VMs usam essa conta.

Para permitir o acesso aos recursos do Kafka, conceda o papel roles/managedkafka.client à conta de serviço do worker do Dataflow. Esse papel inclui a permissão managedkafka.clusters.connect necessária para estabelecer conexões.

Para mais informações sobre a conta de serviço do worker, consulte Segurança e permissões para pipelines no Google Cloud.

Para conceder o papel de cliente do Kafka gerenciado à conta de serviço do Dataflow, siga estas etapas:

Console

  1. No Google Cloud console do, acesse a página IAM.
    Acesse o IAM
  2. Verifique se o projeto está definido como o projeto do consumidor que o cliente do serviço gerenciado para Apache Kafka acessaria.
  3. Clique em Conceder acesso.
  4. Na nova página, em Adicionar principais, insira o endereço de e-mail da conta de serviço do worker do Dataflow que você está usando.
  5. Em Atribuir papéis, selecione o papel Cliente do Kafka gerenciado.
  6. Clique em Salvar.

CLI gcloud

  1. No Google Cloud console, ative o Cloud Shell.

    Ativar o Cloud Shell

    Na parte de baixo do Google Cloud console, uma sessão do Cloud Shell é iniciada e exibe um prompt de linha de comando. O Cloud Shell é um ambiente shell com a Google Cloud CLI já instalada e com valores já definidos para o projeto atual. A inicialização da sessão pode levar alguns segundos.

  2. Execute o gcloud projects add-iam-policy-binding comando:

    gcloud projects add-iam-policy-binding PROJECT_ID \
      --member serviceAccount:SERVICE_ACCOUNT_EMAIL \
      --role roles/managedkafka.client

    Substitua:

    • PROJECT_ID é o ID do projeto.

    • SERVICE_ACCOUNT_EMAIL é o endereço de e-mail da conta de serviço do worker do Dataflow.

Iniciar o modelo Kafka para Cloud Storage Dataflow

É possível iniciar o modelo Kafka para Cloud Storage Dataflow na página de detalhes do cluster no console.

  1. Noconsole, acesse a página Cluster. Google Cloud

    Acessar Clusters

    Os clusters criados em um projeto são listados.

  2. Para acessar a página de detalhes do cluster, clique no nome dele.
  3. Na página de detalhes do cluster, clique em Importar dados.

    A página Criar um job do Dataflow usando o modelo "Kafka para Kafka" é aberta.

  4. No modelo, em Modelo do Dataflow, atualize o modelo para Kafka para Cloud Storage.

Configure os campos no modelo de acordo com as informações incluídas nas seções a seguir.

Inserir um nome de job

No campo Nome do job, insira um nome para o job do Dataflow.

O nome precisa ser exclusivo entre todos os jobs em execução no projeto.

Escolher um endpoint regional para o pipeline

No campo Endpoint regional, defina o endpoint regional como o local do cluster do Kafka para minimizar as taxas de transferência de dados entre regiões.

Os workers do Dataflow podem ser executados independentemente da região do cluster do Kafka. No entanto, você incorre em custos de saída entre regiões se iniciar workers fora da região do cluster do Kafka.

Para conferir o local do cluster, siga as etapas em Listar os clusters do serviço gerenciado para Apache Kafka.

Configurar origem

  1. Em Origem, mantenha o valor padrão de Serviço gerenciado para Apache Kafka.

  2. Em Cluster do Kafka e Modo de autenticação de origem do Kafka, mantenha os valores padrão.

  3. Em Tópico do Kafka, selecione um tópico na lista de tópicos disponíveis.

Configurar o formato de mensagem do Kafka

O modelo do Dataflow oferece suporte aos seguintes três formatos de mensagem:

  • Formato de fio do Avro Confluent: cada mensagem do Kafka inclui um byte mágico, um ID de esquema e o registro codificado em binário do Avro.

    Para formatos Avro (formato de fio do Confluent), é possível usar um único esquema ou vários:

    • Esquema único: todas as mensagens seguem um único esquema Avro predefinido.

    • Vários esquemas: as mensagens podem usar esquemas diferentes. Isso só é aceito para o formato de fio do Avro Confluent.

  • Avro (codificado em binário) : as mensagens contêm apenas o payload do registro sem metadados. É necessário fornecer um arquivo de esquema Avro (.avsc) enviado ao Cloud Storage. Todas as mensagens precisam seguir esse esquema único.

  • JSON: os registros não exigem um esquema predefinido. Os registros que não estão em conformidade com o esquema são enviados para a fila de mensagens inativas (se configurada) ou uma mensagem de erro é registrada. O formato aceito é {"field": "value"}. O formato [{"name": "field", "value": "value"}]não é aceito.

O serviço gerenciado do Google Cloud para Apache Kafka não oferece um registro de esquema. O modelo só oferece suporte à transmissão de credenciais de autenticação para registros de esquema compatíveis com o formato de fio do Confluent.

Formato de fio do Avro Confluent

Se você escolher essa opção como o formato de mensagem do Kafka, configure as seguintes configurações adicionais:

Origem do esquema: esse campo informa ao pipeline onde encontrar o esquema. Escolha uma destas opções:

  • Registro de esquema: seus esquemas são armazenados em um registro de esquema do Confluent. Isso é útil para desenvolver esquemas e gerenciar várias versões. Verifique se o registro de esquema está acessível à rede do cluster do serviço gerenciado para Apache Kafka e se está hospedado na mesma região dos workers do Dataflow. É possível usar um registro de esquema com cenários de esquema único e múltiplo. Configure as seguintes configurações adicionais:

    • URL de conexão do registro de esquema: forneça o URL para se conectar ao registro de esquema.

    • Modo de autenticação: se o registro exigir autenticação, selecione OAuth ou TLS. Caso contrário, selecione Nenhum.

  • Arquivo de esquema único: escolha essa opção se todas as mensagens seguirem um esquema único e fixo definido em um arquivo.

    • Arquivo do Cloud Storage para o arquivo de esquema Avro: o caminho para o arquivo de esquema Avro usado para decodificar todas as mensagens em um tópico.

Codificação binária do Avro

Se você escolher essa opção como o formato de mensagem do Kafka, configure as seguintes configurações adicionais:

  • Arquivo do Cloud Storage para o arquivo de esquema Avro: o caminho para o arquivo de esquema Avro usado para decodificar todas as mensagens em um tópico.

JSON

Se você escolher essa opção como o formato de mensagem do Kafka, nenhuma outra configuração será necessária.

Especificar o deslocamento do Kafka

  1. Para evitar o reprocessamento de mensagens quando workers individuais ou todo o pipeline precisarem ser reiniciados, selecione a opção Confirmar deslocamentos para o Kafka. Isso garante que o pipeline retome o processamento de onde parou, evitando o processamento duplicado e possíveis inconsistências de dados.

  2. No campo Inserir ID do grupo de consumidores, insira um nome exclusivo para o grupo desse pipeline. Na maioria das circunstâncias, você quer que o pipeline leia cada mensagem uma vez e possa ser reiniciado.

  3. Para o campo Deslocamento inicial padrão do Kafka, o pipeline do Dataflow oferece duas opções de deslocamento inicial. Selecione uma destas opções:

    • Mais antigo: processa mensagens do início do tópico do Kafka.

    • Mais recente: processa mensagens a partir do deslocamento disponível mais recente.

Configurar destino

Essas opções controlam como o pipeline de dados grava dados no Cloud Storage.

  1. Em Destino, insira o caminho do bucket e inclua o prefixo do nome do arquivo para seus arquivos de saída. O prefixo do arquivo precisa terminar com uma barra. Por exemplo, gs://test-bucket/test-prefix/

  2. Em Duração da janela, insira a janela de tempo para gravar dados no Cloud Storage. Escolha o formato adequado (Ns para segundos, Nm para minutos, Nh para horas) com base nos requisitos de processamento de dados.

  3. Em Prefixo do nome do arquivo de saída dos arquivos a serem gravados, é possível fornecer um prefixo a ser adicionado a cada arquivo de saída para melhor organização e identificação.

  4. Em Número máximo de fragmentos de saída, defina o número como zero. É possível especificar o número de fragmentos a serem produzidos ao gravar arquivos. Aumentar o número pode gerar maior capacidade de processamento, mas também leva a custos maiores devido a custos de redistribuição mais altos. O serviço seleciona um número ideal quando você define o número como zero.

Configurar fila de mensagens inativas

Às vezes, as mensagens não podem ser processadas devido a corrupção, tipos de dados incompatíveis ou incompatibilidades de esquema.

Para lidar com esses casos, ative a fila de mensagens inativas no modelo e forneça um nome de tabela. O modelo cria a tabela usando um esquema padronizado.

Configurar criptografia

Por padrão, todos os dados em repouso e em trânsito são criptografados por um Google-owned and Google-managed encryption key. Se você tiver chaves de criptografia gerenciadas pelo cliente (CMEK), poderá selecionar suas próprias chaves. Para mais informações sobre como configurar uma CMEK, consulte Configurar a criptografia de mensagens.

Configurar rede

É necessário especificar a rede e a sub-rede do cluster no modelo do Dataflow. A seção Parâmetros opcionais do modelo permite definir a rede para os workers do Dataflow.

O modelo Kafka para Cloud Storage Dataflow provisiona workers do Dataflow na rede padrão do projeto, por padrão. Para permitir que o cluster do serviço gerenciado para Apache Kafka envie dados para o Cloud Storage pelo Dataflow, verifique se os workers do Dataflow podem acessar a rede do cluster.

Recomendamos que, se o cluster do Kafka não estiver conectado a uma sub-rede na rede padrão do projeto, use a rede padrão do projeto para o cluster do Kafka.

Para mais informações sobre como configurar a rede com o pipeline do Dataflow, consulte o seguinte:

Se você tiver dificuldades para configurar a rede do Dataflow, consulte o guia de solução de problemas de rede do Dataflow.

Configurar parâmetros opcionais do Dataflow

Configure os parâmetros opcionais somente se você souber o impacto da configuração nos workers do Dataflow. Configurações incorretas podem afetar a performance ou o custo. Para explicações detalhadas de cada opção, consulte Parâmetros opcionais.

Monitoramento

O modelo do Dataflow para Kafka para Cloud Storage oferece uma experiência de monitoramento que permite explorar registros, métricas e erros no Google Cloud Console. Esse conjunto de ferramentas de monitoramento está disponível como parte da interface do usuário do Dataflow.

A guia Métricas do job permite criar painéis personalizados. Para o modelo do Dataflow Kafka para Cloud Storage , recomendamos configurar um painel de Métricas do job que monitore o seguinte:

  • Capacidade de processamento: o volume de dados processados a qualquer momento. Isso é útil para monitorar o fluxo de dados pelo job e identificar possíveis problemas de desempenho.

    Para mais informações, consulte Monitoramento da capacidade de processamento do Dataflow.

  • Atualização de dados: a diferença em segundos entre o carimbo de data/hora no elemento de dados e o horário em que o evento é processado no pipeline. Isso ajuda a identificar gargalos de performance e de origem de dados ou novas tentativas frequentes.

    Para mais informações, consulte Monitoramento da atualização de dados do Dataflow.

  • Backlog: a quantidade de bytes aguardando processamento. Essas informações informam as decisões de escalonamento automático.

Para mais informações sobre o monitoramento do Dataflow, consulte a documentação de monitoramento do Dataflow.

Solução de problemas

Se você encontrar problemas de performance com o pipeline do Dataflow, o Dataflow vai fornecer um conjunto abrangente de ferramentas de solução de problemas e diagnóstico.

Confira dois cenários comuns e os respectivos guias de solução de problemas:

Para uma visão geral da depuração de pipelines do Dataflow, consulte Resolver problemas e depurar pipelines do Dataflow.

Limitações conhecidas

  • O modelo não oferece suporte à transmissão de credenciais para autenticação no registro de esquema.

  • Ao criar o job do Dataflow Kafka para Cloud Storage, verifique se o Google Cloud projeto está definido como o mesmo projeto que contém o cluster do serviço gerenciado para Apache Kafka.

Apache Kafka® é uma marca registrada da The Apache Software Foundation ou afiliadas nos Estados Unidos e/ou em outros países.

A seguir