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 |
|
|
Tipo de API |
DataStream API e SQL |
|
Atualização ou exclusão de dados em uma tabela de destino |
Suportado |
Pré-requisitos
Crie um índice Elasticsearch. Para mais detalhes, consulte Primeiros passos.
Configure uma lista de permissões de endereços IP públicos ou privados para a instância Elasticsearch. Para mais informações, consulte Gerenciar listas de permissões de endereços IP.
Limitações
-
As tabelas de origem e de dimensão suportam Elasticsearch 6.8.x ou superior.
NotaO 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>' );NotaSe 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.keywordaos nomes dos campos por padrão para garantir compatibilidade. Se isso impedir a correspondência com camposTEXTno Elasticsearch, defina a opçãoignoreKeywordSuffixcomotrue.
-
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 modeouappend 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çõesUPDATEeDELETE.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 mensagensINSERT.
Tipos de dados como
BYTES,ROW,ARRAYeMAPnão possuem representação em string correspondente. Portanto, não é possível usar campos desses tipos como chave primária.Os campos no
DDLcorrespondem 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: Nota
Apenas VVR 11.6 ou posterior suporta o valor |
|
endPoint |
Endereço do servidor do cluster Elasticsearch. |
String |
Sim |
Nenhum |
Nome legado da opção. |
|
hosts |
Usado com |
||||
|
indexName |
Nome do índice. |
String |
Sim |
Nenhum |
Nome legado da opção. |
|
index |
Para uso com |
||||
|
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 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 Nota
Apenas VVR 8.0.5 ou posterior suporta o valor |
|
hosts |
Endereço do servidor do cluster Elasticsearch. |
String |
Sim |
Nenhum |
Exemplo: |
|
index |
Nome do índice. |
String |
Sim |
Nenhum |
A tabela de destino suporta índices estáticos e dinâmicos:
|
|
document-type |
Tipo do documento. |
String |
|
Nenhum |
Quando o tipo de conector for |
|
username |
Nome de usuário para autenticação. |
String |
Não |
Nenhum |
Por padrão, a autenticação está desativada. Se você especificar 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:
|
|
sink.flush-on-checkpoint |
Define se deve haver flush durante o checkpoint. |
Boolean |
Não |
true |
|
|
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 |
|
|
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 |
|
|
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
|
|
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 |
|
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:
|
|
sink.ignore-null-when-update |
Ao atualizar dados, especifica se um campo deve ser atualizado para |
BOOLEAN |
Não |
false |
Valores válidos:
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:
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:
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:
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 |
Duration |
Não |
Nenhum |
Exemplo:
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:
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: Nota
Apenas VVR 11.6 ou posterior suporta o valor |
|
endPoint |
Endereço do servidor do cluster Elasticsearch. |
String |
Sim |
Nenhum |
Nome legado da opção. |
|
hosts |
Para uso com |
||||
|
indexName |
Nome do índice. |
String |
Sim |
Nenhum |
Nome legado da opção. |
|
index |
Para uso com |
||||
|
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 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:
|
|
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:
|
|
ignoreKeywordSuffix |
Define se o sufixo .keyword adicionado automaticamente aos campos STRING deve ser ignorado. |
Boolean |
Não |
false |
Para compatibilidade, o Flink converte tipos Valores válidos:
|
|
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
|
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.
NotaCrie previamente um mapeamento de índice no Elasticsearch. Defina o tipo de dados do campo
embeddingcomodense_vectore 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,ARRAYeMAPno 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;