Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Elasticsearch

Última atualização: Jul 03, 2026

Este tópico explica como usar o conector Elasticsearch.

Contexto

O Alibaba Cloud Elasticsearch é compatível com o Elasticsearch open source e inclui recursos comerciais como Security, Machine Learning, Graph e APM para análise e busca de dados. Ele oferece serviços de nível empresarial, incluindo controle de acesso, monitoramento e alertas de segurança, além de geração automatizada de relatórios.

A tabela a seguir descreve as capacidades do conector Elasticsearch.

Item

Descrição

Tipo de tabela

Tabela de origem, tabela de dimensão e tabela de destino

Modo de execução

Modo em lote e modo streaming

Formato de dados

JSON

Métrica

Métrica

  • Tabela de origem

    • pendingRecords

    • numRecordsIn

    • numRecordsInPerSecond

    • numBytesIn

    • numBytesInPerSecond

  • Tabela de dimensão

    Nenhuma

  • Tabela de destino (para Ververica Runtime (VVR) 6.0.6 e posterior)

    • numRecordsOut

    • numRecordsOutPerSecond

Nota

Para obter mais informações sobre essas métricas, consulte Métricas.

Tipo de API

DataStream API e SQL

Atualização ou exclusão de dados em uma tabela de destino

Suportado

Pré-requisitos

Limitações

  • As tabelas de origem e de dimensão suportam Elasticsearch 6.8.x ou superior.

    Nota

    O uso do Elasticsearch 8.x com tabelas de origem e de dimensão requer VVR 11.6 ou posterior.

  • Tabelas de destino suportam apenas Elasticsearch 6.x, 7.x e 8.x.

  • Somente tabelas de origem Elasticsearch completas são suportadas; tabelas incrementais não são aceitas.

Sintaxe

  • Tabela de origem

    Elasticsearch 8.x

    CREATE TABLE elasticsearch_source(
      name STRING,
      location STRING,
      value FLOAT
    ) WITH (
      'connector' ='elasticsearch-8',
      'hosts' = '<yourHosts>',
      'index' = '<yourIndex>'
    );

    Outras versões

    CREATE TABLE elasticsearch_source(
      name STRING,
      location STRING,
      value FLOAT
    ) WITH (
      'connector' ='elasticsearch',
      'endPoint' = '<yourEndPoint>',
      'indexName' = '<yourIndexName>'
    );
  • Tabela de dimensão

    Elasticsearch 8.x

    CREATE TABLE es_dim(
      field1 STRING, -- Must be of the STRING type when used as a key for a JOIN.
      field2 FLOAT,
      field3 BIGINT,
      PRIMARY KEY (field1) NOT ENFORCED
    ) WITH (
      'connector' ='elasticsearch-8',
      'hosts' = '<yourHosts>',
      'index' = '<yourIndex>'
    );

    Outras versões

    CREATE TABLE es_dim(
      field1 STRING, -- Must be of the STRING type when used as a key for a JOIN.
      field2 FLOAT,
      field3 BIGINT,
      PRIMARY KEY (field1) NOT ENFORCED
    ) WITH (
      'connector' ='elasticsearch',
      'endPoint' = '<yourEndPoint>',
      'indexName' = '<yourIndexName>'
    );
    Nota
    • Se uma chave primária for especificada, apenas um campo poderá ser a chave de junção, e este deve corresponder ao ID do documento no índice Elasticsearch associado.

    • Caso nenhuma chave primária seja definida, use um ou mais campos como chaves de junção. Esses campos devem existir nos documentos correspondentes do Elasticsearch.

    • Para campos STRING, o conector adiciona o sufixo .keyword aos nomes dos campos por padrão para garantir compatibilidade. Se isso impedir a correspondência com campos TEXT no Elasticsearch, defina a opção ignoreKeywordSuffix como true.

  • Tabela de destino

    CREATE TABLE es_sink(
      user_id   STRING,
      user_name   STRING,
      uv BIGINT,
      pv BIGINT,
      PRIMARY KEY (user_id) NOT ENFORCED
    ) WITH (
      'connector' = 'elasticsearch-7', -- If you use Elasticsearch 6.x, set this to 'elasticsearch-6'.
      'hosts' = '<yourHosts>',
      'index' = '<yourIndex>'
    );
    Nota
    • Uma tabela de destino Elasticsearch opera em upsert mode ou append mode, dependendo da definição de uma chave primária.

      • Quando uma chave primária é definida, seu valor serve como ID do documento. A tabela de destino opera então em upsert mode, processando operações UPDATE e DELETE.

      • Na ausência de uma chave primária, o Elasticsearch gera automaticamente um ID de documento aleatório. Nesse caso, a tabela de destino funciona em append mode, consumindo apenas mensagens INSERT.

    • Tipos de dados como BYTES, ROW, ARRAY e MAP não possuem representação em string correspondente. Portanto, não é possível usar campos desses tipos como chave primária.

    • Os campos no DDL correspondem aos campos de um documento Elasticsearch. Não é possível gravar metadados, como o ID do documento, na tabela de destino, pois o cluster Elasticsearch gerencia esses metadados internamente.

