Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Elasticsearch

Dernière mise à jour :Aug 09, 2026

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

Métrique

  • Table source

    • pendingRecords

    • numRecordsIn

    • numRecordsInPerSecond

    • numBytesIn

    • numBytesInPerSecond

  • Table de dimension

    Aucune

  • Table de destination (pour Ververica Runtime (VVR) 6.0.6 et versions ultérieures)

    • numRecordsOut

    • numRecordsOutPerSecond

Remarque

Pour plus d'informations sur ces métriques, consultez Métriques.

Type d'API

API DataStream et SQL

Mise à jour ou suppression de données dans une table de destination

Pris en charge

Prérequis

Limites

  • Les tables sources et les tables de dimension prennent en charge Elasticsearch 6.8.x ou version ultérieure.

    Remarque

    L'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>'
    );
    Remarque
    • Si 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 .keyword aux noms de champs pour des raisons de compatibilité. Si cela empêche la correspondance avec les champs TEXT dans Elasticsearch, définissez l'option ignoreKeywordSuffix sur true.

  • 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 en append 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 mode et peut traiter les opérations UPDATE et DELETE.

      • 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 mode et ne peut consommer que des messages INSERT.

    • Les types de données tels que BYTES, ROW, ARRAY et MAP n'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 DDL correspondent 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 : elasticsearch ou elasticsearch-8.

Remarque

Seule la version VVR 11.6 ou ultérieure prend en charge la valeur elasticsearch-8.

endPoint

Adresse du serveur du cluster Elasticsearch.

String

Oui

Aucune

Nom d'option hérité.

hosts

À utiliser avec elasticsearch-8.

indexName

Nom de l'index.

String

Oui

Aucune

Nom d'option hérité.

index

Pour une utilisation avec elasticsearch-8.

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 elasticsearch-8 a été mise à jour pour utiliser les paramètres username et password.

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 elasticsearch-6, elasticsearch-7 ou elasticsearch-8.

Remarque

Seule la version VVR 8.0.5 ou ultérieure prend en charge la valeur elasticsearch-8.

hosts

L'adresse du serveur du cluster Elasticsearch.

String

Oui

Aucune

Exemple : 127.0.0.1:XXXX.

index

Le nom de l'index.

String

Oui

Aucune

La table de destination prend en charge les index statiques et dynamiques :

  • Pour un index statique, la valeur doit être une chaîne simple, telle que myusers. Tous les enregistrements sont écrits dans l'index myusers.

  • Pour un index dynamique, vous pouvez utiliser {field_name} pour faire référence aux valeurs des champs de l'enregistrement afin de générer dynamiquement l'index cible. Vous pouvez également utiliser {field_name|date_format_string} pour convertir les valeurs des champs de type TIMESTAMP, DATE et TIME au format spécifié par date_format_string. Le paramètre date_format_string est compatible avec DateTimeFormatter de Java. Par exemple, si vous définissez l'index sur myusers-{log_ts|yyyy-MM-dd}, un enregistrement dont la valeur du champ log_ts est 2020-03-27 12:25:55 est écrit dans l'index myusers-2020-03-27.

document-type

Le type de document.

String

  • elasticsearch-6 : Oui

  • elasticsearch-7 : Non pris en charge

Aucune

Lorsque le type de connecteur est elasticsearch-6, la valeur de ce paramètre doit correspondre à la valeur du paramètre type dans Elasticsearch.

username

Le nom d'utilisateur pour l'authentification.

String

Non

Aucune

Par défaut, l'authentification est désactivée. Si vous spécifiez username, vous devez également indiquer un password non vide.

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 :

  • fail (par défaut) : Échec du job si une requête échoue.

  • ignore : Ignore l'échec et abandonne la requête.

  • retry-rejected : Réinsère les requêtes ayant échoué en raison d'une file d'attente pleine.

  • Nom de classe personnalisé : Utilise une sous-classe ActionRequestFailureHandler pour la gestion des échecs.

sink.flush-on-checkpoint

Indique s'il faut effectuer un vidage (flush) lors du checkpoint.

Boolean

Non

