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 |
|
|
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.
CatatanPenggunaan 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.keywordke nama bidang demi kompatibilitas. Jika hal ini mencegah pencocokan dengan bidang tipeTEXTdi Elasticsearch, atur opsiignoreKeywordSuffixketrue.
-
-
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
upsertatauappend, tergantung pada apakah kunci primer didefinisikan.-
Jika kunci primer didefinisikan, nilainya digunakan sebagai ID dokumen. Tabel sink kemudian beroperasi dalam mode
upsertdan dapat memproses operasiUPDATEdanDELETE. -
Jika tidak ada kunci primer yang didefinisikan, Elasticsearch secara otomatis menghasilkan ID dokumen acak. Tabel sink kemudian beroperasi dalam mode
appenddan hanya dapat mengonsumsi pesanINSERT.
-
-
Tipe data seperti
BYTES,ROW,ARRAY, danMAPtidak memiliki representasi string yang sesuai. Oleh karena itu, bidang dengan tipe data tersebut tidak dapat digunakan sebagai kunci primer. -
Bidang dalam
DDLbersesuaian 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: Catatan
Hanya VVR 11.6 atau versi lebih baru yang mendukung nilai |
|
endPoint |
Alamat server kluster Elasticsearch. |
String |
Ya |
Tidak ada |
Nama opsi lama. |
|
hosts |
Digunakan dengan |
||||
|
indexName |
Nama indeks. |
String |
Ya |
Tidak ada |
Nama opsi lama. |
|
index |
Untuk digunakan dengan |
||||
|
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
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 Catatan
Hanya VVR 8.0.5 atau versi lebih baru yang mendukung nilai |
|
hosts |
Alamat server kluster Elasticsearch. |
String |
Ya |
Tidak ada |
Contoh: |
|
index |
Nama indeks. |
String |
Ya |
Tidak ada |
Tabel sink mendukung indeks statis maupun dinamis:
|
|
document-type |
Tipe dokumen. |
String |
|
Tidak ada |
Ketika tipe konektor adalah |
|
username |
Username untuk otentikasi. |
String |
Tidak |
Tidak ada |
Otentikasi dinonaktifkan secara bawaan. Jika Anda menentukan 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:
|
|
sink.flush-on-checkpoint |
Menentukan apakah flush dilakukan saat checkpoint. |
Boolean |
Tidak |
true |
|
|
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 |
|
|
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 |
|
|
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
|
|
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 |
|
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:
|
|
sink.ignore-null-when-update |
Saat memperbarui data, menentukan apakah akan memperbarui bidang menjadi |
BOOLEAN |
Tidak |
false |
Nilai yang valid:
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:
Catatan
Opsi ini hanya didukung di VVR 11.7 atau versi lebih baru. |
|
connect.timeout |
Timeout untuk membuat koneksi. |
Durasi |
Tidak |
Tidak ada |
Contoh:
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:
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 |
Durasi |
Tidak |
Tidak ada |
Contoh:
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:
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: Catatan
Hanya VVR 11.6 atau versi lebih baru yang mendukung nilai |
|
endPoint |
Alamat server kluster Elasticsearch. |
String |
Ya |
Tidak ada |
Nama opsi lama. |
|
hosts |
Untuk digunakan dengan |
||||
|
indexName |
Nama indeks. |
String |
Ya |
Tidak ada |
Nama opsi lama. |
|
index |
Untuk digunakan dengan |
||||
|
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
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:
|
|
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:
|
|
ignoreKeywordSuffix |
Menentukan apakah akan mengabaikan akhiran .keyword yang secara otomatis ditambahkan ke bidang STRING. |
Boolean |
Tidak |
false |
Untuk kompatibilitas, Flink mengonversi tipe Nilai yang valid:
|
|
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
|
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.
CatatanBuat pemetaan indeks di Elasticsearch terlebih dahulu. Atur tipe data bidang
embeddingmenjadidense_vectordan 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, danMAPke 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;