Opções WITH

Tabela de origem

Parâmetro

Descrição

Tipo

Obrigatório

Padrão

Observações

connector

Tipo da tabela de origem.

String

Sim

Nenhum

Valores válidos: elasticsearch ou elasticsearch-8.

Nota

Apenas VVR 11.6 ou posterior suporta o valor elasticsearch-8.

endPoint

Endereço do servidor do cluster Elasticsearch.

String

Sim

Nenhum

Nome legado da opção.

hosts

Usado com elasticsearch-8.

indexName

Nome do índice.

String

Sim

Nenhum

Nome legado da opção.

index

Para uso com elasticsearch-8.

accessId

Nome de usuário para autenticação.

String

Não

Nenhum

Por padrão, este parâmetro está vazio e a autenticação não ocorre. Ao especificar accessId, informe um accessKey não vazio.

Nota

O elasticsearch-8 foi atualizado para usar os parâmetros username e password.

Importante

Para evitar a exposição do seu nome de usuário e senha, recomendamos o uso de variáveis de projeto. Para mais informações, consulte variáveis de projeto.

username

accessKey

Senha para autenticação.

String

Não

Nenhum

password

typeNames

Nome do tipo.

String

Não

_doc

Recomenda-se não configurar esta opção para Elasticsearch 7.0 ou superior.

batchSize

Número máximo de documentos a recuperar do cluster Elasticsearch por requisição scroll.

Int

Não

2000

Nenhum

keepScrollAliveSecs

Tempo máximo para manter o contexto scroll ativo.

Int

Não

3600

Unidade: segundos.

Tabela de destino

Parâmetro

Descrição

Tipo

Obrigatório

Padrão

Observações

connector

Tipo da tabela de destino.

String

Sim

Nenhum

O valor deve ser elasticsearch-6, elasticsearch-7 ou elasticsearch-8.

Nota

Apenas VVR 8.0.5 ou posterior suporta o valor elasticsearch-8.

hosts

Endereço do servidor do cluster Elasticsearch.

String

Sim

Nenhum

Exemplo: 127.0.0.1:XXXX.

index

Nome do índice.

String

Sim

Nenhum

A tabela de destino suporta índices estáticos e dinâmicos:

  • Para índices estáticos, o valor deve ser uma string simples, como myusers. Todos os registros serão gravados no índice myusers.

  • Em índices dinâmicos, use {field_name} para referenciar valores de campos do registro e gerar o índice de destino dinamicamente. Também é possível usar {field_name|date_format_string} para converter valores de campos dos tipos TIMESTAMP, DATE e TIME para o formato especificado por date_format_string. O date_format_string é compatível com o DateTimeFormatter do Java. Por exemplo, se você definir o índice como myusers-{log_ts|yyyy-MM-dd}, um registro com o valor de campo log_ts igual a 2020-03-27 12:25:55 será gravado no índice myusers-2020-03-27.

document-type

Tipo do documento.

String

  • elasticsearch-6: Sim

  • elasticsearch-7: Não suportado

Nenhum

Quando o tipo de conector for elasticsearch-6, o valor deste parâmetro deve ser consistente com o valor do parâmetro type no Elasticsearch.

username

Nome de usuário para autenticação.

