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 |
Jenis yang Didukung | Tabel sumber dan tabel sink |
running mode | Hanya mode streaming |
Metrik khusus konektor | N/A |
Format data | N/A |
Jenis API | SQL, DataStream API, dan API YAML data ingestion |
Pembaruan atau penghapusan data di tabel sink | Tabel sink hanya append-only; Anda tidak dapat memperbarui atau menghapus data. |
Fitur
Konektor sumber SLS membaca langsung 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 lebih lanjut, 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 hal ini 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
Jenis
Wajib
Default
Keterangan
connector
Konektor yang digunakan.
String
Ya
None
Atur ke
sls.endPoint
Titik akhir Log Service (SLS).
String
Ya
None
Tentukan alamat akses VPC Log Service (SLS). Untuk informasi lebih lanjut, lihat Service endpoints.
CatatanSecara 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 lebih lanjut, lihat How do I access the Internet?.
Kami menyarankan agar Anda tidak mengakses SLS melalui internet. Jika harus melakukannya, gunakan HTTPS dan aktifkan transfer acceleration. Untuk informasi lebih lanjut, lihat Manage transfer acceleration.
project
Nama project SLS.
String
Ya
None
N/A
logStore
Nama Logstore atau MetricStore.
String
Ya
None
Data dalam Logstore dikonsumsi dengan cara yang sama seperti data dalam MetricStore.
accessId
ID AccessKey Akun Alibaba Cloud Anda.
String
Ya
None
Untuk informasi lebih lanjut, lihat Obtain an AccessKey pair.
PentingUntuk mencegah Pasangan Kunci Akses Anda terpapar, kami menyarankan agar Anda menggunakan variabel untuk menentukan ID AccessKey dan rahasia AccessKey. Untuk informasi lebih lanjut, lihat Project variables.
accessKey
Rahasia AccessKey Akun Alibaba Cloud Anda.
String
Ya
None
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 semerata mungkin di seluruh subtask sumber.
PentingOpsi ini hanya didukung di VVR 8.0.9 dan versi yang lebih baru. Nilai default adalah
trueuntuk VVR 11.1 dan versi 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 dalam 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, 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.
CatatanNilai 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 versi 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 dicatat oleh kelompok konsumen. Jika kelompok konsumen belum mencatat offset konsumsi untuk suatu shard, konsumsi dimulai dari offset paling awal.
PentingUntuk versi VVR sebelum 11.1, nilai consumer_group tidak didukung. Anda harus mengatur
consumeFromCheckpointketrue. Dalam kasus ini, konsumsi log dimulai dari offset yang dicatat 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
None
Formatnya adalah
yyyy-MM-dd hh:mm:ss.CatatanOpsi 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 keluar setelah semua log dikonsumsi, Anda juga harus mengatur
exitAfterFinishketrue.
consumerGroup
Nama kelompok konsumen.
String
Tidak
None
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 dalam 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 versi yang lebih baru. Untuk versi tersebut, Anda harus mengatur opsi
startupModekeconsumer_group.maxRetries
Jumlah percobaan ulang setelah upaya membaca dari SLS gagal.
String
Tidak
3
N/A
batchGetSize
Jumlah kelompok log yang dibaca dalam satu permintaan.
String
Tidak
100
Pengaturan
batchGetSizetidak boleh melebihi 1.000. Jika dilanggar, sistem akan melaporkan error.exitAfterFinish
Menentukan apakah pekerjaan Flink keluar setelah semua data dikonsumsi.
String
Tidak
false
true: Program Flink keluar setelah semua data dikonsumsi.false(default): Program Flink tidak keluar 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 preprocessing data sebelum konsumsi.
String
Tidak
None
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, sistem terlebih dahulu mencocokkan data di mana nilai bidangrequest_methodadalah 'GET'.CatatanOpsi ini menggunakan bahasa SPL Log Service (SLS). Untuk informasi lebih lanjut, lihat SPL syntax.
PentingOpsi ini hanya didukung di VVR 8.0.1 dan versi yang lebih baru.
Fitur ini dikenai biaya Log Service (SLS). Untuk informasi lebih lanjut, lihat Billing of Log Service.
processor
Nama prosesor SLS untuk preprocessing data. Jika
querydanprocessorkeduanya ditentukan,querymemiliki prioritas lebih tinggi danprocessordiabaikan.String
Tidak
None
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 lebih lanjut, lihat SPL syntax. Untuk informasi tentang cara membuat atau memperbarui prosesor SLS, lihat Manage processors.
PentingOpsi ini hanya didukung di VVR 11.3 dan versi yang lebih baru.
Fitur ini dikenai biaya Log Service (SLS). Untuk informasi lebih lanjut, lihat Billing of Log Service.
Khusus sink
Parameter
Deskripsi
Tipe
Wajib
Default
Keterangan
topicField
Menentukan bidang yang nilainya menimpa atribut
__topic__, yang menunjukkan topik log.String
Tidak
None
Nilai opsi ini harus merupakan bidang yang ada dalam 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 dalam tabel. Jika opsi ini tidak ditentukan, waktu saat ini digunakan.sourceField
Menentukan bidang yang nilainya menimpa atribut
__source__, yang menunjukkan sumber log, seperti alamat IP mesin yang menghasilkan log.String
Tidak
None
Nilai opsi ini harus merupakan bidang yang ada dalam 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
None
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 versi yang lebih baru.
Pemetaan tipe
Tipe Flink | Tipe SLS |
BOOLEAN | STRING |
VARBINARY | |
VARCHAR | |
TINYINT | |
INTEGER | |
BIGINT | |
FLOAT | |
DOUBLE | |
DECIMAL |
Data ingestion (Pratinjau publik)
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 | Jenis sumber data. | String | Ya | None | Nilainya harus | ||||||||||
endpoint | Titik akhir Log Service (SLS). | String | Ya | None | Alamat akses VPC Log Service (SLS). Untuk informasi lebih lanjut, lihat Service Endpoints. Catatan
| ||||||||||
accessId | ID AccessKey untuk Akun Alibaba Cloud Anda. | String | Ya | None | Untuk informasi selengkapnya, lihat Cara melihat informasi ID AccessKey dan Rahasia AccessKey. Penting Untuk mencegah Informasi AccessKey Anda terpapar, kami menyarankan agar Anda menggunakan variabel proyek untuk menentukan nilai AccessKey. Untuk informasi lebih lanjut, lihat Project variables. | ||||||||||
accessKey | Rahasia AccessKey untuk Akun Alibaba Cloud Anda. | String | Ya | None | |||||||||||
project | Nama project Log Service (SLS). | String | Ya | None | None | ||||||||||
logStore | Nama Logstore atau Metricstore SLS. | String | Ya | None | 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 | None |
| ||||||||||
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 | None | Formatnya adalah Catatan Jika Anda ingin pekerjaan Flink keluar setelah semua log dikonsumsi, Anda juga harus mengatur | ||||||||||
consumerGroup | Nama kelompok konsumen. | String | Tidak | None | 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 dilanggar, sistem akan melaporkan error. | ||||||||||
maxRetries | Jumlah percobaan ulang jika pembacaan dari Log Service (SLS) gagal. | Integer | Tidak | 3 | None | ||||||||||
exitAfterFinish | Menentukan apakah pekerjaan Flink keluar setelah semua data dikonsumsi. | Boolean | Tidak | false |
| ||||||||||
query | Pernyataan preprocessing untuk mengonsumsi data dari Log Service (SLS). | String | Tidak | None | Gunakan parameter ini untuk memfilter data di Log Service (SLS) sebelum konsumsi guna menghemat biaya dan meningkatkan kecepatan pemrosesan. Misalnya, Catatan Kueri harus menggunakan sintaks SPL Log Service. Untuk informasi lebih lanjut, lihat SPL syntax. Penting
| ||||||||||
compressType | Jenis kompresi untuk Log Service (SLS). | String | Tidak | None | Jenis kompresi yang didukung meliputi:
| ||||||||||
timeZone | Zona waktu untuk | String | Tidak | None | Secara default, tidak ada offset yang ditambahkan. | ||||||||||
regionId | Wilayah tempat Log Service (SLS) ditempatkan. | String | Tidak | None | Untuk informasi lebih lanjut, lihat Supported regions. | ||||||||||
signVersion | Versi signature permintaan untuk Log Service (SLS). | String | Tidak | None | Untuk informasi lebih lanjut, lihat Request signatures. | ||||||||||
shardModDivisor | Pembagi yang digunakan saat membaca dari shard Logstore. | Int | Tidak | -1 | Untuk informasi lebih lanjut, lihat Shards. | ||||||||||
shardModRemainder | Sisa yang digunakan saat membaca dari shard Logstore. | Int | Tidak | -1 | Untuk informasi lebih lanjut, lihat Shards. | ||||||||||
metadata.list | Kolom metadata yang diteruskan ke pekerjaan downstream. | String | Tidak | None | Bidang metadata yang tersedia meliputi | ||||||||||
decode.table-id.fields | Menentukan bidang yang nilainya digunakan untuk menghasilkan Table ID saat mengurai data log dari Log Service (SLS). | String | Tidak | None | Beberapa bidang dipisahkan dengan 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 | None | 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 parsing 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 SLS Catalog bawaan yang dibuat di halaman Data Management dalam pekerjaan data ingestion Flink CDC. Hal ini menghindari kebutuhan untuk menentukan properti koneksi secara manual.
source:
type: sls
using.built-in-catalog: sls_catalogSaat ini, pekerjaan data ingestion dapat secara otomatis menggunakan kembali parameter berikut dari SLS Catalog 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. Konektor mengurai skema setiap entri log dan 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 untuk mengurai skema log.
Informasi primary key
Log Log Service (SLS) tidak berisi informasi primary key. Anda dapat menambahkan primary key secara manual ke tabel dengan menggunakan aturan transformasi:
transform: - source-table: <project>.<logstore> projection: * primary-keys: key1, key2Inferensi 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, dan membandingkannya dengan skema saat ini. Jika skema yang diinferensi tidak konsisten dengan skema saat ini, skema digabungkan sesuai aturan berikut:Jika kolom fisik yang diinferensi berisi bidang yang tidak ada dalam 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 dalam 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 di 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 primary key ke tabel.
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
# Arahkan 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: trueDataStream API
Untuk membaca atau menulis data dengan DataStream API, gunakan konektor DataStream. Untuk informasi lebih lanjut, lihat Usage of DataStream connectors.
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 {
// Sets up the streaming execution environment
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// Creates an SLS source and prints the data to the console.
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");
// The batch get size is required.
accessInfo.setBatchGetSize(10);
// Optional parameters
accessInfo.setConsumerGroup("yourConsumerGroup");
accessInfo.setMaxRetries(3);
// Start time for consumption, set to the current time.
int startInSec = (int) (new Date().getTime() / 1000);
// Stop time for consumption, where -1 means never stop.
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>