All Products
Search
Document Center

Realtime Compute for Apache Flink:Log Service (SLS)

Last Updated:Apr 25, 2026

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 "__tag__:__receive_time__":"1616742274", '__receive_time__' dan '1616742274' disimpan sebagai pasangan kunci-nilai dalam map. Untuk mengakses nilai tersebut dalam SQL, gunakan __tag__['__receive_time__'].

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 parallelism lebih 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 fitur failover otomatis 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.

    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 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.

    Penting

    Untuk 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.

    Penting
    • Opsi ini hanya didukung di VVR 8.0.9 dan versi yang lebih baru. Nilai default adalah true untuk 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 consumerGroup untuk mencatat progres konsumsi dalam kelompok konsumen SLS. Kemudian, atur opsi consumeFromCheckpoint ke true dan 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.

    Catatan
    • Nilai harus lebih besar dari atau sama dengan 60.000 milidetik (1 menit).

    • Opsi ini hanya berlaku ketika enableNewSource diatur ke true.

    • 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.

    Penting
    • Untuk versi VVR sebelum 11.1, nilai consumer_group tidak didukung. Anda harus mengatur consumeFromCheckpoint ke true. 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 startupMode diatur ke timestamp.

    Catatan

    Opsi startTime dan stopTime didasarkan 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.

    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 keluar setelah semua log dikonsumsi, Anda juga harus mengatur exitAfterFinish ke true.

    consumerGroup

    Nama kelompok konsumen.

    String

    Tidak

    None

    Kelompok konsumen digunakan untuk mencatat progres konsumsi. Anda dapat menentukan nama kustom tanpa format tetap.

    Catatan

    Pekerjaan 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.

    Penting

    Parameter ini tidak lagi didukung di VVR 11.1 dan versi yang lebih baru. Untuk versi tersebut, Anda harus mengatur opsi startupMode ke consumer_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 batchGetSize tidak 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

    Penting

    Opsi 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 bidang request_method adalah 'GET'.

    Catatan

    Opsi ini menggunakan bahasa SPL Log Service (SLS). Untuk informasi lebih lanjut, lihat SPL syntax.

    Penting
    • Opsi 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 query dan processor keduanya ditentukan, query memiliki prioritas lebih tinggi dan processor diabaikan.

    String

    Tidak

    None

    Opsi ini memfilter data SLS sebelum Flink mengonsumsinya, yang mengurangi biaya dan meningkatkan kecepatan pemrosesan. Kami menyarankan agar Anda menggunakan processor daripada query.

    Misalnya, 'processor' = 'test-filter-processor' menunjukkan bahwa prosesor SLS memfilter data sebelum Flink membaca data dari SLS.

    Catatan

    Opsi 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.

    Penting

    Opsi 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 INT yang 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 partitionField ditentukan, 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.

    Catatan

    Opsi 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 sls.

endpoint

Titik akhir Log Service (SLS).

String

Ya

None

Alamat akses VPC Log Service (SLS). Untuk informasi lebih lanjut, lihat Service Endpoints.

Catatan
  • Secara default, Realtime Compute for Apache Flink tidak dapat mengakses internet. Anda dapat menggunakan NAT Gateway untuk mengaktifkan komunikasi antara Virtual Private Cloud (VPC) Anda dan internet. Untuk informasi lebih lanjut, lihat How can I access the internet?.

  • Kami tidak menyarankan mengakses Log Service (SLS) melalui internet. Jika Anda harus melakukannya, gunakan HTTPS dan aktifkan transfer acceleration untuk SLS.

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

  • continuous: Melakukan inferensi skema untuk setiap catatan data. Jika skema tidak kompatibel, skema yang lebih luas diinferensi dan event perubahan skema dihasilkan.

  • static: Melakukan inferensi skema hanya sekali saat pekerjaan dimulai. Data selanjutnya diparse berdasarkan skema awal, dan tidak ada event perubahan skema yang dihasilkan.

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

  • timestamp (default): Mengonsumsi log dari timestamp tertentu.

  • latest: Mengonsumsi log dari offset terbaru.

  • earliest: Mengonsumsi log dari offset paling awal.

  • consumer_group: Mengonsumsi log dari offset yang dicatat dalam kelompok konsumen. Jika tidak ada offset yang dicatat untuk suatu shard, konsumsi dimulai dari offset paling awal.

