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ê:
-
Ativou o Realtime Compute for Apache Flink e criou um projeto
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).
Criou uma instância do Alibaba Cloud Elasticsearch. Para mais informações, consulte Criar um cluster do Alibaba Cloud Elasticsearch
Ativou o SLS e criou um projeto e um Logstore. Para mais informações, consulte Ativar o Alibaba Cloud Simple Log Service, Gerenciar projetos e Criar um Logstore básico
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.

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
Faça login no console do Realtime Compute for Apache Flink.
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 |
|
|
Sim |
Endpoint público do SLS. Exemplo: |
|
|
Sim |
Seu AccessKey ID. |
|
|
Sim |
Seu AccessKey secret. |
|
|
Sim |
Momento inicial para consumo de logs. Ao executar o job, o horário selecionado deve ser posterior a este valor. |
|
|
Sim |
Nome do projeto do SLS. |
|
|
Sim |
Nome do Logstore no projeto. |
|
|
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 |
|
|
Sim |
— |
Endpoint público da instância do Elasticsearch, no formato |
|
|
Sim |
|
Nome de usuário da instância do Elasticsearch. O padrão é |
|
|
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. |
|
|
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. |
|
|
Sim |
— |
Tipo de índice. Para instâncias do Elasticsearch V7.0 ou posteriores, use |
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 |
|
Valor do campo da chave primária |
|
Append |
Sem |
Geração aleatória de ID de documento pelo sistema |
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 |
|
|
Os registros recebidos substituem completamente os documentos existentes. |
|
|
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.