true

  • true : Valeur par défaut.

  • false : Si cette option est désactivée, le connecteur n'attend pas que toutes les requêtes en attente soient acquittées lors d'un checkpoint. Cela signifie que le connecteur ne fournit pas de garantie de livraison « au moins une fois » (at-least-once).

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

  • DISABLED (par défaut) : Aucune nouvelle tentative n'est effectuée. Le job échoue dès la première erreur de requête.

  • CONSTANT : Stratégie de backoff constante où le temps d'attente entre les nouvelles tentatives reste identique.

  • EXPONENTIAL : Stratégie de backoff exponentiel où le temps d'attente entre les nouvelles tentatives augmente de manière exponentielle.

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

  • Pour une stratégie de backoff constant, cette valeur représente le délai entre chaque nouvelle tentative.

  • Pour une stratégie de backoff exponentiel, cette valeur correspond au délai de base initial.

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
  • Cette option est prise en charge uniquement dans la version VVR 4.0.13 ou ultérieure.

  • Cette option n'est effective que si une clé primaire est définie.

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 elasticsearch-7 et elasticsearch-8.

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 :

  • DELETE_ROW_ON_PK (par défaut) : Ignore les messages -U mais supprime la ligne (document) correspondant à la clé primaire lors de la réception d'un message -D.

  • IGNORE_DELETE : Ignore à la fois les messages -U et -D. Aucune rétractation ne se produit dans le sink Elasticsearch.

  • NON_PK_FIELD_TO_NULL : Ignore les messages -U. Toutefois, lors de la réception d'un message -D, ce paramètre modifie la ligne (document) associée à la clé primaire : la valeur de la clé primaire reste inchangée et toutes les autres valeurs non clés primaires du schéma de table sont définies sur NULL. Cette option est principalement utilisée pour les mises à jour partielles lorsque plusieurs sinks écrivent simultanément dans la même table Elasticsearch.

  • CHANGELOG_STANDARD : Similaire à DELETE_ROW_ON_PK, mais supprime également la ligne (document) correspondant à la clé primaire lors de la réception d'un message -U.

    Remarque

    Cette option est prise en charge uniquement dans la version VVR 8.0.8 ou ultérieure.

sink.ignore-null-when-update

Lors de la mise à jour des données, indique s'il faut mettre à jour un champ avec la valeur null ou le laisser inchangé si la valeur du champ entrant est null.

BOOLEAN

Non

false

Valeurs valides :

  • true : Le champ n'est pas mis à jour. Cette valeur est prise en charge uniquement lorsqu'une clé primaire est définie pour la table Flink et que le format des données Elasticsearch est JSON.

  • false : Le champ est mis à jour avec la valeur null.

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 :

'connection.request-timeout' = '1 min' -- 1 minute
'connection.request-timeout' = '500ms' -- 500 milliseconds
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 :

'connect.timeout' = '1 min' -- 1 minute
'connect.timeout' = '500ms' -- 500 milliseconds
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 :

'socket.timeout' = '1 min' -- 1 minute
'socket.timeout' = '500ms' -- 500 milliseconds
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 Keep-Alive du serveur. Si le serveur n'envoie pas d'en-tête Keep-Alive, la connexion reste active indéfiniment.

Durée

Non

Aucune

Exemple :

'connection.keep-alive' = '1 min' -- 1 minute
'connection.keep-alive' = '500ms' -- 500 milliseconds
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 :

  • true : définit le champ doc_as_upsert de la requête de mise à jour sur true.

  • false : remplit le champ upsert de la requête de mise à jour avec le document.

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 : elasticsearch ou elasticsearch-8.

Remarque

Seule la version VVR 11.6 ou ultérieure prend en charge la valeur elasticsearch-8.

endPoint

L'adresse du serveur du cluster Elasticsearch.

String

Oui

Aucune

Nom d'option hérité.

hosts

À utiliser avec elasticsearch-8.

indexName

Le nom de l'index.

String

Oui

Aucune

Nom d'option hérité.

index

À utiliser avec elasticsearch-8.

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

elasticsearch-8 utilise désormais les paramètres username et password.

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 :

  • ALL : Met en cache toutes les données de la table de dimension. Avant le démarrage de la tâche, le système charge toutes les données de la table de dimension dans le cache. Les recherches suivantes sont servies directement depuis le cache. Si une clé est introuvable, le système la considère comme inexistante. Le système recharge l'intégralité du cache lorsque le délai TTL expire.

  • LRU : Met en cache une partie des données de la table de dimension. Lorsqu'un enregistrement de la table source arrive, le système recherche d'abord les données dans le cache. En cas d'échec du cache, il interroge la table de dimension physique.

  • None : Aucune mise en cache.

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 :

  • Lorsque cache est défini sur LRU, cacheTTLMs représente le TTL des entrées de cache. Par défaut, les entrées n'expirent pas.

  • Lorsque cache est défini sur ALL, cacheTTLMs représente l'intervalle de rechargement du cache. Par défaut, le cache n'est pas rechargé.

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 Text d'Elasticsearch en STRING et ajoute un suffixe .keyword au nom du champ par défaut.

Valeurs valides :

  • true : Ignore le suffixe.

    Si le suffixe empêche la correspondance avec les champs de type Text dans Elasticsearch, définissez cette option sur true.

  • false : N'ignore pas le suffixe.

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
  • Cette option est prise en charge uniquement dans VVR 8.0.8 ou version ultérieure.

  • Cette option prend effet uniquement pour les tables de dimension sans clé primaire, car les données des tables avec clé primaire sont uniques.

  • Une valeur par défaut élevée permet de garantir l'exactitude des requêtes, mais augmente l'utilisation de la mémoire lors des requêtes Elasticsearch. Si vous rencontrez des problèmes de mémoire, vous pouvez réduire cette valeur pour optimiser l'utilisation de la mémoire.

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.

    Remarque

    Créez un mappage d'index dans Elasticsearch au préalable. Définissez le type de données du champ embedding sur dense_vector et 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, ARRAY et MAP dans 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;