All Products
Search
Document Center

Realtime Compute for Apache Flink:Elasticsearch

Last Updated:May 23, 2026

Topik ini menjelaskan cara menggunakan konektor Elasticsearch.

Latar Belakang

Alibaba Cloud Elasticsearch kompatibel dengan Elasticsearch open source dan mencakup fitur komersial seperti Security, Machine Learning, Graph, dan APM untuk analisis serta pencarian data. Layanan ini menyediakan layanan tingkat enterprise, termasuk kontrol akses, pemantauan keamanan dan peringatan, serta pembuatan laporan otomatis.

Tabel berikut menjelaskan kemampuan konektor Elasticsearch.

Item

Deskripsi

Tipe tabel

Tabel sumber, tabel dimensi, dan tabel sink

Mode eksekusi

Mode batch dan mode streaming

Format data

JSON

Metrik

Metrik

  • Tabel sumber

    • pendingRecords

    • numRecordsIn

    • numRecordsInPerSecond

    • numBytesIn

    • numBytesInPerSecond

  • Tabel dimensi

    Tidak ada

  • Tabel sink (untuk Ververica Runtime (VVR) 6.0.6 dan versi lebih baru)

    • numRecordsOut

    • numRecordsOutPerSecond

Catatan

Untuk informasi lebih lanjut mengenai metrik ini, lihat Metrik.

Tipe API

DataStream API dan SQL

Pembaruan atau penghapusan data pada tabel sink

Didukung

Prasyarat

  • Anda telah membuat indeks Elasticsearch. Untuk informasi selengkapnya, lihat Memulai.

  • Anda telah mengonfigurasi daftar putih alamat IP publik atau privat untuk instans Elasticsearch. Untuk informasi selengkapnya, lihat Mengelola daftar putih alamat IP.

Batasan

  • Tabel sumber dan tabel dimensi mendukung Elasticsearch 6.8.x atau versi lebih baru.

    Catatan

    Penggunaan Elasticsearch 8.x dengan tabel sumber dan dimensi memerlukan VVR 11.6 atau versi lebih baru.

  • Tabel sink hanya mendukung Elasticsearch 6.x, 7.x, dan 8.x.

  • Hanya tabel sumber Elasticsearch lengkap yang didukung; tabel inkremental tidak didukung.

Sintaksis

  • Tabel sumber

    Elasticsearch 8.x

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

    Versi lainnya

    CREATE TABLE elasticsearch_source(
      name STRING,
      location STRING,
      value FLOAT
    ) WITH (
      'connector' ='elasticsearch',
      'endPoint' = '<yourEndPoint>',
      'indexName' = '<yourIndexName>'
    );
  • Tabel dimensi

    Elasticsearch 8.x

    CREATE TABLE es_dim(
      field1 STRING, -- Harus bertipe STRING jika digunakan sebagai kunci untuk JOIN.
      field2 FLOAT,
      field3 BIGINT,
      PRIMARY KEY (field1) NOT ENFORCED
    ) WITH (
      'connector' ='elasticsearch-8',
      'hosts' = '<yourHosts>',
      'index' = '<yourIndex>'
    );

    Versi lainnya

    CREATE TABLE es_dim(
      field1 STRING, -- Harus bertipe STRING jika digunakan sebagai kunci untuk JOIN.
      field2 FLOAT,
      field3 BIGINT,
      PRIMARY KEY (field1) NOT ENFORCED
    ) WITH (
      'connector' ='elasticsearch',
      'endPoint' = '<yourEndPoint>',
      'indexName' = '<yourIndexName>'
    );
    Catatan
    • Jika kunci primer ditentukan, hanya satu bidang yang dapat menjadi kunci gabungan, dan harus merupakan ID dokumen di indeks Elasticsearch yang sesuai.

    • Jika tidak ada kunci primer yang ditentukan, Anda dapat menggunakan satu atau beberapa bidang sebagai kunci gabungan. Kunci-kunci tersebut harus merupakan bidang dalam dokumen Elasticsearch yang sesuai.

    • Untuk bidang STRING, konektor secara default menambahkan akhiran .keyword ke nama bidang demi kompatibilitas. Jika hal ini mencegah pencocokan dengan bidang tipe TEXT di Elasticsearch, atur opsi ignoreKeywordSuffix ke true.

  • Sink table

    CREATE TABLE es_sink(
      user_id   STRING,
      user_name   STRING,
      uv BIGINT,
      pv BIGINT,
      PRIMARY KEY (user_id) NOT ENFORCED
    ) WITH (
      'connector' = 'elasticsearch-7', -- Jika Anda menggunakan Elasticsearch 6.x, atur nilai ini menjadi 'elasticsearch-6'.
      'hosts' = '<yourHosts>',
      'index' = '<yourIndex>'
    );
    Catatan
    • Tabel sink Elasticsearch beroperasi dalam mode upsert atau append, tergantung pada apakah kunci primer didefinisikan.

      • Jika kunci primer didefinisikan, nilainya digunakan sebagai ID dokumen. Tabel sink kemudian beroperasi dalam mode upsert dan dapat memproses operasi UPDATE dan DELETE.

      • Jika tidak ada kunci primer yang didefinisikan, Elasticsearch secara otomatis menghasilkan ID dokumen acak. Tabel sink kemudian beroperasi dalam mode append dan hanya dapat mengonsumsi pesan INSERT.

    • Tipe data seperti BYTES, ROW, ARRAY, dan MAP tidak memiliki representasi string yang sesuai. Oleh karena itu, bidang dengan tipe data tersebut tidak dapat digunakan sebagai kunci primer.

    • Bidang dalam DDL bersesuaian dengan bidang dalam dokumen Elasticsearch. Anda tidak dapat menulis metadata, seperti ID dokumen, ke tabel sink karena metadata tersebut dikelola oleh kluster Elasticsearch.

