Todos os produtos
Search
Central de documentação

ApsaraMQ for Kafka:Create an Elasticsearch sink connector

Última atualização: Jun 27, 2026

Um conector de sink do Elasticsearch lê mensagens de um tópico na instância do ApsaraMQ for Kafka e as grava em um índice do Elasticsearch. O conector usa o Function Compute como intermediário: consome mensagens do tópico de origem, encaminha-as para uma função do Function Compute e essa função as grava no Elasticsearch pela Bulk API. Cada mensagem torna-se um documento no índice de destino, com metadados como tópico, partição, offset e timestamp.

Antes de começar

Conclua as configurações a seguir antes de criar o conector.

ApsaraMQ for Kafka

Function Compute

Elasticsearch

Informações a coletar

Reúna os detalhes abaixo antes de iniciar o assistente.

Informação

Onde encontrar

Exemplo

ID da instância do Elasticsearch

Console do Elasticsearch

es-cn-oew1o67x0000****

Endpoint do Elasticsearch (público ou privado)

Informações básicas do cluster

es-cn-oew1o67x0000****.elasticsearch.aliyuncs.com

Porta do Elasticsearch

Informações básicas do cluster

9200 (HTTP/HTTPS) ou 9300 (TCP)

Nome de usuário e senha do Elasticsearch

Definidos durante a criação do cluster; redefina se necessário

elastic / ****

Nome do índice do Elasticsearch

Console do Elasticsearch

elastic_test

Nome do tópico de origem

Console do ApsaraMQ for Kafka

elasticsearch-test-input

Limites

  • A instância do ApsaraMQ for Kafka e o cluster do Elasticsearch devem estar na mesma região.

  • O ApsaraMQ for Kafka serializa mensagens como strings UTF-8. Não há suporte para dados binários.

  • Se você especificar o endpoint privado do cluster do Elasticsearch, o Function Compute não poderá acessá-lo por padrão. Para habilitar a conectividade, configure o serviço do Function Compute para usar a mesma VPC e o mesmo vSwitch do cluster do Elasticsearch. Consulte Configurar o serviço do Function Compute.

  • Para outros limites de conectores, consulte Limites.

Faturamento

O conector usa o Function Compute para exportar dados. O Function Compute oferece um nível gratuito. O uso que exceder esse nível será cobrado conforme descrito em Faturamento do Function Compute.

Criar e implantar o conector

  1. Faça login no console do ApsaraMQ for Kafka.

  2. Na seção Resource Distribution da página Overview, selecione sua região.

  3. No painel de navegação à esquerda, clique em Connectors.

  4. Na página Connectors, selecione sua instância na lista suspensa Select Instance e clique em Create Connector.

Etapa 1: Configurar informações básicas

Na etapa Configure Basic Information, defina o nome do conector e revise os detalhes da instância.

Parâmetro

Descrição

Exemplo

Name

Nome exclusivo na instância do ApsaraMQ for Kafka. Use de 1 a 48 caracteres: dígitos, letras minúsculas e hifens (-). Não pode começar com hífen. O conector cria automaticamente um grupo de consumidores chamado connect-<connector-name>.

kafka-elasticsearch-sink

Instance

Exibe o nome e o ID da instância atual.

demo alikafka_post-cn-st21p8vj****

Importante

Por padrão, a opção Authorize to Create Service Linked Role vem selecionada. O ApsaraMQ for Kafka cria uma função vinculada ao serviço caso ainda não exista.

Clique em Next.

Etapa 2: Configurar o serviço de origem

Na etapa Configure Source Service, selecione Message Queue for Apache Kafka como serviço de origem e configure os parâmetros a seguir.

Parâmetro

Descrição

Exemplo

Data Source Topic

Tópico do qual os dados são consumidos.

elasticsearch-test-input

Consumer Thread Concurrency

Número de threads de consumidor simultâneas. Valores válidos: 1, 2, 3, 6, 12. Padrão: 6.

6

Consumer Offset

Ponto de início do consumo. Earliest Offset lê desde o início. Latest Offset lê apenas novas mensagens.

Earliest Offset

Clique em Configure Runtime Environment para expandir parâmetros adicionais.

ParâmetroDescriçãoExemplo
VPC IDVPC da instância de origem. Preenchido automaticamente; nenhuma alteração necessária.vpc-bp1xpdnd3l***
vSwitch IDvSwitch da instância de origem. Deve estar na mesma VPC.vsw-bp1d2jgg81***
Failure Handling PolicyAção a ser tomada quando uma mensagem falhar. Continue Subscription registra o erro e continua o consumo. Stop Subscription registra o erro e interrompe a partição. Consulte Gerenciar um conector para detalhes sobre logs e Códigos de erro para solução de problemas.
Nota
Continue Subscription
Resource Creation MethodAuto cria os tópicos internos necessários automaticamente. Manual permite que você mesmo os crie.Auto
Connector Consumer GroupGrupo de consumidores para a tarefa do conector. Formato: connect-<connector-name>.connect-kafka-elasticsearch-sink

Tópicos internos (apenas para criação manual)

Se você definir o Resource Creation Method como Manual, crie os tópicos a seguir. Todos os tópicos que exigem Local Storage estão disponíveis apenas em instâncias da Professional Edition.

Parâmetro

Convenção de nomenclatura

Partições

Mecanismo de armazenamento

**cleanup.policy**

Task Offset Topic

connect-offset-*

Mais de 1

Local Storage