startTime

Waktu mulai untuk konsumsi log.

String

Tidak

Waktu saat ini

Formatnya adalah yyyy-MM-dd HH:mm:ss.

Parameter ini hanya berlaku ketika startupMode diatur ke timestamp.

Catatan

Parameter startTime dan stopTime didasarkan pada atribut __receive_time__ di Log Service (SLS), bukan atribut __timestamp__.

stopTime

Waktu akhir untuk konsumsi log.

String

Tidak

None

Formatnya adalah yyyy-MM-dd HH:mm:ss.

Catatan

Jika Anda ingin pekerjaan Flink keluar setelah semua log dikonsumsi, Anda juga harus mengatur exitAfterFinish=true.

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

  • true: Pekerjaan Flink keluar setelah semua data dikonsumsi.

  • false (default): Pekerjaan Flink tidak keluar setelah semua data dikonsumsi.

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, 'query' = '*| where request_method = ''GET''' memfilter data di mana bidang request_method bernilai 'GET' sebelum data dibaca oleh Flink.

Catatan

Kueri harus menggunakan sintaks SPL Log Service. Untuk informasi lebih lanjut, lihat SPL syntax.

Penting
  • Untuk informasi tentang wilayah tempat fitur ini tersedia di Log Service (SLS), lihat Consume logs based on rules.

  • Fitur ini sedang dalam pratinjau publik dan gratis. Anda mungkin dikenai biaya untuk fitur ini di masa mendatang. Untuk informasi lebih lanjut, lihat Pricing.

compressType

Jenis kompresi untuk Log Service (SLS).

String

Tidak

None

Jenis kompresi yang didukung meliputi:

  • lz4

  • deflate

  • zstd

timeZone

Zona waktu untuk startTime dan stopTime.

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 __source__, __topic__, __timestamp__, dan __tag__. Pisahkan beberapa bidang dengan koma.

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 ,. Misalnya, jika catatan log SLS upstream adalah {"col0":"a", "col1":"b", "col2":"c"}, hasil untuk konfigurasi parameter yang berbeda adalah sebagai berikut:

Konfigurasi

Table ID

None

Semua pesan adalah Project.Logstore

col0

a

col0,col1

a.b

col0,col1,col2

a.b.c

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 , untuk memisahkan beberapa definisi bidang. Misalnya, id BIGINT, name VARCHAR(10) menentukan tipe bidang id sebagai BIGINT dan tipe bidang name sebagai VARCHAR(10).

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:

  • SQL: Mengurai timestamp input dalam format yyyy-MM-dd HH:mm:ss.s{precision} (misalnya, 2020-12-30 12:13:14.123) dan mengeluarkannya dalam format yang sama.

  • ISO-8601: Mengurai timestamp input dalam format yyyy-MM-ddTHH:mm:ss.s{precision} (misalnya, 2020-12-30T12:13:14.123) dan mengeluarkannya dalam format yang sama.

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 ingestion.ignore-errors diaktifkan, pekerjaan gagal ketika jumlah kumulatif error melebihi nilai ini.

Integer

Tidak

-1

Parameter ini hanya berlaku ketika ingestion.ignore-errors diaktifkan. Nilai default -1 berarti pekerjaan mengabaikan semua exception parsing.

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_catalog

Saat 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 maxPreFetchLogGroups kelompok 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.

    Catatan

    Untuk 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, 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 ke continuous, 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: true

DataStream API

Penting

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>

FAQ

Bagaimana cara mengatasi TaskManager OOM (java.lang.OutOfMemoryError: Java heap space) saat memulihkan program Flink yang gagal?