Opsi WITH

Tabel sumber

Parameter

Deskripsi

Tipe

Wajib

Bawaan

Keterangan

connector

Tipe tabel sumber.

String

Ya

Tidak ada

Nilai valid: elasticsearch atau elasticsearch-8.

Catatan

Hanya VVR 11.6 atau versi lebih baru yang mendukung nilai elasticsearch-8.

endPoint

Alamat server kluster Elasticsearch.

String

Ya

Tidak ada

Nama opsi lama.

hosts

Digunakan dengan elasticsearch-8.

indexName

Nama indeks.

String

Ya

Tidak ada

Nama opsi lama.

index

Untuk digunakan dengan elasticsearch-8.

accessId

Username untuk otentikasi.

String

Tidak

Tidak ada

Secara bawaan, parameter ini kosong dan otentikasi tidak dilakukan. Jika Anda menentukan accessId, Anda harus menentukan accessKey yang tidak kosong.

Catatan

elasticsearch-8 telah diperbarui untuk menggunakan parameter username dan password.

Penting

Untuk mencegah eksposur username dan password Anda, kami menyarankan agar Anda menggunakan variabel proyek. Untuk informasi selengkapnya, lihat variabel proyek.

username

accessKey

Password untuk otentikasi.

String

Tidak

Tidak ada

password

typeNames

Nama tipe.

String

Tidak

_doc

Kami menyarankan agar Anda tidak mengonfigurasi opsi ini untuk Elasticsearch 7.0 atau versi lebih baru.

batchSize

Jumlah maksimum dokumen yang diambil dari kluster Elasticsearch per permintaan scroll.

Int

Tidak

2000

Tidak ada

keepScrollAliveSecs

Waktu maksimum untuk menjaga konteks scroll tetap aktif.

Int

Tidak

3600

Unit: detik.

Tabel sink

Parameter

Deskripsi

Type

Wajib

Bawaan

Keterangan

connector

Tipe tabel sink.

String

Ya

Tidak ada

Nilai harus berupa elasticsearch-6, elasticsearch-7 , atau elasticsearch-8.

Catatan

Hanya VVR 8.0.5 atau versi lebih baru yang mendukung nilai elasticsearch-8.

hosts

Alamat server kluster Elasticsearch.

String

Ya

Tidak ada

