Todos os produtos
Search
Central de documentação

Elasticsearch:Use o Real-time Computing para processar e sincronizar dados no ES

Última atualização: Jun 27, 2026

Este tutorial demonstra como construir um pipeline de recuperação de logs que lê dados do Alibaba Cloud Simple Log Service (SLS), processa-os com o Realtime Compute for Apache Flink usando Flink SQL e grava os resultados no Elasticsearch para busca full-text.

Pré-requisitos

Antes de começar, verifique se você:

Como funciona

Os dados de log do SLS fluem para o Realtime Compute for Apache Flink. O Flink aplica as transformações SQL e grava os registros processados no Elasticsearch pela API REST. O Elasticsearch indexa cada registro como um documento, tornando-o imediatamente pesquisável.

Flink+ES data link

Nota

O conector sink do Elasticsearch usa APIs REST e é compatível com todas as versões do Elasticsearch. O Realtime Compute for Apache Flink versão 3.2.2 e posteriores oferece suporte a tabelas sink do Elasticsearch.

Configure o pipeline Flink SQL

Etapa 1: Crie um job do Flink

  1. Faça login no console do Realtime Compute for Apache Flink.

  2. Crie um job do Realtime Compute for Apache Flink. Para mais informações, consulte a seção Job Development > Development do Blink SQL Development Guide no documento do Alibaba Cloud Blink Exclusive Mode (Descontinuado para Alibaba Cloud).

Etapa 2: Crie a tabela source do SLS

Defina uma tabela source para leitura do SLS:

create table sls_stream(
  a int,
  b int,
  c VARCHAR
)
WITH (
  type ='sls',
  endPoint ='<yourEndpoint>',
  accessId ='<yourAccessId>',
  accessKey ='<yourAccessKey>',
  startTime = '<yourStartTime>',
  project ='<yourProjectName>',
  logStore ='<yourLogStoreName>',
  consumerGroup ='<yourConsumerGroupName>'
);

Substitua os placeholders pelos valores reais:

Parâmetro

Obrigatório

Descrição

endPoint

Sim

Endpoint público do SLS. Exemplo: http://cn-hangzhou.log.aliyuncs.com. Deve começar com http://. Para a lista completa de endpoints, consulte Endpoints.

accessId

Sim

Seu AccessKey ID.

accessKey

Sim

Seu AccessKey secret.

startTime

Sim

Momento inicial para consumo de logs. Ao executar o job, o horário selecionado deve ser posterior a este valor.

project

Sim

Nome do projeto do SLS.

logStore

Sim

Nome do Logstore no projeto.

consumerGroup

Sim

Nome do grupo de consumidores do SLS.

Etapa 3: Crie a tabela sink do Elasticsearch

Defina uma tabela sink para gravação no Elasticsearch:

CREATE TABLE es_stream_sink(
  a int,
  cnt BIGINT,
  PRIMARY KEY(a)
)
WITH(
  type ='elasticsearch-7',
  endPoint = 'http://<instanceid>.public.elasticsearch.aliyuncs.com:<port>',
  accessId = '<yourAccessId>',
  accessKey = '<yourAccessSecret>',
  index = '<yourIndex>',
  typeName = '<yourTypeName>'
);

Substitua os placeholders pelos valores reais:

Parâmetro

Obrigatório

Padrão

Descrição

endPoint

Sim

Endpoint público da instância do Elasticsearch, no formato http://<instanceid>.public.elasticsearch.aliyuncs.com:9200. Disponível na página de informações básicas da instância. Consulte Visualize as informações básicas de uma instância.

accessId

Sim

elastic

Nome de usuário da instância do Elasticsearch. O padrão é elastic.

accessKey

Sim

Senha do usuário, definida na criação da instância. Se esquecer a senha, consulte Redefinir a senha de acesso de uma instância.

index

Sim

Nome do índice. Crie o índice antes de executar o job ou ative a criação automática. Consulte Guia para iniciantes: Da criação da instância à recuperação de dados e Configure parâmetros YML.

typeName

Sim

Tipo de índice. Para instâncias do Elasticsearch V7.0 ou posteriores, use _doc.

Etapa 4: Escrever a lógica de negócios

Insira os resultados agregados na tabela sink:

INSERT INTO es_stream_sink
SELECT
  a,
  count(*) as cnt
FROM sls_stream GROUP BY a

Etapa 5: Publicar e execute o job

Publique e execute o job. O sistema agrega os dados do SLS e os grava no Elasticsearch.

Tratamento de chaves e modos de atualização

A gravação de documentos no Elasticsearch pelo Flink depende da definição de chave primária no DDL da tabela sink.

Modo

Gatilho

ID do Documento

Upsert

PRIMARY KEY definida

Valor do campo da chave primária

Append

Sem PRIMARY KEY

Geração aleatória de ID de documento pelo sistema

Importante

Especifique apenas um campo como PRIMARY KEY.

No modo upsert, o parâmetro updateMode controla a atualização de documentos existentes:

Valor de updateMode

Comportamento

full

Os registros recebidos substituem completamente os documentos existentes.

inc

Atualiza apenas os campos presentes no registro recebido; os demais permanecem inalterados.

Por padrão, todas as atualizações usam semântica upsert (INSERT ou UPDATE). Para mais informações, consulte Index API.

Próximos passos

Para implementar lógica de gravação personalizada além do suporte do conector integrado do Elasticsearch, use o recurso de sink personalizado do Realtime Compute for Apache Flink.