Cette rubrique explique comment utiliser le connecteur Elasticsearch.
Contexte
Alibaba Cloud Elasticsearch est compatible avec Elasticsearch open source et inclut des fonctionnalités commerciales telles que Security, Machine Learning, Graph et APM pour l'analyse et la recherche de données. Il fournit des services de niveau entreprise, notamment le contrôle d'accès, la surveillance de sécurité et les alertes, ainsi que la génération automatique de rapports.
Le tableau suivant décrit les capacités du connecteur Elasticsearch.
|
Élément |
Description |
|
Type de table |
Table source, table de dimension et table de destination (sink) |
|
Mode d'exécution |
Mode par lots et mode en continu |
|
Format de données |
JSON |
|
Métrique |
|
|
Type d'API |
API DataStream et SQL |
|
Mise à jour ou suppression de données dans une table de destination |
Pris en charge |
Prérequis
Vous avez créé un index Elasticsearch. Pour plus d'informations, consultez Prise en main.
Vous avez configuré une liste d'autorisation d'adresses IP publiques ou privées pour l'instance Elasticsearch. Pour plus d'informations, consultez Gérer les listes d'autorisation d'adresses IP.
Limites
-
Les tables sources et les tables de dimension prennent en charge Elasticsearch 6.8.x ou version ultérieure.
RemarqueL'utilisation d'Elasticsearch 8.x avec des tables sources et de dimension nécessite VVR 11.6 ou une version ultérieure.
Les tables de destination ne prennent en charge que les versions 6.x, 7.x et 8.x d'Elasticsearch.
Seules les tables sources Elasticsearch complètes sont prises en charge ; les tables incrémentielles ne le sont pas.
Syntaxe
-
Table source
Elasticsearch 8.x
CREATE TABLE elasticsearch_source( name STRING, location STRING, value FLOAT ) WITH ( 'connector' ='elasticsearch-8', 'hosts' = '<yourHosts>', 'index' = '<yourIndex>' );Autres versions
CREATE TABLE elasticsearch_source( name STRING, location STRING, value FLOAT ) WITH ( 'connector' ='elasticsearch', 'endPoint' = '<yourEndPoint>', 'indexName' = '<yourIndexName>' ); -
Table de dimension
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>' );Autres versions
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>' );RemarqueSi une clé primaire est spécifiée, un seul champ peut servir de clé de jointure et il doit correspondre à l'ID de document dans l'index Elasticsearch associé.
Si aucune clé primaire n'est spécifiée, vous pouvez utiliser un ou plusieurs champs comme clés de jointure. Ces clés doivent être des champs présents dans les documents Elasticsearch correspondants.
Pour les champs
STRING, le connecteur ajoute par défaut le suffixe.keywordaux noms de champs pour des raisons de compatibilité. Si cela empêche la correspondance avec les champsTEXTdans Elasticsearch, définissez l'optionignoreKeywordSuffixsurtrue.
-
Table de destination
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>' );Remarque-
Une table de destination Elasticsearch fonctionne soit en
upsert mode, soit enappend mode, selon qu'une clé primaire est définie.Si une clé primaire est définie, sa valeur est utilisée comme ID de document. La table de destination fonctionne alors en
upsert modeet peut traiter les opérationsUPDATEetDELETE.Si aucune clé primaire n'est définie, Elasticsearch génère automatiquement un ID de document aléatoire. La table de destination fonctionne alors en
append modeet ne peut consommer que des messagesINSERT.
Les types de données tels que
BYTES,ROW,ARRAYetMAPn'ont pas de représentation sous forme de chaîne correspondante. Par conséquent, vous ne pouvez pas utiliser des champs de ces types de données comme clé primaire.Les champs de la
DDLcorrespondent aux champs d'un document Elasticsearch. Vous ne pouvez pas écrire de métadonnées, telles que l'ID de document, dans la table de destination car le cluster Elasticsearch conserve ces métadonnées.
-
Options WITH
Table source
|
Paramètre |
Description |
Type |
Obligatoire |
Valeur par défaut |
Remarques |
|
connector |
Type de la table source. |
String |
Oui |
Aucune |
Valeurs valides : Remarque
Seule la version VVR 11.6 ou ultérieure prend en charge la valeur |
|
endPoint |
Adresse du serveur du cluster Elasticsearch. |
String |
Oui |
Aucune |
Nom d'option hérité. |
|
hosts |
À utiliser avec |
||||
|
indexName |
Nom de l'index. |
String |
Oui |
Aucune |
Nom d'option hérité. |
|
index |
Pour une utilisation avec |
||||
|
accessId |
Nom d'utilisateur pour l'authentification. |
String |
Non |
Aucune |
Par défaut, ce paramètre est vide et aucune authentification n'est effectuée. Si vous spécifiez accessId, vous devez indiquer une valeur non vide pour accessKey. Remarque
La version Important
Pour éviter toute exposition de votre nom d'utilisateur et de votre mot de passe, nous vous recommandons d'utiliser des variables de projet. Pour plus d'informations, consultez la rubrique variables de projet. |
|
username |
|||||
|
accessKey |
Mot de passe pour l'authentification. |
String |
Non |
Aucune |
|
|
password |
|||||
|
typeNames |
Nom du type. |
String |
Non |
_doc |
Nous vous déconseillons de configurer cette option pour Elasticsearch 7.0 ou versions ultérieures. |
|
batchSize |
Nombre maximal de documents à récupérer depuis le cluster Elasticsearch par requête de défilement (scroll). |
Int |
Non |
2000 |
Aucune |
|
keepScrollAliveSecs |
Durée maximale de conservation du contexte de défilement (scroll) actif. |
Int |
Non |
3600 |
Unité : secondes. |
Table de destination
|
Paramètre |
Description |
Type |
Obligatoire |
Valeur par défaut |
Remarques |
|
connector |
Le type de la table de destination (sink). |
String |
Oui |
Aucune |
La valeur doit être Remarque
Seule la version VVR 8.0.5 ou ultérieure prend en charge la valeur |
|
hosts |
L'adresse du serveur du cluster Elasticsearch. |
String |
Oui |
Aucune |
Exemple : |
|
index |
Le nom de l'index. |
String |
Oui |
Aucune |
La table de destination prend en charge les index statiques et dynamiques :
|
|
document-type |
Le type de document. |
String |
|
Aucune |
Lorsque le type de connecteur est |
|
username |
Le nom d'utilisateur pour l'authentification. |
String |
Non |
Aucune |
Par défaut, l'authentification est désactivée. Si vous spécifiez Important
Pour éviter l'exposition de votre nom d'utilisateur et de votre mot de passe, nous vous recommandons d'utiliser des variables de projet. Pour plus d'informations, consultez la section variables de projet. |
|
password |
Le mot de passe pour l'authentification. |
String |
Non |
Aucune |
|
|
document-id.key-delimiter |
Le délimiteur pour l'ID de document. |
String |
Non |
_ |
Le connecteur utilise la clé primaire pour générer l'ID de document. Il concatène tous les champs de la clé primaire dans l'ordre défini dans le DDL, en utilisant le délimiteur spécifié par document-id.key-delimiter, afin de créer une chaîne d'ID de document pour chaque ligne. Remarque
Un ID de document est une chaîne pouvant contenir jusqu'à 512 octets et ne comportant pas d'espaces. |
|
failure-handler |
La politique de gestion des échecs pour les requêtes Elasticsearch ayant échoué. |
String |
Non |
fail |
Politiques valides :
|
|
sink.flush-on-checkpoint |
Indique s'il faut effectuer un vidage (flush) lors du checkpoint. |
Boolean |
Non |
true |
|
|
sink.bulk-flush.backoff.strategy |
Si l'opération de vidage échoue en raison d'une erreur de requête temporaire, définissez sink.bulk-flush.backoff.strategy pour spécifier la stratégie de nouvelle tentative. |
Enum |
Non |
DISABLED |
|
|
sink.bulk-flush.backoff.max-retries |
Le nombre maximal de nouvelles tentatives. |
Int |
Non |
Aucune |
Aucune |
|
sink.bulk-flush.backoff.delay |
Le délai entre les nouvelles tentatives. |
Duration |
Non |
Aucune |
|
|
sink.bulk-flush.max-actions |
Le nombre maximal d'actions mises en mémoire tampon pour chaque requête bulk. |
Int |
Non |
1000 |
Une valeur de 0 désactive cette fonctionnalité. |
|
sink.bulk-flush.max-size |
La taille mémoire maximale du tampon de requêtes. |
String |
Non |
2 MB |
L'unité est le Mo. La valeur par défaut est 2 Mo. Une valeur de 0 désactive cette fonctionnalité. |
|
sink.bulk-flush.interval |
L'intervalle de vidage (flush). |
Duration |
Non |
1s |
L'unité est la seconde. La valeur par défaut est 1 s. Une valeur de 0 s désactive cette fonctionnalité. |
|
connection.path-prefix |
Une chaîne à ajouter au début de chaque chemin de communication REST. |
String |
Non |
Aucune |
Aucune |
|
retry-on-conflict |
Le nombre maximal de nouvelles tentatives pour une opération de mise à jour en cas de conflit de version. Si le nombre de nouvelles tentatives dépasse cette valeur, le job échoue avec une exception. |
Int |
Non |
0 |
Remarque
|
|
routing-fields |
Spécifie un ou plusieurs noms de champs Elasticsearch utilisés pour router un document vers un shard spécifique. |
String |
Non |
Aucune |
Séparez plusieurs noms de champs par un point-virgule (;). Si les données d'un champ sont vides, le champ est défini sur null. Remarque
Cette option est prise en charge uniquement dans la version VVR 8.0.6 ou ultérieure, pour |
|
sink.delete-strategy |
Configure la façon dont le sink gère un message de rétractation (-D pour DELETE ou -U pour UPDATE_BEFORE). |
Enum |
Non |
DELETE_ROW_ON_PK |
Stratégies valides :
|
|
sink.ignore-null-when-update |
Lors de la mise à jour des données, indique s'il faut mettre à jour un champ avec la valeur |
BOOLEAN |
Non |
false |
Valeurs valides :
Remarque
Cette option est prise en charge uniquement dans la version VVR 11.1 ou ultérieure. |
|
connection.request-timeout |
Le délai d'expiration pour demander une connexion au gestionnaire de connexions. |
Duration |
Non |
Aucune |
Exemple :
Remarque
Cette option est prise en charge uniquement dans la version VVR 11.7 ou ultérieure. |
|
connect.timeout |
Le délai d'expiration pour l'établissement d'une connexion. |
Duration |
Non |
Aucune |
Exemple :
Remarque
Cette option est prise en charge uniquement dans la version VVR 11.7 ou ultérieure. |
|
socket.timeout |
Délai d'attente pour la réception des données, correspondant à la période maximale d'inactivité entre deux paquets de données consécutifs. |
Durée |
Non |
Aucune |
Exemple :
Remarque
Cette option est prise en charge uniquement dans les versions VVR 11.7 et ultérieures. |
|
connection.keep-alive |
Durée maximale pendant laquelle une connexion peut rester inactive avant que le système ne la ferme. Si cette option n'est pas définie, la durée est déterminée par l'en-tête de réponse |
Durée |
Non |
Aucune |
Exemple :
Remarque
Cette option est prise en charge uniquement dans les versions VVR 11.7 et ultérieures. |
|
sink.bulk-flush.update.doc_as_upsert |
Indique si le document doit être traité comme un document upsert dans une requête de mise à jour. |
BOOLEAN |
Non |
false |
Valeurs valides :
Selon https://github.com/elastic/elasticsearch/issues/105804, les pipelines d'ingestion Elasticsearch ne prennent pas en charge les mises à jour partielles pour les requêtes de mise à jour en bloc. Si vous souhaitez utiliser un pipeline d'ingestion, définissez cette option sur true. Remarque
Cette option est prise en charge uniquement dans les versions VVR 11.5 et ultérieures. |
Table de dimension
|
Paramètre |
Description |
Type |
Obligatoire |
Valeur par défaut |
Remarques |
|
connector |
Le type de la table de dimension. |
String |
Oui |
Aucune |
Valeurs valides : Remarque
Seule la version VVR 11.6 ou ultérieure prend en charge la valeur |
|
endPoint |
L'adresse du serveur du cluster Elasticsearch. |
String |
Oui |
Aucune |
Nom d'option hérité. |
|
hosts |
À utiliser avec |
||||
|
indexName |
Le nom de l'index. |
String |
Oui |
Aucune |
Nom d'option hérité. |
|
index |
À utiliser avec |
||||
|
accessId |
Le nom d'utilisateur pour l'authentification. |
String |
Non |
Aucune |
Par défaut, ce paramètre est vide et aucune authentification n'est effectuée. Si vous spécifiez accessId, vous devez spécifier une valeur non vide pour accessKey. Remarque
Important
Pour éviter l'exposition de votre nom d'utilisateur et de votre mot de passe, nous vous recommandons d'utiliser des variables de projet. Pour plus d'informations, consultez les variables de projet. |
|
username |
|||||
|
accessKey |
Le mot de passe pour l'authentification. |
String |
Non |
Aucune |
|
|
password |
|||||
|
typeNames |
Le nom du type. |
String |
Non |
_doc |
Nous vous recommandons de ne pas configurer cette option pour Elasticsearch 7.0 ou version ultérieure. |
|
maxJoinRows |
Le nombre maximal de lignes à joindre pour une seule recherche. |
Integer |
Non |
1024 |
Aucune |
|
cache |
La stratégie de mise en cache. |
String |
Non |
Aucune |
Valeurs valides :
|
|
cacheSize |
La taille du cache, spécifiée sous forme de nombre de lignes. |
Long |
Non |
100000 |
Le paramètre cacheSize prend effet uniquement lorsque la stratégie de cache LRU est sélectionnée pour cache. |
|
cacheTTLMs |
Le délai d'expiration (TTL) du cache. |
Long |
Non |
Long.MAX_VALUE |
Unité : millisecondes. Le comportement de cacheTTLMs dépend du paramètre cache :
|
|
ignoreKeywordSuffix |
Indique s'il faut ignorer le suffixe .keyword qui est automatiquement ajouté aux champs STRING. |
Boolean |
Non |
false |
Pour des raisons de compatibilité, Flink convertit les types Valeurs valides :
|
|
cacheEmpty |
Indique s'il faut mettre en cache les résultats vides provenant des recherches dans la table de dimension physique. |
Boolean |
Non |
true |
Le paramètre cacheEmpty est effectif uniquement lorsque cache utilise la stratégie de cache LRU. |
|
queryMaxDocs |
Pour les tables de dimension sans clé primaire, il s'agit du nombre maximal de documents que le serveur Elasticsearch renvoie pour chaque requête de recherche. |
Integer |
Non |
10000 |
La valeur par défaut de 10 000 correspond au nombre maximal de documents qu'un serveur Elasticsearch peut renvoyer par requête. Cette valeur ne peut pas dépasser cette limite. Remarque
|
Mappage des types
Flink analyse les données Elasticsearch au format JSON. Pour plus de détails, consultez le mappage des types de données.
Exemples
-
Exemple de table source
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; -
Exemple de table de dimension
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; -
Exemple de table de destination 1
Cet exemple écrit du contenu texte dans Elasticsearch après la vectorisation de texte.
RemarqueCréez un mappage d'index dans Elasticsearch au préalable. Définissez le type de données du champ
embeddingsurdense_vectoret spécifiez les dimensions. Sinon, Elasticsearch pourrait l'interpréter comme un type de tableau standard.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; -
Exemple de table de destination 2
Le connecteur prend en charge l'écriture de types complexes tels que
ROW,ARRAYetMAPdans 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;