String

Não

Nenhum

Por padrão, a autenticação está desativada. Se você especificar username, também deverá especificar uma password não vazia.

Importante

Para evitar a exposição do seu nome de usuário e senha, recomendamos o uso de variáveis de projeto. Para mais informações, consulte variáveis de projeto.

password

Senha para autenticação.

String

Não

Nenhum

document-id.key-delimiter

Delimitador para o ID do documento.

String

Não

_

O conector usa a chave primária para gerar o ID do documento. Ele concatena todos os campos da chave primária na ordem definida no DDL, usando o delimitador especificado por document-id.key-delimiter, para criar uma string de ID de documento para cada linha.

Nota

Um ID de documento é uma string de até 512 bytes que não contém espaços.

failure-handler

Política de tratamento de falhas para requisições Elasticsearch malsucedidas.

String

Não

fail

Políticas válidas:

  • fail (padrão): Falha o job se uma requisição falhar.

  • ignore: Ignora a falha e descarta a requisição.

  • retry-rejected: Readiciona requisições que falharam devido a fila cheia.

  • Nome de classe personalizada: Usa uma subclasse ActionRequestFailureHandler para tratar falhas.

sink.flush-on-checkpoint

Define se deve haver flush durante o checkpoint.

Boolean

Não

true

  • true: Valor padrão.

  • false: Se desativado, o conector não aguarda o reconhecimento de todas as requisições pendentes durante um checkpoint. Isso significa que o conector não garante entrega at-least-once.

sink.bulk-flush.backoff.strategy

Se a operação de flush falhar devido a um erro temporário de requisição, defina sink.bulk-flush.backoff.strategy para especificar a estratégia de nova tentativa.

Enum

Não

DISABLED

  • DISABLED (padrão): Nenhuma nova tentativa é realizada. O job falha no primeiro erro de requisição.

  • CONSTANT: Estratégia de backoff constante em que o tempo de espera entre tentativas é sempre o mesmo.

  • EXPONENTIAL: Estratégia de backoff exponencial em que o tempo de espera entre tentativas aumenta exponencialmente.

sink.bulk-flush.backoff.max-retries

Número máximo de novas tentativas.

Int

Não

Nenhum

Nenhum

sink.bulk-flush.backoff.delay

Atraso entre as tentativas.

Duration

Não

Nenhum

  • Para a estratégia constant backoff, este valor representa o atraso entre cada tentativa.

  • Na estratégia exponential backoff, este valor corresponde ao atraso base inicial.

sink.bulk-flush.max-actions

Número máximo de ações armazenadas em buffer para cada requisição bulk.

Int

Não

1000

Um valor igual a 0 desativa este recurso.

sink.bulk-flush.max-size

Tamanho máximo de memória do buffer de requisições.

String

Não

2 MB

A unidade é MB. O valor padrão é 2 MB. Um valor igual a 0 desativa este recurso.

sink.bulk-flush.interval

Intervalo de flush.

Duration

Não

1s

A unidade é segundos. O valor padrão é 1s. Um valor de 0s desativa este recurso.

connection.path-prefix

String a prefixar em todo caminho de comunicação REST.

String

Não

Nenhum

Nenhum

retry-on-conflict

Número máximo de tentativas para uma operação de atualização em caso de conflito de versão. Se o número de tentativas exceder esse valor, o job falhará com uma exceção.

Int

Não

0

Nota
  • Esta opção é suportada apenas no VVR 4.0.13 ou posterior.

  • Esta opção só tem efeito quando uma chave primária é definida.

routing-fields

Especifica um ou mais nomes de campos do Elasticsearch usados para rotear um documento para um shard específico.

String

Não

Nenhum

Separe vários nomes de campos com ponto e vírgula (;). Se os dados de um campo estiverem vazios, o campo será definido como null.

Nota

Esta opção é suportada apenas no VVR 8.0.6 ou posterior, para elasticsearch-7 e elasticsearch-8.

sink.delete-strategy

Configura como o destino lida com uma mensagem de retração (-D para DELETE ou -U para UPDATE_BEFORE).

Enum

Não

DELETE_ROW_ON_PK