Contoh: 127.0.0.1:XXXX.

index

Nama indeks.

String

Ya

Tidak ada

Tabel sink mendukung indeks statis maupun dinamis:

  • Untuk indeks statis, nilainya harus berupa string biasa, seperti myusers. Semua catatan ditulis ke indeks myusers.

  • Untuk indeks dinamis, Anda dapat menggunakan {field_name} untuk mereferensikan nilai bidang dari catatan guna menghasilkan indeks target secara dinamis. Anda juga dapat menggunakan {field_name|date_format_string} untuk mengonversi nilai bidang dengan tipe data TIMESTAMP, DATE, dan TIME ke format yang ditentukan oleh date_format_string. date_format_string kompatibel dengan DateTimeFormatter Java. Misalnya, jika Anda mengatur indeks menjadi myusers-{log_ts|yyyy-MM-dd}, catatan dengan nilai bidang log_ts 2020-03-27 12:25:55 akan ditulis ke indeks myusers-2020-03-27.

document-type

Tipe dokumen.

String

  • elasticsearch-6: Ya

  • elasticsearch-7: Tidak didukung

Tidak ada

Ketika tipe konektor adalah elasticsearch-6, nilai parameter ini harus konsisten dengan nilai parameter type di Elasticsearch.

username

Username untuk otentikasi.

String

Tidak

Tidak ada

Otentikasi dinonaktifkan secara bawaan. Jika Anda menentukan username, Anda juga harus menentukan password yang tidak kosong.

Penting

Untuk mencegah eksposur username dan password Anda, kami menyarankan agar Anda menggunakan variabel proyek. Untuk informasi selengkapnya, lihat variabel proyek.

password

Password untuk otentikasi.

String

Tidak

Tidak ada

document-id.key-delimiter

Pembatas untuk ID dokumen.

String

Tidak

_

Konektor menggunakan kunci primer untuk menghasilkan ID dokumen. Konektor menggabungkan semua bidang kunci primer sesuai urutan yang didefinisikan dalam DDL, menggunakan pembatas yang ditentukan oleh document-id.key-delimiter, untuk membuat string ID dokumen untuk setiap baris.

Catatan

ID dokumen adalah string hingga 512 byte yang tidak mengandung spasi.

failure-handler

Kebijakan penanganan kegagalan untuk permintaan Elasticsearch yang gagal.

String

Tidak

fail

Kebijakan yang valid:

  • fail (bawaan): Pekerjaan gagal jika permintaan gagal.

  • ignore: Mengabaikan kegagalan dan membuang permintaan.

  • retry-rejected: Menambahkan kembali permintaan yang gagal karena antrian penuh.

  • Nama kelas kustom: Menggunakan subclass ActionRequestFailureHandler untuk penanganan kegagalan.

sink.flush-on-checkpoint

Menentukan apakah flush dilakukan saat checkpoint.

Boolean

Tidak

true

  • true: Nilai bawaan.

  • false: Jika dinonaktifkan, konektor tidak menunggu semua permintaan tertunda diakui selama checkpoint. Artinya, konektor tidak memberikan jaminan pengiriman at-least-once.

sink.bulk-flush.backoff.strategy

Jika operasi flush gagal karena kesalahan permintaan sementara, atur sink.bulk-flush.backoff.strategy untuk menentukan strategi percobaan ulang.

Enum

Tidak

DISABLED

  • DISABLED (bawaan): Tidak ada percobaan ulang. Pekerjaan gagal pada kesalahan permintaan pertama.

  • CONSTANT: Strategi backoff konstan di mana waktu tunggu antar percobaan ulang selalu sama.

  • EXPONENTIAL: Strategi backoff eksponensial di mana waktu tunggu antar percobaan ulang meningkat secara eksponensial.

sink.bulk-flush.backoff.max-retries

Jumlah maksimum percobaan ulang.

Int

Tidak

Tidak ada

Tidak ada

sink.bulk-flush.backoff.delay

Penundaan antar percobaan ulang.

Durasi

Tidak

