All Products
Search
Document Center

Realtime Compute for Apache Flink:Log Service (SLS)

Last Updated:Aug 14, 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

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 "__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 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 parallelism lebih 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 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

    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.

    Penting

    Untuk 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 true untuk 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 consumerGroup untuk mencatat progres konsumsi di 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 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 enableNewSource diatur ke true.

    • 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 consumeFromCheckpoint ke true. 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 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

    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 exitAfterFinish ke true.

    consumerGroup

    Nama kelompok konsumen.

    String

    Tidak

    Tidak ada

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

    Penting

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

    Penting

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

    Catatan

    Opsi 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 query dan processor keduanya ditentukan, query yang diutamakan dan processor diabaikan.

    String

    Tidak

    Tidak ada

    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 selengkapnya, lihat Sintaksis SPL. Untuk informasi tentang cara membuat atau memperbarui prosesor SLS, lihat Kelola prosesor.

    Penting

    Opsi 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 enableNewSource diatur ke true.

    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 enableNewSource diatur ke true.

  • 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 INT yang 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 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 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 sls.

endpoint

Titik akhir Log Service (SLS).

String

Ya

Tidak ada

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. Anda dapat menggunakan NAT Gateway untuk mengaktifkan komunikasi antara Virtual Private Cloud (VPC) Anda dan internet. Untuk informasi selengkapnya, lihat Bagaimana cara mengakses internet?.

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

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

  • 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

Tidak ada

  • 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 direkam di kelompok konsumen. Jika tidak ada offset yang direkam 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

Tidak ada

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

Catatan

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

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

  • true: Pekerjaan Flink berhenti setelah semua data dikonsumsi.

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

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

Catatan

Kueri harus menggunakan sintaksis SPL Log Service. Untuk informasi selengkapnya, lihat Sintaksis SPL.

Penting
  • Untuk informasi tentang wilayah tempat fitur ini tersedia di Log Service (SLS), lihat Konsumsi log berdasarkan aturan.

  • Fitur ini dalam status Beta dan gratis. Anda mungkin dikenai biaya untuk fitur ini di masa mendatang. Untuk informasi selengkapnya, lihat Harga.

compressType

Tipe kompresi untuk Log Service (SLS).

String

Tidak

Tidak ada

Tipe kompresi yang didukung meliputi:

  • lz4

  • deflate

  • zstd

timeZone

Zona waktu untuk startTime dan stopTime.

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

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

Konfigurasi

ID Tabel

Tidak ada

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

Tidak ada

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

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

    Catatan

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

Penting

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>

FAQ

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