Estratégias válidas:

  • DELETE_ROW_ON_PK (padrão): Ignora mensagens -U, mas exclui a linha (documento) correspondente à chave primária ao receber uma mensagem -D.

  • IGNORE_DELETE: Ignora mensagens -U e -D. Nenhuma retração ocorre no destino Elasticsearch.

  • NON_PK_FIELD_TO_NULL: Ignora mensagens -U. No entanto, ao receber uma mensagem -D, essa configuração modifica a linha (documento) da chave primária: o valor da chave primária permanece o mesmo, e todos os outros valores de campos não primários no esquema da tabela são definidos como NULL. Isso é usado principalmente para atualizações parciais quando vários destinos gravam na mesma tabela Elasticsearch simultaneamente.

  • CHANGELOG_STANDARD: Semelhante a DELETE_ROW_ON_PK, mas também exclui a linha (documento) correspondente à chave primária quando uma mensagem -U é recebida.

    Nota

    Esta opção é suportada apenas no VVR 8.0.8 ou posterior.

sink.ignore-null-when-update

Ao atualizar dados, especifica se um campo deve ser atualizado para null ou mantido inalterado caso o valor recebido seja null.

BOOLEAN

Não

false

Valores válidos:

  • true: O campo não é atualizado. Este valor é suportado apenas quando uma chave primária é definida para a tabela Flink e o formato de dados do Elasticsearch é JSON.

  • false: O campo é atualizado para null.

Nota

Esta opção é suportada apenas no VVR 11.1 ou posterior.

connection.request-timeout

Tempo limite para solicitar uma conexão do gerenciador de conexões.

Duration

Não

Nenhum

Exemplo:

'connection.request-timeout' = '1 min' -- 1 minute
'connection.request-timeout' = '500ms' -- 500 milliseconds
Nota

Esta opção é suportada apenas no VVR 11.7 ou posterior.

connect.timeout

Tempo limite para estabelecer uma conexão.

Duration

Não

Nenhum

Exemplo:

'connect.timeout' = '1 min' -- 1 minute
'connect.timeout' = '500ms' -- 500 milliseconds
Nota

Esta opção é suportada apenas no VVR 11.7 ou posterior.

socket.timeout

Tempo limite de espera por dados, correspondendo ao período máximo de inatividade entre dois pacotes consecutivos.

Duration

Não

Nenhum

Exemplo:

'socket.timeout' = '1 min' -- 1 minute
'socket.timeout' = '500ms' -- 500 milliseconds
Nota

Esta opção é suportada apenas no VVR 11.7 ou posterior.

connection.keep-alive

Duração máxima que uma conexão pode permanecer ociosa antes de ser fechada pelo sistema. Se esta opção não for definida, o cabeçalho de resposta Keep-Alive do servidor determinará a duração. Caso o servidor não envie um cabeçalho Keep-Alive, a conexão permanecerá ativa indefinidamente.

Duration

Não

Nenhum

Exemplo:

'connection.keep-alive' = '1 min' -- 1 minute
'connection.keep-alive' = '500ms' -- 500 milliseconds
Nota

Esta opção é suportada apenas no VVR 11.7 ou posterior.

sink.bulk-flush.update.doc_as_upsert

Define se o documento deve ser tratado como um documento upsert em uma requisição de atualização.

BOOLEAN

Não

false

Valores válidos:

  • true: Define o campo doc_as_upsert da requisição de atualização como true.

  • false: Preenche o campo upsert da requisição de atualização com o documento.

De acordo com https://github.com/elastic/elasticsearch/issues/105804, os pipelines de ingestão do Elasticsearch não suportam atualizações parciais para requisições de atualização em massa. Se você deseja usar um pipeline de ingestão, defina esta opção como true.

Nota

Esta opção é suportada apenas no VVR 11.5 ou posterior.

Tabela de dimensão

Parâmetro

Descrição

Tipo

Obrigatório

Padrão

Observações

connector

Tipo da tabela de dimensão.

String

Sim

Nenhum

Valores válidos: elasticsearch ou elasticsearch-8.

Nota

Apenas VVR 11.6 ou posterior suporta o valor elasticsearch-8.

endPoint

Endereço do servidor do cluster Elasticsearch.