Tidak ada

  • Untuk strategi constant backoff, nilai ini adalah penundaan antar setiap percobaan ulang.

  • Untuk strategi exponential backoff, nilai ini adalah penundaan dasar awal.

sink.bulk-flush.max-actions

Jumlah maksimum aksi yang dibuffer untuk setiap permintaan bulk.

Int

Tidak

1000

Nilai 0 menonaktifkan fitur ini.

sink.bulk-flush.max-size

Ukuran memori maksimum buffer permintaan.

String

Tidak

2 MB

Unitnya adalah MB. Nilai bawaannya adalah 2 MB. Nilai 0 menonaktifkan fitur ini.

sink.bulk-flush.interval

Interval flush.

Durasi

Tidak

1s

Unitnya adalah detik. Nilai bawaannya adalah 1s. Nilai 0s menonaktifkan fitur ini.

connection.path-prefix

String yang ditambahkan di awal setiap path komunikasi REST.

String

Tidak

Tidak ada

Tidak ada

retry-on-conflict

Jumlah maksimum percobaan ulang untuk operasi pembaruan jika terjadi konflik versi. Jika jumlah percobaan ulang melebihi nilai ini, pekerjaan gagal dengan exception.

Int

Tidak

0

Catatan
  • Opsi ini hanya didukung di VVR 4.0.13 atau versi lebih baru.

  • Opsi ini hanya berlaku ketika kunci primer didefinisikan.

routing-fields

Menentukan satu atau beberapa nama bidang Elasticsearch yang digunakan untuk merutekan dokumen ke shard tertentu.

String

Tidak

Tidak ada

Pisahkan beberapa nama bidang dengan titik koma (;). Jika data bidang kosong, bidang tersebut diatur menjadi null.

Catatan

Opsi ini hanya didukung di VVR 8.0.6 atau versi lebih baru, untuk elasticsearch-7 dan elasticsearch-8.

sink.delete-strategy

Mengonfigurasi cara sink menangani pesan retraction (-D untuk DELETE atau -U untuk UPDATE_BEFORE).

Enum

Tidak

DELETE_ROW_ON_PK

Strategi yang valid:

  • DELETE_ROW_ON_PK (bawaan): Mengabaikan pesan -U tetapi menghapus baris (dokumen) yang sesuai dengan kunci primer ketika menerima pesan -D.

  • IGNORE_DELETE: Mengabaikan pesan -U dan -D. Tidak ada retraction yang terjadi di sink Elasticsearch.

  • NON_PK_FIELD_TO_NULL: Mengabaikan pesan -U. Namun, ketika pesan -D diterima, pengaturan ini memodifikasi baris (dokumen) untuk kunci primer: nilai kunci primer tetap sama, dan semua nilai bidang non-kunci primer lainnya dalam skema tabel diatur menjadi NULL. Ini terutama digunakan untuk pembaruan parsial ketika beberapa sink menulis ke tabel Elasticsearch yang sama secara bersamaan.

  • CHANGELOG_STANDARD: Mirip dengan DELETE_ROW_ON_PK, tetapi juga menghapus baris (dokumen) yang sesuai dengan kunci primer ketika pesan -U diterima.

    Catatan

    Opsi ini hanya didukung di VVR 8.0.8 atau versi lebih baru.

sink.ignore-null-when-update

Saat memperbarui data, menentukan apakah akan memperbarui bidang menjadi null atau membiarkannya tidak berubah jika nilai bidang masuk adalah null.

BOOLEAN

Tidak

false

Nilai yang valid:

  • true: Bidang tidak diperbarui. Nilai ini hanya didukung ketika kunci primer ditetapkan untuk tabel Flink dan format data Elasticsearch adalah JSON.

  • false: Bidang diperbarui menjadi null.

Catatan

Opsi ini hanya didukung di VVR 11.1 atau versi lebih baru.

connection.request-timeout

Timeout untuk meminta koneksi dari connection manager.

Durasi

Tidak

Tidak ada

Contoh:

'connection.request-timeout' = '1 min' -- 1 menit
'connection.request-timeout' = '500ms' -- 500 milidetik
Catatan

Opsi ini hanya didukung di VVR 11.7 atau versi lebih baru.

