Pelajari cara menggunakan konektor Log Service (SLS).
Latar Belakang
Simple Log Service adalah layanan end-to-end untuk data log yang membantu Anda mengumpulkan, mengonsumsi, mengirimkan, mengkueri, dan menganalisis data log secara efisien. Layanan ini meningkatkan efisiensi O&M serta memungkinkan pemrosesan volume data log yang sangat besar.
Tabel berikut mencantumkan kemampuan konektor SLS.
|
Kategori |
Deskripsi |
|
Tipe yang Didukung |
tabel sumber dan tabel sink |
|
running mode |
Hanya mode streaming |
|
Metrik khusus konektor |
N/A |
|
Format data |
N/A |
|
Tipe API |
SQL, DataStream API, dan data ingestion YAML API |
|
Pembaruan atau penghapusan data di tabel sink |
Tabel sink hanya append-only; Anda tidak dapat memperbarui atau menghapus data. |
Fitur
Konektor sumber SLS langsung membaca bidang atribut pesan. Tabel berikut mencantumkan bidang yang didukung.
|
Parameter |
Tipe |
Deskripsi |
|
__source__ |
STRING METADATA VIRTUAL |
Sumber pesan. |
|
__topic__ |
STRING METADATA VIRTUAL |
Topik pesan. |
|
__timestamp__ |
BIGINT METADATA VIRTUAL |
Waktu log. |
|
__tag__ |
MAP<VARCHAR, VARCHAR> METADATA VIRTUAL |
Tag pesan. Misalnya, untuk atribut |
Prasyarat
Pastikan Anda telah membuat Project Log Service dan Logstore. Untuk informasi selengkapnya, lihat Buat Project dan Logstore.
Batasan
-
Hanya Ververica Runtime (VVR) 11.1 dan versi yang lebih baru yang mendukung penggunaan SLS sebagai sumber sinkron untuk data ingestion yang didefinisikan dalam YAML.
-
Konektor SLS hanya menjamin semantik at-least-once.
-
Jangan mengatur
source parallelismlebih tinggi dari jumlah shard karena akan membuang sumber daya. Selain itu, pada Ververica Runtime (VVR) 8.0.5 dan versi sebelumnya, perubahan jumlah shard dapat menyebabkan fiturfailoverotomatis gagal, sehingga beberapa shard tidak dikonsumsi.
SQL
Sintaksis
CREATE TABLE sls_table(
a INT,
b INT,
c VARCHAR
) WITH (
'connector' = 'sls',
'endPoint' = '<yourEndPoint>',
'project' = '<yourProjectName>',
'logStore' = '<yourLogStoreName>',
'accessId' = '${secret_values.ak_id}',
'accessKey' = '${secret_values.ak_secret}'
);
Opsi
-
Umum
Parameter
Deskripsi
Tipe
Wajib
Default
Keterangan
connector
Konektor yang digunakan.
String
Ya
Tidak ada
Atur ke
sls.endPoint
Titik akhir Log Service (SLS).
String
Ya
Tidak ada
Tentukan alamat akses VPC Log Service (SLS). Untuk informasi selengkapnya, lihat Titik akhir layanan.
Catatan-
Secara default, Realtime Compute for Apache Flink tidak dapat mengakses Internet. Untuk mengaktifkan akses internet dari Virtual Private Cloud (VPC) Anda, gunakan NAT Gateway. Untuk informasi selengkapnya, lihat Bagaimana cara mengakses Internet?.
-
Kami menyarankan agar Anda tidak mengakses SLS melalui internet. Jika harus melakukannya, gunakan HTTPS dan aktifkan akselerasi transfer. Untuk informasi selengkapnya, lihat Kelola akselerasi transfer.
project
Nama project SLS.
String
Ya
Tidak ada
N/A
logStore
Nama Logstore atau MetricStore.
String
Ya
Tidak ada
Data dalam Logstore dikonsumsi dengan cara yang sama seperti data dalam MetricStore.
accessId
ID AccessKey Akun Alibaba Cloud Anda.
String
Ya
Tidak ada
Untuk informasi selengkapnya, lihat Dapatkan Pasangan Kunci Akses.
PentingUntuk mencegah Pasangan Kunci Akses Anda terpapar, kami menyarankan agar Anda menggunakan variabel untuk menentukan ID AccessKey dan rahasia AccessKey. Untuk informasi selengkapnya, lihat Variabel project.
accessKey
Rahasia AccessKey Akun Alibaba Cloud Anda.
String
Ya
Tidak ada
-
-
Khusus sumber
Parameter
Deskripsi
Tipe
Wajib
Default
Keterangan
enableNewSource
Menentukan apakah akan menggunakan sumber data baru yang mengimplementasikan antarmuka FLIP-27.
Boolean
Tidak
false
Sumber baru dapat secara otomatis menyesuaikan perubahan shard dan mendistribusikan shard seimbang di semua subtask sumber.
Penting-
Opsi ini hanya didukung di VVR 8.0.9 dan yang lebih baru. Nilai default adalah
trueuntuk VVR 11.1 dan yang lebih baru. -
Jika Anda mengubah nilai opsi ini, Anda tidak dapat memulihkan pekerjaan dari state yang disimpan. Untuk melanjutkan konsumsi dari offset historis, pertama-tama mulai pekerjaan dengan opsi
consumerGroupuntuk mencatat progres konsumsi di kelompok konsumen SLS. Kemudian, atur opsiconsumeFromCheckpointketruedan restart pekerjaan tanpa state. -
Jika Logstore memiliki shard read-only, beberapa subtask mungkin terus meminta data dari shard lain setelah menyelesaikan tugasnya sendiri. Hal ini dapat menyebabkan workload tidak seimbang dan memengaruhi performa. Untuk mengatasi masalah ini, Anda dapat menyesuaikan parallelism, mengoptimalkan strategi penjadwalan, atau menggabungkan shard kecil untuk mengurangi jumlah shard dan menyederhanakan alokasi tugas.
shardDiscoveryIntervalMs
Interval untuk penemuan shard dinamis.
Long
Tidak
60000
Atur opsi ini ke nilai negatif untuk menonaktifkan penemuan dinamis. Satuan: milidetik.
Catatan-
Nilai harus lebih besar dari atau sama dengan 60.000 milidetik (1 menit).
-
Opsi ini hanya berlaku ketika
enableNewSourcediatur ketrue. -
Opsi ini hanya didukung di VVR 8.0.9 dan yang lebih baru.
startupMode
Mode startup untuk tabel sumber.
String
Tidak
timestamp
-
timestamp(default): Mengonsumsi log mulai dari waktu mulai yang ditentukan. -
latest: Mulai mengonsumsi log dari offset terbaru. -
earliest: Mulai mengonsumsi log dari offset paling awal. -
consumer_group: Mulai konsumsi log dari offset yang direkam oleh kelompok konsumen. Jika kelompok konsumen belum merekam offset konsumsi untuk suatu shard, konsumsi dimulai dari offset paling awal.
Penting-
Untuk versi VVR sebelum 11.1, nilai consumer_group tidak didukung. Anda harus mengatur
consumeFromCheckpointketrue. Dalam kasus ini, konsumsi log dimulai dari offset yang direkam oleh kelompok konsumen yang ditentukan, dan pengaturan mode startup tidak berlaku.
startTime
Waktu mulai untuk konsumsi log.
String
Tidak
Waktu saat ini
Formatnya adalah
yyyy-MM-dd hh:mm:ss.Opsi ini hanya berlaku ketika
startupModediatur ketimestamp.CatatanOpsi
startTimedanstopTimedidasarkan pada atribut__receive_time__di SLS, bukan atribut__timestamp__.stopTime
Waktu akhir untuk konsumsi log.
String
Tidak
Tidak ada
Formatnya adalah
yyyy-MM-dd hh:mm:ss.Catatan-
Opsi ini hanya digunakan untuk mengonsumsi log historis dan harus diatur ke waktu di masa lalu. Jika Anda mengaturnya ke waktu di masa depan, konsumsi mungkin berhenti lebih awal jika tidak ada log baru yang ditulis, sehingga mengakibatkan gangguan aliran data tanpa pesan error.
-
Jika Anda ingin pekerjaan Flink berhenti setelah semua log dikonsumsi, Anda juga harus mengatur
exitAfterFinishketrue.
consumerGroup
Nama kelompok konsumen.
String
Tidak
Tidak ada
Kelompok konsumen digunakan untuk mencatat progres konsumsi. Anda dapat menentukan nama kustom tanpa format tetap.
CatatanPekerjaan Flink yang berbeda harus menggunakan kelompok konsumen yang berbeda. Jika beberapa pekerjaan Flink menggunakan kelompok konsumen yang sama, mereka tidak berkoordinasi dan masing-masing pekerjaan mengonsumsi semua data. Hal ini karena Flink tidak menggunakan kelompok konsumen SLS untuk penugasan partisi saat mengonsumsi data dari SLS. Akibatnya, setiap konsumen mengonsumsi pesan secara independen, meskipun mereka berbagi kelompok konsumen yang sama.
consumeFromCheckpoint
Menentukan apakah akan mengonsumsi dari checkpoint kelompok konsumen.
String
Tidak
false
-
true: Anda juga harus menentukan kelompok konsumen. Program Flink mulai mengonsumsi log dari checkpoint yang disimpan di kelompok konsumen. Jika kelompok konsumen tidak memiliki checkpoint yang sesuai, konsumsi dimulai dari nilai konfigurasi startTime. -
false(nilai default): Tidak mulai mengonsumsi log dari checkpoint yang disimpan untuk kelompok konsumen yang ditentukan.
PentingParameter ini tidak lagi didukung di VVR 11.1 dan yang lebih baru. Untuk versi tersebut, Anda harus mengatur opsi
startupModekeconsumer_group.maxRetries
Jumlah percobaan ulang setelah gagal membaca dari SLS.
String
Tidak
3
N/A
batchGetSize
Jumlah kelompok log yang dibaca dalam satu permintaan.
String
Tidak
100
Pengaturan
batchGetSizetidak boleh melebihi 1000, atau akan muncul error.exitAfterFinish
Menentukan apakah pekerjaan Flink berhenti setelah semua data dikonsumsi.
String
Tidak
false
-
true: Program Flink berhenti setelah semua data dikonsumsi. -
false(default): Program Flink tidak berhenti setelah konsumsi data selesai.
query
PentingOpsi ini sudah tidak digunakan lagi di VVR 11.3, tetapi versi yang lebih baru tetap kompatibel.
Pernyataan kueri untuk pra-pemrosesan data sebelum konsumsi.
String
Tidak
Tidak ada
Gunakan opsi ini untuk memfilter data SLS sebelum Flink mengonsumsinya. Hal ini mengurangi biaya dan meningkatkan kecepatan pemrosesan.
Misalnya,
'query' = '*| where request_method = ''GET'''menunjukkan bahwa sebelum Flink membaca data dari SLS, data terlebih dahulu dicocokkan dengan nilai bidangrequest_methodyang bernilai 'GET'.CatatanOpsi ini menggunakan bahasa SPL Log Service (SLS). Untuk informasi selengkapnya, lihat Sintaksis SPL.
Penting-
Opsi ini hanya didukung di VVR 8.0.1 dan yang lebih baru.
-
Fitur ini dikenai biaya Log Service (SLS). Untuk informasi selengkapnya, lihat Penagihan Log Service.
processor
Nama prosesor SLS untuk pra-pemrosesan data. Jika
querydanprocessorkeduanya ditentukan,queryyang diutamakan danprocessordiabaikan.String
Tidak
Tidak ada
Opsi ini memfilter data SLS sebelum Flink mengonsumsinya, yang mengurangi biaya dan meningkatkan kecepatan pemrosesan. Kami menyarankan agar Anda menggunakan
processordaripadaquery.Misalnya,
'processor' = 'test-filter-processor'menunjukkan bahwa prosesor SLS memfilter data sebelum Flink membaca data dari SLS.CatatanOpsi ini menggunakan bahasa SPL Log Service (SLS). Untuk informasi selengkapnya, lihat Sintaksis SPL. Untuk informasi tentang cara membuat atau memperbarui prosesor SLS, lihat Kelola prosesor.
PentingOpsi ini hanya didukung di VVR 11.3 dan yang lebih baru.
Fitur ini dikenai biaya Log Service (SLS). Untuk informasi selengkapnya, lihat Penagihan Log Service.
preserveRawBytes
Menentukan apakah akan langsung mengambil dan menyimpan byte mentah yang dibawa oleh SLS ketika
enableNewSourcediatur ketrue.Boolean
Tidak
false
Ketika opsi ini diaktifkan, bidang BINARY/VARBINARY membaca byte[] mentah secara langsung, dan bidang CHAR/VARCHAR membuat data string dari byte mentah tersebut. Bidang lain tetap menggunakan logika konversi default. Perilaku ini konsisten dengan sumber lama.
Catatan-
Opsi ini hanya didukung di VVR 11.8 dan yang lebih baru.
-
Opsi ini hanya berlaku ketika
enableNewSourcediatur ketrue.
-
-
Khusus sink
Parameter
Deskripsi
Tipe
Wajib
Default
Keterangan
topicField
Menentukan bidang yang nilainya menimpa atribut
__topic__, yang menunjukkan topik log.String
Tidak
Tidak ada
Nilai opsi ini harus merupakan bidang yang ada di tabel.
timeField
Menentukan bidang yang nilainya menimpa atribut
__timestamp__, yang menunjukkan waktu penulisan log.String
Tidak
Waktu saat ini
Nilai opsi ini harus merupakan bidang
INTyang ada di tabel. Jika opsi ini tidak ditentukan, waktu saat ini yang digunakan.sourceField
Menentukan bidang yang nilainya menimpa atribut
__source__, yang menunjukkan sumber log, seperti alamat IP mesin yang menghasilkan log.String
Tidak
Tidak ada
Nilai opsi ini harus merupakan bidang yang ada di tabel.
partitionField
Menentukan bidang untuk partisi. Hash dari nilai bidang ini menentukan shard mana yang menerima data, memastikan catatan dengan hash yang sama masuk ke shard yang sama.
String
Tidak
Tidak ada
Jika opsi ini tidak ditentukan, setiap catatan ditulis secara acak ke shard yang tersedia.
buckets
Jika
partitionFieldditentukan, opsi ini menentukan jumlah bucket untuk memetakan nilai hash.String
Tidak
64
Nilai harus merupakan pangkat 2 dalam rentang [1, 256]. Jumlah bucket harus lebih besar dari atau sama dengan jumlah shard. Jika tidak, beberapa shard mungkin tidak menerima data apa pun.
flushIntervalMs
Interval penulisan data.
String
Tidak
2000
Satuan: milidetik.
writeNullProperties
Menentukan apakah nilai null ditulis sebagai string kosong ke SLS.
Boolean
Tidak
true
-
true(nilai default): Menulis nilai null ke log sebagai string kosong. -
false: Bidang yang bernilai null tidak ditulis ke log.
CatatanOpsi ini hanya didukung di VVR 8.0.6 dan yang lebih baru.
-
Pemetaan tipe
|
Tipe Flink |
Tipe SLS |
|
BOOLEAN |
STRING |
|
VARBINARY |
|
|
VARCHAR |
|
|
TINYINT |
|
|
INTEGER |
|
|
BIGINT |
|
|
FLOAT |
|
|
DOUBLE |
|
|
DECIMAL |
Data ingestion (Beta)
Batasan
Fitur ini hanya didukung oleh Realtime Compute for Apache Flink versi 11.1 dan yang lebih baru.
Sintaksis
source:
type: sls
name: SLS Source
endpoint: <endpoint>
project: <project>
logstore: <logstore>
accessId: <accessId>
accessKey: <accessKey>
Parameter
|
Parameter |
Deskripsi |
Tipe |
Wajib |
Default |
Keterangan |
||||||||||
|
type |
Tipe sumber data. |
String |
Ya |
Tidak ada |
Nilainya harus |
||||||||||
|
endpoint |
Titik akhir Log Service (SLS). |
String |
Ya |
Tidak ada |
Alamat akses VPC Log Service (SLS). Untuk informasi selengkapnya, lihat Titik Akhir Layanan. Catatan
|
||||||||||
|
accessId |
ID AccessKey untuk Akun Alibaba Cloud Anda. |
String |
Ya |
Tidak ada |
Untuk informasi selengkapnya, lihat Bagaimana cara melihat Informasi AccessKey ID dan Rahasia AccessKey?. Penting
Untuk mencegah Informasi AccessKey Anda terpapar, kami menyarankan agar Anda menggunakan variabel project untuk menentukan nilai AccessKey. Untuk informasi selengkapnya, lihat Variabel project. |
||||||||||
|
accessKey |
Rahasia AccessKey untuk Akun Alibaba Cloud Anda. |
String |
Ya |
Tidak ada |
|||||||||||
|
project |
Nama project Log Service (SLS). |
String |
Ya |
Tidak ada |
Tidak ada |
||||||||||
|
logStore |
Nama Logstore atau Metricstore SLS. |
String |
Ya |
Tidak ada |
Data dalam Logstore dikonsumsi dengan cara yang sama seperti data dalam Metricstore. |
||||||||||
|
schema.inference.strategy |
Strategi inferensi skema. |
String |
Tidak |
continuous |
|
||||||||||
|
maxPreFetchLogGroups |
Jumlah maksimum kelompok log yang dibaca dari setiap shard untuk inferensi skema awal. |
Integer |
Tidak |
50 |
Sebelum pekerjaan membaca dan memproses data, konektor melakukan pre-consume sejumlah kelompok log yang ditentukan dari setiap shard untuk menginisialisasi informasi skema. |
||||||||||
|
shardDiscoveryIntervalMs |
Interval, dalam milidetik, untuk menemukan perubahan shard secara dinamis. |
Long |
Tidak |
60000 |
Atur parameter ini ke nilai negatif untuk menonaktifkan penemuan dinamis. Catatan
Nilai harus lebih besar dari atau sama dengan 60.000 milidetik (1 menit). |
||||||||||
|
startupMode |
Mode startup. |
String |
Tidak |
Tidak ada |
|
||||||||||
|
startTime |
Waktu mulai untuk konsumsi log. |
String |
Tidak |
Waktu saat ini |
Formatnya adalah Parameter ini hanya berlaku ketika Catatan
Parameter |
||||||||||
|
stopTime |
Waktu akhir untuk konsumsi log. |
String |
Tidak |
Tidak ada |
Formatnya adalah Catatan
Jika Anda ingin pekerjaan Flink berhenti setelah semua log dikonsumsi, Anda juga harus mengatur |
||||||||||
|
consumerGroup |
Nama kelompok konsumen. |
String |
Tidak |
Tidak ada |
Kelompok konsumen mencatat progres konsumsi. Anda dapat menentukan nama kustom apa pun. |
||||||||||
|
batchGetSize |
Jumlah kelompok log yang dibaca per permintaan. |
Integer |
Tidak |
100 |
Nilai batchGetSize tidak boleh melebihi 1.000. Jika tidak, akan terjadi error. |
||||||||||
|
maxRetries |
Jumlah percobaan ulang jika pembacaan dari Log Service (SLS) gagal. |
Integer |
Tidak |
3 |
Tidak ada |
||||||||||
|
exitAfterFinish |
Menentukan apakah pekerjaan Flink berhenti setelah semua data dikonsumsi. |
Boolean |
Tidak |
false |
|
||||||||||
|
query |
Pernyataan pra-pemrosesan untuk mengonsumsi data dari Log Service (SLS). |
String |
Tidak |
Tidak ada |
Gunakan parameter ini untuk memfilter data di Log Service (SLS) sebelum konsumsi guna menghemat biaya dan meningkatkan kecepatan pemrosesan. Misalnya, Catatan
Kueri harus menggunakan sintaksis SPL Log Service. Untuk informasi selengkapnya, lihat Sintaksis SPL. Penting
|
||||||||||
|
compressType |
Tipe kompresi untuk Log Service (SLS). |
String |
Tidak |
Tidak ada |
Tipe kompresi yang didukung meliputi:
|
||||||||||
|
timeZone |
Zona waktu untuk |
String |
Tidak |
Tidak ada |
Secara default, tidak ada offset yang ditambahkan. |
||||||||||
|
regionId |
Wilayah tempat Log Service (SLS) ditempatkan. |
String |
Tidak |
Tidak ada |
Untuk informasi selengkapnya, lihat Wilayah yang didukung. |
||||||||||
|
signVersion |
Versi signature permintaan untuk Log Service (SLS). |
String |
Tidak |
Tidak ada |
Untuk informasi selengkapnya, lihat Signature permintaan. |
||||||||||
|
shardModDivisor |
Divisor yang digunakan saat membaca dari shard Logstore. |
Int |
Tidak |
-1 |
Untuk informasi selengkapnya, lihat Shard. |
||||||||||
|
shardModRemainder |
Sisa yang digunakan saat membaca dari shard Logstore. |
Int |
Tidak |
-1 |
Untuk informasi selengkapnya, lihat Shards. |
||||||||||
|
metadata.list |
Kolom metadata yang diteruskan ke pekerjaan downstream. |
String |
Tidak |
Tidak ada |
Bidang metadata yang tersedia meliputi |
||||||||||
|
decode.table-id.fields |
Menentukan bidang yang nilainya digunakan untuk menghasilkan ID Tabel saat mengurai data log dari Log Service (SLS). |
String |
Tidak |
Tidak ada |
Beberapa bidang dipisahkan oleh koma Inggris
Catatan
Parameter ini didukung di Realtime Compute for Apache Flink versi 11.6 dan yang lebih baru. |
||||||||||
|
fixed-types |
Menentukan tipe data untuk bidang tertentu saat mengurai data log dari Log Service (SLS). |
String |
Tidak |
Tidak ada |
Saat mengurai data, tentukan tipe untuk bidang tertentu. Gunakan koma Catatan
Parameter ini didukung di Realtime Compute for Apache Flink versi 11.6 dan yang lebih baru. |
||||||||||
|
timestamp-format.standard |
Format untuk bidang timestamp dalam data log dari Log Service (SLS). |
String |
Tidak |
SQL |
Nilai yang valid:
Catatan
Parameter ini didukung di Realtime Compute for Apache Flink versi 11.6 dan yang lebih baru. |
||||||||||
|
ingestion.ignore-errors |
Menentukan apakah akan mengabaikan error yang terjadi selama penguraian data. |
Boolean |
Tidak |
false |
Catatan
Parameter ini didukung di Realtime Compute for Apache Flink versi 11.6 dan yang lebih baru. |
||||||||||
|
ingestion.error-tolerance.max-count |
Jika |
Integer |
Tidak |
-1 |
Parameter ini hanya berlaku ketika Catatan
Parameter ini didukung di Realtime Compute for Apache Flink versi 11.6 dan yang lebih baru. |
Gunakan katalog yang sudah ada
Mulai dari Realtime Compute for Apache Flink versi 11.5, Anda dapat mereferensikan Katalog SLS bawaan yang dibuat di halaman Data Management dalam pekerjaan data ingestion Flink CDC, sehingga tidak perlu menentukan properti koneksi secara manual.
source:
type: sls
using.built-in-catalog: sls_catalog
Saat ini, pekerjaan data ingestion dapat secara otomatis menggunakan kembali parameter berikut dari Katalog SLS bawaan:
-
endpoint
-
project
-
accessId
-
accessKey
Untuk mengganti parameter yang digunakan kembali secara otomatis ini, Anda dapat mendefinisikannya secara eksplisit dalam konfigurasi YAML. Parameter yang didefinisikan dalam file YAML memiliki prioritas lebih tinggi.
Pemetaan tipe data
Jika fixed-types tidak dikonfigurasi, pemetaan tipe data berikut berlaku:
|
Tipe SLS |
Tipe CDC |
|
STRING |
STRING |
Jika fixed-types dikonfigurasi, sistem mengurai data dengan menggunakan tipe yang ditentukan.
Inferensi dan evolusi skema
-
Pre-konsumsi data shard dan inisialisasi skema
Konektor SLS mempertahankan skema Logstore yang sedang dibacanya. Sebelum membaca data dari Logstore, konektor melakukan pre-consume hingga
maxPreFetchLogGroupskelompok log dari setiap shard, mengurai skema setiap entri log, lalu menggabungkannya untuk menginisialisasi skema tabel. Event pembuatan tabel kemudian dihasilkan berdasarkan skema awal ini sebelum konsumsi data dimulai.CatatanUntuk setiap shard, konektor mencoba mengonsumsi data mulai dari satu jam sebelum waktu saat ini guna mengurai skema log.
-
Informasi kunci primer
Log Log Service (SLS) tidak berisi informasi kunci primer. Anda dapat menambahkan kunci primer ke tabel secara manual dengan menggunakan aturan transformasi:
transform: - source-table: <project>.<logstore> projection: * primary-keys: key1, key2 -
Inferensi skema dan perubahan skema
Setelah skema diinisialisasi, jika schema.inference.strategy diatur ke
static, konektor SLS mengurai setiap entri log berdasarkan skema awal dan tidak menghasilkan event perubahan skema. Jika schema.inference.strategy diatur kecontinuous, konektor mengurai setiap entri log, menginferensi kolom fisik, lalu membandingkannya dengan skema saat ini. Jika skema yang diinferensi tidak konsisten dengan skema saat ini, skema tersebut digabungkan sesuai aturan berikut:-
Jika kolom fisik yang diinferensi berisi bidang yang tidak ada di skema saat ini, konektor menambahkan bidang tersebut ke skema dan menghasilkan event untuk menambahkan kolom nullable.
-
Jika kolom fisik yang diinferensi tidak memiliki bidang yang ada di skema saat ini, konektor mempertahankan bidang tersebut, mengisi datanya dengan NULL, dan tidak menghasilkan event penghapusan kolom.
Konektor SLS menginferensi tipe data semua bidang dalam setiap entri log sebagai String. Saat ini, hanya penambahan kolom baru yang didukung. Konektor menambahkan kolom baru ke akhir skema dan mengaturnya sebagai nullable.
-
Contoh kode
-
SQL untuk tabel sumber dan sink
CREATE TEMPORARY TABLE sls_input( `time` BIGINT, url STRING, dt STRING, float_field FLOAT, double_field DOUBLE, boolean_field BOOLEAN, `__topic__` STRING METADATA VIRTUAL, `__source__` STRING METADATA VIRTUAL, `__timestamp__` STRING METADATA VIRTUAL, __tag__ MAP<VARCHAR, VARCHAR> METADATA VIRTUAL, proctime as PROCTIME() ) WITH ( 'connector' = 'sls', 'endpoint' ='cn-hangzhou-intranet.log.aliyuncs.com', 'accessId' = '${secret_values.ak_id}', 'accessKey' = '${secret_values.ak_secret}', 'starttime' = '2023-08-30 00:00:00', 'project' ='sls-test', 'logstore' ='sls-input' ); CREATE TEMPORARY TABLE sls_sink( `time` BIGINT, url STRING, dt STRING, float_field FLOAT, double_field DOUBLE, boolean_field BOOLEAN, `__topic__` STRING, `__source__` STRING, `__timestamp__` BIGINT , receive_time BIGINT ) WITH ( 'connector' = 'sls', 'endpoint' ='cn-hangzhou-intranet.log.aliyuncs.com', 'accessId' = '${ak_id}', 'accessKey' = '${ak_secret}', 'project' ='sls-test', 'logstore' ='sls-output' ); INSERT INTO sls_sink SELECT `time`, url, dt, float_field, double_field, boolean_field, `__topic__` , `__source__` , `__timestamp__` , cast(__tag__['__receive_time__'] as bigint) as receive_time FROM sls_input; -
Data ingestion dengan sumber data SLS
Gunakan SLS sebagai sumber data untuk mengingesti data secara real-time ke sistem downstream yang didukung. Misalnya, konfigurasi berikut mendefinisikan pekerjaan data ingestion yang menulis data dari logstore ke data lake berformat Paimon di Data Lake Formation (DLF). Pekerjaan ini secara otomatis menginferensi skema tabel sink dan mendukung evolusi skema saat runtime.
source:
type: sls
name: SLS Source
endpoint: ${endpoint}
project: test_project
logstore: test_log
accessId: ${accessId}
accessKey: ${accessKey}
# Tambahkan kunci primer ke tabel.
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
# Rutekan semua data dari test_project.test_log ke tabel test_database.inventory.
route:
- source-table: test_project.test_log
sink-table: test_database.inventory
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (Opsional) Aktifkan deletion vectors untuk meningkatkan performa baca.
table.properties.deletion-vectors.enabled: true
DataStream API
Untuk membaca atau menulis data dengan DataStream API, gunakan konektor DataStream. Untuk informasi selengkapnya, lihat Penggunaan konektor DataStream.
Jika Anda menggunakan versi VVR sebelum 8.0.10, pekerjaan Anda mungkin gagal dimulai karena dependensi yang hilang. Untuk mengatasi masalah ini, tambahkan uber-JAR yang sesuai sebagai dependensi tambahan.
Baca dari SLS
Realtime Compute for Apache Flink menyediakan kelas SlsSourceFunction, implementasi dari SourceFunction, untuk membaca data dari Simple Log Service (SLS). Contoh berikut membaca data dari SLS.
public class SlsDataStreamSource {
public static void main(String[] args) throws Exception {
// Menyiapkan lingkungan eksekusi streaming
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// Membuat sumber SLS dan mencetak data ke konsol.
env.addSource(createSlsSource())
.map(SlsDataStreamSource::convertMessages)
.print();
env.execute("SLS Stream Source");
}
private static SlsSourceFunction createSlsSource() {
SLSAccessInfo accessInfo = new SLSAccessInfo();
accessInfo.setEndpoint("yourEndpoint");
accessInfo.setProjectName("yourProject");
accessInfo.setLogstore("yourLogStore");
accessInfo.setAccessId("yourAccessId");
accessInfo.setAccessKey("yourAccessKey");
// Ukuran batch get wajib diatur.
accessInfo.setBatchGetSize(10);
// Parameter opsional
accessInfo.setConsumerGroup("yourConsumerGroup");
accessInfo.setMaxRetries(3);
// Waktu mulai konsumsi, diatur ke waktu saat ini.
int startInSec = (int) (new Date().getTime() / 1000);
// Waktu berhenti konsumsi, di mana -1 berarti tidak pernah berhenti.
int stopInSec = -1;
return new SlsSourceFunction(accessInfo, startInSec, stopInSec);
}
private static List<String> convertMessages(SourceRecord input) {
List<String> res = new ArrayList<>();
for (FastLogGroup logGroup : input.getLogGroups()) {
int logsCount = logGroup.getLogsCount();
for (int i = 0; i < logsCount; i++) {
FastLog log = logGroup.getLogs(i);
int fieldCount = log.getContentsCount();
for (int idx = 0; idx < fieldCount; idx++) {
FastLogContent f = log.getContents(idx);
res.add(String.format("key: %s, value: %s", f.getKey(), f.getValue()));
}
}
}
return res;
}
}
Tulis ke SLS
Realtime Compute for Apache Flink menyediakan kelas SLSOutputFormat, implementasi dari OutputFormat, untuk menulis data ke SLS. Contoh berikut menulis data ke SLS.
public class SlsDataStreamSink {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.fromSequence(0, 100)
.map((MapFunction<Long, SinkRecord>) aLong -> getSinkRecord(aLong))
.addSink(createSlsSink())
.name(SlsDataStreamSink.class.getSimpleName());
env.execute("SLS Stream Sink");
}
private static OutputFormatSinkFunction createSlsSink() {
Configuration conf = new Configuration();
conf.setString(SLSOptions.ENDPOINT, "yourEndpoint");
conf.setString(SLSOptions.PROJECT, "yourProject");
conf.setString(SLSOptions.LOGSTORE, "yourLogStore");
conf.setString(SLSOptions.ACCESS_ID, "yourAccessId");
conf.setString(SLSOptions.ACCESS_KEY, "yourAccessKey");
SLSOutputFormat outputFormat = new SLSOutputFormat(conf);
return new OutputFormatSinkFunction<>(outputFormat);
}
private static SinkRecord getSinkRecord(Long seed) {
SinkRecord record = new SinkRecord();
LogItem logItem = new LogItem((int) (System.currentTimeMillis() / 1000));
logItem.PushBack("level", "info");
logItem.PushBack("name", String.valueOf(seed));
logItem.PushBack("message", "it's a test message for " + seed.toString());
record.setContent(logItem);
return record;
}
}
XML
Konektor DataStream SLS tersedia di repositori pusat Maven.
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-sls</artifactId>
<version>${vvr-version}</version>
<exclusions>
<exclusion>
<groupId>org.apache.flink</groupId>
<artifactId>flink-format-common</artifactId>
</exclusion>
</exclusions>
</dependency>