String

Sim

Nenhum

Nome legado da opção.

hosts

Para uso com elasticsearch-8.

indexName

Nome do índice.

String

Sim

Nenhum

Nome legado da opção.

index

Para uso com elasticsearch-8.

accessId

Nome de usuário para autenticação.

String

Não

Nenhum

Por padrão, este parâmetro está vazio e a autenticação não ocorre. Ao especificar accessId, informe um accessKey não vazio.

Nota

O elasticsearch-8 agora usa os parâmetros username e password.

Importante

Para evitar a exposição do seu nome de usuário e senha, recomendamos o uso de variáveis de projeto. Para mais informações, consulte variáveis de projeto.

username

accessKey

Senha para autenticação.

String

Não

Nenhum

password

typeNames

Nome do tipo.

String

Não

_doc

Recomenda-se não configurar esta opção para Elasticsearch 7.0 ou superior.

maxJoinRows

Número máximo de linhas para unir em uma única consulta.

Integer

Não

1024

Nenhum

cache

Estratégia de cache.

String

Não

Nenhum

Valores válidos:

  • ALL: Armazena todos os dados da tabela de dimensão em cache. Antes do início do job, o sistema carrega todos os dados da tabela de dimensão para o cache. Consultas subsequentes são atendidas diretamente do cache. Se uma chave não for encontrada, o sistema considera que ela não existe. O sistema recarrega todo o cache quando o TTL expira.

  • LRU: Armazena uma parte dos dados da tabela de dimensão em cache. Quando um registro da tabela de origem chega, o sistema primeiro busca os dados no cache. Em caso de cache miss, ele consulta a tabela de dimensão física.

  • None: Sem cache.

cacheSize

Tamanho do cache, especificado como número de linhas.

Long

Não

100000

O parâmetro cacheSize só tem efeito quando a política de cache LRU é selecionada para cache.

cacheTTLMs

Tempo de vida (TTL) do cache.

Long

Não

Long.MAX_VALUE

Unidade: milissegundos. O comportamento de cacheTTLMs depende da configuração de cache:

  • Quando cache é LRU, cacheTTLMs é o TTL das entradas de cache. Por padrão, as entradas não expiram.

  • Quando cache é ALL, cacheTTLMs é o intervalo para recarregar o cache. Por padrão, o cache não é recarregado.

ignoreKeywordSuffix

Define se o sufixo .keyword adicionado automaticamente aos campos STRING deve ser ignorado.

Boolean

Não

false

Para compatibilidade, o Flink converte tipos Text do Elasticsearch para STRING e anexa um sufixo .keyword ao nome do campo por padrão.

Valores válidos:

  • true: Ignora o sufixo.

    Se o sufixo impedir a correspondência com campos do tipo Text no Elasticsearch, defina esta opção como true.

  • false: Não ignora o sufixo.

cacheEmpty

Define se resultados vazios de consultas na tabela de dimensão física devem ser armazenados em cache.

Boolean

Não

true

O parâmetro cacheEmpty só é efetivo quando cache usa a política de cache LRU.

queryMaxDocs

Para tabelas de dimensão sem chave primária, este é o número máximo de documentos que o servidor Elasticsearch retorna para cada consulta.

Integer

Não

10000

O valor padrão de 10.000 corresponde ao número máximo de documentos que um servidor Elasticsearch pode retornar por consulta. Este valor não pode exceder esse limite.

Nota
  • Esta opção é suportada apenas no VVR 8.0.8 ou posterior.

  • Esta opção só tem efeito para tabelas de dimensão sem chave primária, já que os dados em tabelas com chave primária são únicos.

  • Um valor padrão alto ajuda a garantir a correção da consulta, mas aumenta o uso de memória durante as consultas ao Elasticsearch. Se você encontrar problemas de memória, reduza este valor para otimizar o consumo.

Mapeamento de tipos

O Flink analisa dados do Elasticsearch como JSON. Para detalhes, consulte mapeamento de tipos de dados.