connect.timeout

Timeout untuk membuat koneksi.

Durasi

Tidak

Tidak ada

Contoh:

'connect.timeout' = '1 min' -- 1 menit
'connect.timeout' = '500ms' -- 500 milidetik
Catatan

Opsi ini hanya didukung di VVR 11.7 atau versi lebih baru.

socket.timeout

Timeout untuk menunggu data, yaitu periode maksimum ketidakaktifan antara dua paket data berturut-turut.

Durasi

Tidak

Tidak ada

Contoh:

'socket.timeout' = '1 min' -- 1 menit
'socket.timeout' = '500ms' -- 500 milidetik
Catatan

Opsi ini hanya didukung di VVR 11.7 atau versi lebih baru.

connection.keep-alive

Durasi maksimum koneksi dapat tetap idle sebelum sistem menutupnya. Jika opsi ini tidak diatur, header respons Keep-Alive dari server menentukan durasinya. Jika server tidak mengirim header Keep-Alive, koneksi tetap aktif tanpa batas.

Durasi

Tidak

Tidak ada

Contoh:

'connection.keep-alive' = '1 min' -- 1 menit
'connection.keep-alive' = '500ms' -- 500 milidetik
Catatan

Opsi ini hanya didukung di VVR 11.7 atau versi lebih baru.

sink.bulk-flush.update.doc_as_upsert

Menentukan apakah dokumen diperlakukan sebagai dokumen upsert dalam permintaan pembaruan.

BOOLEAN

Tidak

false

Nilai yang valid:

  • true: Mengatur bidang doc_as_upsert dari permintaan pembaruan menjadi true.

  • false: Mengisi bidang upsert dari permintaan pembaruan dengan dokumen.

Menurut https://github.com/elastic/elasticsearch/issues/105804, pipeline ingest Elasticsearch tidak mendukung pembaruan parsial untuk permintaan bulk update. Jika Anda ingin menggunakan pipeline ingest, atur opsi ini menjadi true.

Catatan

Opsi ini hanya didukung di VVR 11.5 atau versi lebih baru.

Tabel dimensi

Parameter

Deskripsi

Tipe

Wajib

Bawaan

Keterangan

connector

Tipe tabel dimensi.

String

Ya

Tidak ada

Nilai valid: elasticsearch atau elasticsearch-8.

Catatan

Hanya VVR 11.6 atau versi lebih baru yang mendukung nilai elasticsearch-8.

endPoint

Alamat server kluster Elasticsearch.

String

Ya

Tidak ada

Nama opsi lama.

hosts

Untuk digunakan dengan elasticsearch-8.

indexName

Nama indeks.

String

Ya

Tidak ada

Nama opsi lama.

index

Untuk digunakan dengan elasticsearch-8.

accessId

Username untuk otentikasi.

String

Tidak

Tidak ada

Secara bawaan, parameter ini kosong dan otentikasi tidak dilakukan. Jika Anda menentukan accessId, Anda harus menentukan accessKey yang tidak kosong.

Catatan

elasticsearch-8 kini menggunakan parameter username dan password.

Penting

Untuk mencegah eksposur username dan password Anda, kami menyarankan agar Anda menggunakan variabel proyek. Untuk informasi selengkapnya, lihat variabel proyek.

username

accessKey

Password untuk otentikasi.

String

Tidak

Tidak ada

password

typeNames

Nama tipe.

String

Tidak

_doc

Kami menyarankan agar Anda tidak mengonfigurasi opsi ini untuk Elasticsearch 7.0 atau versi lebih baru.

maxJoinRows

Jumlah maksimum baris yang digabungkan untuk satu lookup.

Integer

Tidak

1024

Tidak ada

cache

Strategi caching.

String

Tidak

Tidak ada

