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
Ative o recurso de conector da sua instância.
Crie um tópico para usar como fonte de dados.
Function Compute
Elasticsearch
Crie um cluster e um índice do Elasticsearch no console do Elasticsearch. Use a versão 7.0 ou posterior para garantir compatibilidade com o cliente do Function Compute (versão 7.7.0).
Adicione o bloco CIDR do endpoint do Function Compute à lista de permissões do Elasticsearch. Para testes iniciais, especifique
0.0.0.0/0para permitir todos os endereços IP na VPC; depois, restrinja o intervalo após verificar a conectividade.
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) |
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
Faça login no console do ApsaraMQ for Kafka.
Na seção Resource Distribution da página Overview, selecione sua região.
No painel de navegação à esquerda, clique em Connectors.
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 |
kafka-elasticsearch-sink |
|
Instance |
Exibe o nome e o ID da instância atual. |
demo alikafka_post-cn-st21p8vj**** |
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âmetro | Descrição | Exemplo |
|---|---|---|
| VPC ID | VPC da instância de origem. Preenchido automaticamente; nenhuma alteração necessária. | vpc-bp1xpdnd3l*** |
| vSwitch ID | vSwitch da instância de origem. Deve estar na mesma VPC. | vsw-bp1d2jgg81*** |
| Failure Handling Policy | Açã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 Method | Auto cria os tópicos internos necessários automaticamente. Manual permite que você mesmo os crie. | Auto |
| Connector Consumer Group | Grupo 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 |
** |
|
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 |
-- |
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 |
|
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 |
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.
Na página Connectors, localize o conector. Na coluna Actions, escolha More > Configure Function.
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
Na página Connectors, localize o conector e clique em Test na coluna Actions.
No painel Send Message, defina Method of Sending como Console.
No campo Message Key, insira uma chave, por exemplo,
demo.-
No campo Message Content, insira um corpo JSON, por exemplo:
{"key": "test"} 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
-
Execute a consulta a seguir para pesquisar no índice de destino:
GET /<index_name>/_search -
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.