Exemplos

  • Exemplo de tabela de origem

    CREATE TEMPORARY TABLE elasticsearch_source (
      name STRING,
      location STRING,
      `value` FLOAT
    ) WITH (
      'connector' ='elasticsearch',
      'endPoint' = '<yourEndPoint>',
      'accessId' = '${secret_values.ak_id}',
      'accessKey' = '${secret_values.ak_secret}',
      'indexName' = '<yourIndexName>',
      'typeNames' = '<yourTypeName>'
    );
    
    CREATE TEMPORARY TABLE blackhole_sink (
      name STRING,
      location STRING,
      `value` FLOAT
    ) WITH (
      'connector' ='blackhole'
    );
    
    INSERT INTO blackhole_sink
    SELECT name, location, `value`
    FROM elasticsearch_source;
  • Exemplo de tabela de dimensão

    CREATE TEMPORARY TABLE datagen_source (
      id STRING, 
      data STRING,
      proctime as PROCTIME()
    ) WITH (
      'connector' = 'datagen' 
    );
    
    CREATE TEMPORARY TABLE es_dim (
      id STRING,
      `value` FLOAT,
      PRIMARY KEY (id) NOT ENFORCED
    ) WITH (
      'connector' ='elasticsearch',
      'endPoint' = '<yourEndPoint>',
      'accessId' = '${secret_values.ak_id}',
      'accessKey' = '${secret_values.ak_secret}',
      'indexName' = '<yourIndexName>',
      'typeNames' = '<yourTypeName>'
    );
    
    CREATE TEMPORARY TABLE blackhole_sink (
      id STRING,
      data STRING,
      `value` FLOAT
    ) WITH (
      'connector' = 'blackhole' 
    );
    
    INSERT INTO blackhole_sink
    SELECT e.*, w.*
    FROM datagen_source AS e
    JOIN es_dim FOR SYSTEM_TIME AS OF e.proctime AS w
    ON e.id = w.id;
  • Exemplo de tabela de destino 1

    Este exemplo grava conteúdo de texto no Elasticsearch após a vetorização de texto.

    Nota

    Crie previamente um mapeamento de índice no Elasticsearch. Defina o tipo de dados do campo embedding como dense_vector e especifique as dimensões. Caso contrário, o Elasticsearch pode inferi-lo como um tipo de array comum.

        CREATE TEMPORARY TABLE datagen_source (
          id STRING,
          content STRING,
          embedding ARRAY<FLOAT>
        ) WITH (
          'connector' = 'datagen'
        );
    
        CREATE TEMPORARY TABLE es_sink (
          id STRING,
          content STRING,
          embedding ARRAY<FLOAT>,
          PRIMARY KEY (id) NOT ENFORCED -- The primary key is optional. If you define a primary key, its value becomes the document ID. Otherwise, a random document ID is generated.
        ) WITH (
          'connector' = 'elasticsearch-8',
          'hosts' = '<yourHosts>',
          'index' = '<yourIndex>',
          'username' ='${secret_values.ak_id}',
          'password' ='${secret_values.ak_secret}'
        );
    
        INSERT INTO es_sink
        SELECT id, content, embedding
        FROM datagen_source;
  • Exemplo de tabela de destino 2

    O conector suporta a gravação de tipos complexos como ROW, ARRAY e MAP no Elasticsearch.

    CREATE TEMPORARY TABLE datagen_source(  
      id STRING,
        details ROW<  
            name STRING,  
            ages ARRAY<INT>,  
            attributes MAP<STRING, STRING>  
        >
    ) WITH (  
        'connector' = 'datagen'
    );
    
    CREATE TEMPORARY TABLE es_sink (
      id STRING,
        details ROW<  
            name STRING,  
            ages ARRAY<INT>,  
            attributes MAP<STRING, STRING>  
        >, 
      PRIMARY KEY (id) NOT ENFORCED  -- The primary key is optional. If you define a primary key, its value becomes the document ID. Otherwise, a random document ID is generated.
    ) WITH (
      'connector' = 'elasticsearch-6',
      'hosts' = '<yourHosts>',
      'index' = '<yourIndex>',
      'document-type' = '<yourElasticsearch.types>',
      'username' ='${secret_values.ak_id}',
      'password' ='${secret_values.ak_secret}'
    );
    
    INSERT INTO es_sink
    SELECT id, details
    FROM datagen_source;