Nilai yang valid:

  • ALL: Menyimpan cache semua data dari tabel dimensi. Sebelum pekerjaan dimulai, sistem memuat semua data dari tabel dimensi ke cache. Lookup berikutnya dilayani langsung dari cache. Jika kunci tidak ditemukan, sistem menganggapnya tidak ada. Sistem memuat ulang seluruh cache ketika TTL berakhir.

  • LRU: Menyimpan cache sebagian data dari tabel dimensi. Ketika catatan dari tabel sumber tiba, sistem terlebih dahulu mencari data di cache. Jika terjadi cache miss, sistem melakukan query ke tabel dimensi fisik.

  • None: Tidak ada caching.

cacheSize

Ukuran cache, ditentukan sebagai jumlah baris.

Long

Tidak

100000

Parameter cacheSize hanya berlaku ketika kebijakan cache LRU dipilih untuk cache.

cacheTTLMs

Waktu hidup (TTL) untuk cache.

Long

Tidak

Long.MAX_VALUE

Unit: milidetik. Perilaku cacheTTLMs bergantung pada pengaturan cache:

  • Ketika cache adalah LRU, cacheTTLMs adalah TTL untuk entri cache. Secara bawaan, entri tidak kedaluwarsa.

  • Ketika cache adalah ALL, cacheTTLMs adalah interval untuk memuat ulang cache. Secara bawaan, cache tidak dimuat ulang.

ignoreKeywordSuffix

Menentukan apakah akan mengabaikan akhiran .keyword yang secara otomatis ditambahkan ke bidang STRING.

Boolean

Tidak

false

Untuk kompatibilitas, Flink mengonversi tipe Text dari Elasticsearch ke STRING dan secara bawaan menambahkan akhiran .keyword ke nama bidang.

Nilai yang valid:

  • true: Mengabaikan akhiran tersebut.

    Jika akhiran tersebut mencegah pencocokan dengan bidang tipe Text di Elasticsearch, atur opsi ini menjadi true.

  • false: Tidak mengabaikan akhiran tersebut.

cacheEmpty

Menentukan apakah hasil kosong dari lookup di tabel dimensi fisik disimpan dalam cache.

Boolean

Tidak

true

Parameter cacheEmpty hanya berlaku ketika cache menggunakan kebijakan cache LRU.

queryMaxDocs

Untuk tabel dimensi tanpa kunci primer, ini adalah jumlah maksimum dokumen yang dikembalikan oleh server Elasticsearch untuk setiap query lookup.

Integer

Tidak

10000

Nilai bawaan 10.000 sesuai dengan jumlah maksimum dokumen yang dapat dikembalikan oleh server Elasticsearch per query. Nilai ini tidak boleh melebihi batas tersebut.

Catatan
  • Opsi ini hanya didukung di VVR 8.0.8 atau versi lebih baru.

  • Opsi ini hanya berlaku untuk tabel dimensi tanpa kunci primer, karena data dalam tabel dengan kunci primer bersifat unik.

  • Nilai bawaan yang besar membantu memastikan kebenaran query tetapi meningkatkan penggunaan memori selama query Elasticsearch. Jika Anda mengalami masalah memori, Anda dapat mengurangi nilai ini untuk mengoptimalkan penggunaan memori.

Pemetaan tipe

Flink mengurai data Elasticsearch sebagai JSON. Untuk detailnya, lihat pemetaan tipe data.

Contoh

  • Contoh tabel sumber

    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;
  • Contoh tabel dimensi

    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;
  • Contoh tabel sink 1

    Contoh ini menulis konten teks ke Elasticsearch setelah vektorisasi teks.

    Catatan

    Buat pemetaan indeks di Elasticsearch terlebih dahulu. Atur tipe data bidang embedding menjadi dense_vector dan tentukan dimensinya. Jika tidak, Elasticsearch mungkin menginferensinya sebagai tipe array biasa.

        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 -- Kunci primer opsional. Jika Anda mendefinisikan kunci primer, nilainya menjadi ID dokumen. Jika tidak, ID dokumen acak dihasilkan.
        ) 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;
  • Contoh tabel sink 2

    Konektor mendukung penulisan tipe kompleks seperti ROW, ARRAY, dan MAP ke 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  -- Kunci primer opsional. Jika Anda mendefinisikan kunci primer, nilainya menjadi ID dokumen. Jika tidak, ID dokumen acak dihasilkan.
    ) 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;