Compact

Task Configuration Topic

connect-config-*

Exatamente 1

Local Storage

Compact

Task Status Topic

connect-status-*

6 (recomendado)

Local Storage

Compact

Dead-letter Queue Topic

connect-error-*

6 (recomendado)

Local Storage ou Cloud Storage

--

Error Data Topic

connect-error-*

6 (recomendado)

Local Storage ou Cloud Storage

--

Nota

Para economizar recursos de tópicos, use o mesmo tópico tanto para a fila de mensagens mortas quanto para o tópico de dados de erro.

Clique em Next.

Etapa 3: Configurar o serviço de destino

Na etapa Configure Destination Service, selecione Elasticsearch como serviço de destino e configure os parâmetros a seguir.

Parâmetro

Descrição

Exemplo

Elasticsearch Instance ID

ID do cluster do Elasticsearch.

es-cn-oew1o67x0000****

Endpoint

Endpoint público ou privado do cluster. Consulte Visualizar informações básicas do cluster.

es-cn-oew1o67x0000****.elasticsearch.aliyuncs.com

Port

9200 para HTTP/HTTPS ou 9300 para TCP.

9300

Username

Nome de usuário do Elasticsearch. Padrão: elastic. Personalize pelo X-Pack RBAC se necessário. A conta deve ter permissões de gravação no índice de destino.

elastic

Password

Senha definida durante a criação do cluster. Redefina a senha caso a tenha esquecido.

****

Index

Nome do índice do Elasticsearch de destino.

elastic_test

Nota
  • O nome de usuário e a senha são passados ao Function Compute como variáveis de ambiente quando a tarefa do conector é criada. O ApsaraMQ for Kafka não armazena essas credenciais após a criação da tarefa.

  • A conta precisa ter permissões para gravar no índice, pois as mensagens são enviadas pela Bulk API do Elasticsearch.

Clique em Create.

Implantar o conector

Após a criação, o conector aparece na página Connectors. Clique em Deploy na coluna Actions para iniciar o conector.

Configurar o serviço do Function Compute

Depois de implantar o conector, o Function Compute cria automaticamente um serviço chamado kafka-service-<connector-name>-<random-string>. Se o cluster do Elasticsearch usar um endpoint privado, configure o serviço do Function Compute para utilizar a mesma VPC e o mesmo vSwitch do cluster do Elasticsearch.

  1. Na página Connectors, localize o conector. Na coluna Actions, escolha More > Configure Function.

  2. No console do Function Compute, localize o serviço criado automaticamente e atualize as configurações de VPC e vSwitch para corresponder às do cluster do Elasticsearch.

Verificar o fluxo de dados

Envie uma mensagem de teste para confirmar que os dados fluem do ApsaraMQ for Kafka para o Elasticsearch.

Enviar uma mensagem de teste

  1. Na página Connectors, localize o conector e clique em Test na coluna Actions.

  2. No painel Send Message, defina Method of Sending como Console.

  3. No campo Message Key, insira uma chave, por exemplo, demo.

  4. No campo Message Content, insira um corpo JSON, por exemplo:

       {"key": "test"}
  5. Para Send to Specified Partition, clique em Yes e insira um Partition ID (por exemplo, 0) para direcionar a uma partição específica, ou clique em No para deixar o sistema atribuir uma. Para consultar IDs de partição, veja Visualizar status da partição.

Também é possível enviar mensagens de teste via Docker ou um SDK. Selecione a opção correspondente no campo Method of Sending e siga as instruções na tela.

Verificar o índice do Elasticsearch

  1. Faça login no console do Kibana.

  2. Execute a consulta a seguir para pesquisar no índice de destino:

       GET /<index_name>/_search
  3. Confirme se a resposta contém a mensagem enviada. Uma resposta bem-sucedida tem a seguinte aparência:

       {
         "took": 8,
         "timed_out": false,
         "_shards": {
           "total": 5,
           "successful": 5,
           "skipped": 0,
           "failed": 0
         },
         "hits": {
           "total": {
             "value": 1,
             "relation": "eq"
           },
           "max_score": 1.0,
           "hits": [
             {
               "_index": "product_****",
               "_type": "_doc",
               "_id": "TX3TZHgBfHNEDGoZ****",
               "_score": 1.0,
               "_source": {
                 "msg_body": {
                   "key": "test",
                   "offset": 2,
                   "overflowFlag": false,
                   "partition": 2,
                   "timestamp": 1616599282417,
                   "topic": "dv****",
                   "value": "test1",
                   "valueSize": 8
                 },
                 "doc_as_upsert": true
               }
             }
           ]
         }
       }

Solução de problemas

Sintoma

Causa possível

Resolução

Falha do conector ao gravar no Elasticsearch

O Function Compute não consegue alcançar o cluster do Elasticsearch

Verifique se o serviço do Function Compute usa a mesma VPC e o mesmo vSwitch do cluster do Elasticsearch. Consulte Configurar o serviço do Function Compute.

Mensagens não são consumidas

Configuração incorreta do offset do consumidor

Verifique a configuração de Consumer Offset. Use Earliest Offset para ler todas as mensagens existentes ou Latest Offset apenas para novas mensagens.

Erros de autenticação

Credenciais inválidas ou permissões insuficientes

Confirme se o nome de usuário e a senha estão corretos e se a conta tem permissões de gravação no índice de destino.

Para logs de chamadas de função do Function Compute, consulte Configurar o recurso de log.