All Products
Search
Document Center

Realtime Compute for Apache Flink:Ingesti log Real-time ke data lake

Last Updated:Apr 25, 2026

Topik ini menjelaskan praktik terbaik untuk menulis log data ke Alibaba Cloud Data Lake Formation (DLF) menggunakan pekerjaan YAML ingesti data Flink CDC.

Ingesti log Kafka ke data lake

Gunakan ingesti data Flink CDC dan pekerjaan YAML sederhana untuk mengingesti log data ke data lake secara real-time. Sistem secara otomatis melakukan inferensi skema dan mendukung evolusi skema.

Asumsikan topik inventory di Apache Kafka menyimpan data untuk tabel log dalam format JSON. Contoh pekerjaan berikut menyinkronkan data ini ke tabel sink yang sesuai di Data Lake Formation (DLF):

source:
  type: kafka
  name: Kafka Source
  # Alamat broker Kafka.
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  # Topik yang dikonsumsi.
  topic: inventory
  # Menentukan bahwa konsumsi dimulai dari offset paling awal.
  scan.startup.mode: earliest-offset
  # Format nilai pesan Kafka.
  value.format: json
  # Mengatur strategi inferensi skema ke continuous. Strategi ini mendeteksi skema setiap pesan dan menyinkronkan perubahan skema.
  schema.inference.strategy: continuous
  # (Opsional) Meratakan kolom bersarang secara rekursif dalam data JSON.
  json.infer-schema.flatten-nested-columns.enable: true
  # (Opsional) Melewatkan 100 pengecualian parsing pertama. Jika jumlah pengecualian melebihi 100, pekerjaan gagal.
  ingestion.ignore-errors: true
  ingestion.error-tolerance.max-count: 100

sink:
  type: paimon
  # Jenis metastore. Nilainya tetap rest.
  catalog.properties.metastore: rest
  # Penyedia token. Nilainya tetap dlf.
  catalog.properties.token.provider: dlf
  # URI untuk mengakses DLF Rest Catalog Server. Formatnya adalah http://[region-id]-vpc.dlf.aliyuncs.com. Misalnya, http://cn-hangzhou-vpc.dlf.aliyuncs.com.
  catalog.properties.uri: dlf_uri
  # Nama katalog DLF.
  catalog.properties.warehouse: your_warehouse
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true

# Menambahkan informasi primary key ke tabel. 
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id
    
# Menulis semua data dari topik inventory ke tabel test_database.inventory.
route:
  - source-table: inventory
    sink-table: test_database.inventory
    
pipeline:
  # (Opsional) Mencatat data kotor yang menyebabkan pengecualian pemrosesan dalam log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Asumsikan topik inventory di Apache Kafka menyimpan data untuk beberapa tabel log dalam format JSON. Bidang databaseName dan tableName dalam muatan JSON menyediakan nama database dan tabel. Contoh pekerjaan berikut menyinkronkan data dari tabel-tabel tersebut ke tabel sink yang sesuai di DLF:

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  topic: inventory
  scan.startup.mode: earliest-offset
  value.format: json
  # (Opsional) Meratakan kolom bersarang secara rekursif dalam data JSON.
  json.infer-schema.flatten-nested-columns.enable: true
  # Menggunakan nilai bidang databaseName sebagai nama database dan nilai bidang tableName sebagai nama tabel.
  json.decode.parser-table-id.fields: databaseName,tableName
  # (Opsional) Melewatkan 100 pengecualian parsing pertama. Jika jumlah pengecualian melebihi 100, pekerjaan gagal.
  ingestion.ignore-errors: true
  ingestion.error-tolerance.max-count: 100

sink:
  type: paimon
  # Jenis metastore. Nilainya tetap rest.
  catalog.properties.metastore: rest
  # Penyedia token. Nilainya tetap dlf.
  catalog.properties.token.provider: dlf
  # URI untuk mengakses DLF Rest Catalog Server. Formatnya adalah http://[region-id]-vpc.dlf.aliyuncs.com. Misalnya, http://cn-hangzhou-vpc.dlf.aliyuncs.com.
  catalog.properties.uri: dlf_uri
  # Nama katalog DLF.
  catalog.properties.warehouse: your_warehouse
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true

# Menambahkan informasi primary key ke tabel. 
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id
    
# Menulis data dari ods.inventory, ods.customer, dan ods.user ke tabel test_database.inventory, test_database.customer, dan test_database.user, masing-masing.
route:
  - source-table: ods.inventory
    sink-table: test_database.inventory
  - source-table: ods.customer
    sink-table: test_database.customer
  - source-table: ods.user
    sink-table: test_database.user   
    
pipeline:
  # (Opsional) Mencatat data kotor yang menyebabkan pengecualian pemrosesan dalam log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Untuk detail kebijakan inferensi dan evolusi skema untuk tabel sumber Kafka berformat JSON, lihat Kebijakan Parsing dan Evolusi Skema.

Kasus penggunaan

Bagian berikut menjelaskan konfigurasi pekerjaan untuk kasus penggunaan umum. Untuk konfigurasi lebih rinci, lihat Konektor Kafka Ingesti Data.

Menyelesaikan konflik nama bidang

Pesan Apache Kafka terdiri dari Key dan Value. Anda dapat menentukan format untuk masing-masing bagian dengan mengatur key.format dan value.format. Skema akhir merupakan gabungan semua bidang dari kedua bagian tersebut.

Jika terdapat konflik nama bidang antara bagian Key dan Value, gunakan key.fields-prefix dan value.fields-prefix untuk menambahkan awalan pada nama bidang.

Misalnya, jika bagian Key pesan Kafka berisi bidang id dan name, dan bagian Value berisi bidang id dan price, konfigurasi pekerjaan berikut menghasilkan skema dengan bidang key_id, key_name, val_id, dan val_price.

source:
  type: kafka
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  key.format: json
  # Menambahkan awalan pada nama bidang di Key.
  key.fields-prefix: key_
  # Menambahkan awalan pada nama bidang di Value.
  value.fields-prefix: val_
  # Mengatur strategi inferensi skema ke continuous. Strategi ini mendeteksi skema setiap pesan dan menyinkronkan perubahan skema.
  schema.inference.strategy: continuous
  # Mengabaikan error selama parsing data.
  ingestion.ignore-errors: true
  # Jika terjadi 100 error parsing data, pekerjaan gagal.
  ingestion.error-tolerance.max-count: 100

# Menulis semua data dari topik test_topic ke tabel test_database.test_topic.
route:
  - source-table: test_topic
    sink-table: test_database.test_topic

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Mengaktifkan pengumpulan data kotor. Data kotor ditulis ke file log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Membaca metadata

Gunakan konfigurasi metadata.list untuk membaca dan meneruskan metadata tambahan pesan Kafka. Kolom metadata yang ditambahkan di metadata.list dapat digunakan langsung di modul transform. Untuk daftar metadata yang didukung, lihat Kolom metadata yang tersedia.

Konfigurasi berikut menambahkan metadata partition dan offset ke data. Kemudian menggunakan metadata ini di modul transform untuk memfilter data di mana partisi lebih besar dari 1 dan offset lebih besar dari 100.

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Menambahkan kolom metadata.
  metadata.list: partition,offset
  # Mengatur strategi inferensi skema ke continuous. Strategi ini mendeteksi skema setiap pesan dan menyinkronkan perubahan skema.
  schema.inference.strategy: continuous
  # Mengabaikan error selama parsing data.
  ingestion.ignore-errors: true
  # Jika terjadi 100 error parsing data, pekerjaan gagal.
  ingestion.error-tolerance.max-count: 100
  
transform:
  - source-table: \.*.\.*
    filter: '`partition` > 1 and `offset` > 100'
    
# Menulis semua data dari topik test_topic ke tabel test_database.test_topic.
route:
  - source-table: test_topic
    sink-table: test_database.test_topic

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Mengaktifkan pengumpulan data kotor. Data kotor ditulis ke file log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Menangani data kotor

Data log mungkin berisi data kotor dengan format salah, yang dapat menyebabkan pekerjaan gagal dan restart berulang kali. Ingesti data Flink CDC mendukung pengabaian error parsing dan pengumpulan data yang gagal diparse. Untuk informasi lebih lanjut, lihat Pengumpulan Data Kotor.

Pekerjaan berikut menulis data kotor yang gagal diparse ke file log. Pekerjaan gagal jika terjadi lebih dari 100 error parsing.

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Mengatur strategi inferensi skema ke continuous. Strategi ini mendeteksi skema setiap pesan dan menyinkronkan perubahan skema.
  schema.inference.strategy: continuous
  # Mengabaikan error selama parsing data.
  ingestion.ignore-errors: true
  # Jika terjadi 100 error parsing data, pekerjaan gagal.
  ingestion.error-tolerance.max-count: 100
  
# Menulis semua data dari topik test_topic ke tabel test_database.test_topic.
route:
  - source-table: test_topic
    sink-table: test_database.test_topic

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Mengaktifkan pengumpulan data kotor. Data kotor ditulis ke file log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Jika Anda tidak perlu mengumpulkan data kotor dan ingin mencegah kegagalan pekerjaan akibat data kotor, Anda dapat menggunakan json.ignore-parse-errors, debezium-json.ignore-parse-errors, atau canal-json.ignore-parse-errors untuk langsung mengabaikan error parsing.

Menentukan TableID

Format Debezium JSON dan Canal JSON bersifat tetap, dengan bidang tertentu yang digunakan untuk menyimpan TableID. Untuk data JSON generik yang tidak memiliki format tetap, nama topik berfungsi sebagai TableID secara default. Untuk menggunakan bidang tertentu dari data sebagai TableID, konfigurasikan json.decode.parser-table-id.fields. Misalnya, diberikan data JSON {"col0":"a", "col1":"b", "col2":"c"}, konfigurasi berbeda menghasilkan TableID berikut:

Konfigurasi

TableID

col0

a

col0,col1

a.b

col0,col1,col2

a.b.c

Konfigurasi pekerjaan berikut menggabungkan kolom db dan table dari data untuk membuat TableID.

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Menggunakan bidang db dan table dari data sebagai TableID.
  json.decode.parser-table-id.fields: db,table
  # Mengatur strategi inferensi skema ke continuous. Strategi ini mendeteksi skema setiap pesan dan menyinkronkan perubahan skema.
  schema.inference.strategy: continuous
  # Mengabaikan error selama parsing data.
  ingestion.ignore-errors: true
  # Jika terjadi 100 error parsing data, pekerjaan gagal.
  ingestion.error-tolerance.max-count: 100

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Mengaktifkan pengumpulan data kotor. Data kotor ditulis ke file log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Inferensi tipe data

Tipe data untuk log data diinferensi dengan mengurai data. Untuk informasi rinci tentang inferensi tipe, lihat Kebijakan Parsing dan Evolusi Skema. Bagian berikut menjelaskan cara menyesuaikan konfigurasi untuk mengontrol parsing tipe dalam skenario umum.

(Umum) Mengatur semua tipe bidang menjadi string

Jika pemrosesan dan penyimpanan downstream tidak memerlukan tipe bidang tertentu, Anda dapat mengaktifkan json.infer-schema.primitive-as-string, debezium-json.infer-schema.primitive-as-string, atau canal-json.infer-schema.primitive-as-string. Ini melewatkan inferensi tipe dan mengatur semua tipe bidang menjadi String.

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Mengatur tipe semua bidang dalam data JSON menjadi String.
  json.infer-schema.primitive-as-string: true
  # Mengatur strategi inferensi skema ke continuous. Strategi ini mendeteksi skema setiap pesan dan menyinkronkan perubahan skema.
  schema.inference.strategy: continuous
  # Mengabaikan error selama parsing data.
  ingestion.ignore-errors: true
  # Jika terjadi 100 error parsing data, pekerjaan gagal.
  ingestion.error-tolerance.max-count: 100

# Menulis semua data dari topik test_topic ke tabel test_database.test_topic.
route:
  - source-table: test_topic
    sink-table: test_database.test_topic

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Mengaktifkan pengumpulan data kotor. Data kotor ditulis ke file log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

(Umum) Menentukan skema awal

Dalam beberapa skenario, Anda mungkin perlu menentukan skema awal, misalnya saat menulis data Kafka ke tabel downstream yang sudah ada. Anda dapat menentukan skema awal dengan menambahkan parameter scan.value.initial-schemas.ddls.

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Menggunakan bidang db dan table dari data sebagai TableID.
  json.decode.parser-table-id.fields: db,table
  # Menentukan skema awal.
  scan.value.initial-schemas.ddls: |
    CREATE TABLE db1.t1 (id BIGINT, name VARCHAR(10));
    CREATE TABLE db1.t2 (id BIGINT);
  # Mengatur strategi inferensi skema ke continuous. Strategi ini mendeteksi skema setiap pesan dan menyinkronkan perubahan skema.
  schema.inference.strategy: continuous
  # Mengabaikan error selama parsing data.
  ingestion.ignore-errors: true
  # Jika terjadi 100 error parsing data, pekerjaan gagal.
  ingestion.error-tolerance.max-count: 100
  
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Mengaktifkan pengumpulan data kotor. Data kotor ditulis ke file log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Konfigurasi di atas menentukan tipe awal BIGINT untuk bidang id dan VARCHAR(10) untuk bidang name di tabel db1.t1. Konfigurasi ini juga menentukan tipe awal BIGINT untuk bidang id di tabel db1.t2.

(JSON) Menentukan tipe bidang tetap

Untuk data JSON, tipe bidang diinferensi dengan mengurai tipe node JSON. Terkadang tipe yang diinferensi tidak sesuai harapan Anda. Untuk mengatasi hal ini, gunakan konfigurasi json.infer-schema.fixed-types untuk menentukan tipe bidang tertentu.

Konfigurasi pekerjaan berikut menentukan tipe bidang id sebagai BIGINT dan tipe bidang name sebagai VARCHAR(10).

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Menentukan tipe tetap untuk bidang tertentu.
  json.infer-schema.fixed-types: id BIGINT, name VARCHAR(10)
  # Diperlukan untuk versi 11.5 dan sebelumnya.
  scan.max.pre.fetch.records: 0
  # Mengatur strategi inferensi skema ke continuous. Strategi ini mendeteksi skema setiap pesan dan menyinkronkan perubahan skema.
  schema.inference.strategy: continuous
  # Mengabaikan error selama parsing data.
  ingestion.ignore-errors: true
  # Jika terjadi 100 error parsing data, pekerjaan gagal.
  ingestion.error-tolerance.max-count: 100

# Menulis semua data dari topik test_topic ke tabel test_database.test_topic.
route:
  - source-table: test_topic
    sink-table: test_database.test_topic
    
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Mengaktifkan pengumpulan data kotor. Data kotor ditulis ke file log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

(Canal JSON) Menentukan sumber inferensi

Data Canal JSON berisi informasi lebih banyak daripada data JSON standar. Jika bidang sqlType atau mysqlType ada dalam data Canal JSON, Anda dapat menggunakan informasi ini untuk mengurai tipe data yang lebih tepat.

  • Mengurai skema dari sqlType

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: canal-json
  # Menginferensi skema dari bidang sqlType.
  canal-json.infer-schema.strategy: SQL_TYPE
  # Mengatur strategi inferensi skema ke continuous. Strategi ini mendeteksi skema setiap pesan dan menyinkronkan perubahan skema.
  schema.inference.strategy: continuous
  # Mengabaikan error selama parsing data.
  ingestion.ignore-errors: true
  # Jika terjadi 100 error parsing data, pekerjaan gagal.
  ingestion.error-tolerance.max-count: 100
  
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Mengaktifkan pengumpulan data kotor. Data kotor ditulis ke file log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger
  • Mengurai skema dari mysqlType

    source:
      type: kafka
      name: Kafka source
      properties.bootstrap.servers: localhost:9092
      topic: test_topic
      properties.group.id: test_group
      scan.startup.mode: earliest-offset
      value.format: canal-json
      # Menginferensi skema dari bidang mysqlType.
      canal-json.infer-schema.strategy: MYSQL_TYPE
      # Mengatur strategi inferensi skema ke continuous. Strategi ini mendeteksi skema setiap pesan dan menyinkronkan perubahan skema.
      schema.inference.strategy: continuous
      # Mengabaikan error selama parsing data.
      ingestion.ignore-errors: true
      # Jika terjadi 100 error parsing data, pekerjaan gagal.
      ingestion.error-tolerance.max-count: 100
      
    sink:
      type: paimon
      catalog.properties.metastore: rest
      catalog.properties.uri: dlf_uri
      catalog.properties.warehouse: your_warehouse
      catalog.properties.token.provider: dlf
      # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
      table.properties.deletion-vectors.enabled: true
      
    pipeline:
      # Mengaktifkan pengumpulan data kotor. Data kotor ditulis ke file log.
      dirty-data.collector:
        name: Logger Dirty Data Collector
        type: logger

Menggunakan skema statis

Jika skema data di topik Anda bersifat tetap, atur schema.inference.strategy ke static. Pekerjaan ingesti data kemudian hanya melakukan inferensi skema sekali saat startup dan tidak mengurai skema data berikutnya.

source:
  type: kafka
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Mengatur strategi inferensi skema ke static. Skema hanya diinferensi sekali saat pekerjaan dimulai.
  schema.inference.strategy: static
  # Mencoba mengonsumsi 20 catatan dari setiap partisi untuk menginferensi skema.
  scan.max.pre.fetch.records: 20
  # Mengabaikan error selama parsing data.
  ingestion.ignore-errors: true
  # Jika terjadi 100 error parsing data, pekerjaan gagal.
  ingestion.error-tolerance.max-count: 100
  
# Menulis semua data dari topik test_topic ke tabel test_database.test_topic.
route:
  - source-table: test_topic
    sink-table: test_database.test_topic
    
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Mengaktifkan pengumpulan data kotor. Data kotor ditulis ke file log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Anda juga dapat menambahkan parameter scan.value.initial-schemas.ddls untuk menentukan skema awal dan melewatkan inferensi skema untuk tabel tertentu. Contoh berikut menentukan skema awal untuk tabel db1.t1 dan db1.t2.

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Menggunakan bidang db dan table dari data sebagai TableID.
  json.decode.parser-table-id.fields: db,table
  # Menentukan skema awal.
  scan.value.initial-schemas.ddls: | 
    CREATE TABLE db1.t1 (id BIGINT, name VARCHAR(10));
    CREATE TABLE db1.t2 (id BIGINT);
  # Mengatur strategi inferensi skema ke static. Skema hanya diinferensi sekali saat pekerjaan dimulai.
  schema.inference.strategy: static
  # Mengabaikan error selama parsing data.
  ingestion.ignore-errors: true
  # Jika terjadi 100 error parsing data, pekerjaan gagal.
  ingestion.error-tolerance.max-count: 100
  
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Mengaktifkan pengumpulan data kotor. Data kotor ditulis ke file log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Mempercepat parsing data log Kafka

Selain konfigurasi akselerasi konektor Kafka umum, pekerjaan ingesti data Flink CDC menawarkan pengaturan khusus untuk mempercepat parsing berdasarkan kasus penggunaan Anda.

  1. Parsing skema bisa memakan waktu. Jika aplikasi downstream Anda hanya memerlukan tipe String, Anda dapat mengaktifkan json.infer-schema.primitive-as-string, debezium-json.infer-schema.primitive-as-string, atau canal-json.infer-schema.primitive-as-string. Ini melewatkan inferensi tipe dan mengatur semua tipe bidang menjadi String, sehingga mempercepat parsing.

  2. Untuk data Canal JSON, Anda dapat menggunakan canal-json.database.include dan canal-json.table.include untuk memfilter data dari tabel yang tidak diperlukan.

  3. Jika skema data tidak akan berubah atau jika Anda tidak memerlukan evolusi skema, Anda dapat mengubah schema.inference.strategy menjadi static. Ini hanya melakukan inferensi skema sekali saat pekerjaan dimulai.

Ingesti log SLS ke data lake

Log Service (SLS) adalah layanan satu atap untuk data log. Ingesti data Flink CDC memungkinkan Anda mengingesti data log dari SLS ke data lake secara real-time. Layanan ini secara otomatis melakukan inferensi skema dan mendukung evolusi skema.

source:
  type: sls
  endpoint: localhost
  project: test_pj
  logstore: test_log
  accessId: access_id
  accessKey: access_key
  
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Mengaktifkan pengumpulan data kotor. Data kotor ditulis ke file log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger
  • TableID: Secara default, project dan logstore digabungkan untuk membentuk TableID. Dalam pekerjaan di atas, TableID-nya adalah test_pj.test_log.

  • Tipe data: Konektor SLS secara default memperlakukan semua bidang dalam setiap entri log sebagai tipe String.

  • Evolusi skema: Evolusi skema saat ini terbatas pada penambahan kolom baru. Kolom baru ditambahkan di akhir skema secara default.

Kasus penggunaan

Bagian berikut menjelaskan konfigurasi pekerjaan untuk kasus penggunaan umum. Untuk konfigurasi lebih rinci, lihat Konektor SLS Ingesti Data.

Membaca metadata

Anda dapat menggunakan konfigurasi metadata.list untuk membaca dan meneruskan metadata SLS tambahan. Kolom metadata yang ditambahkan di metadata.list dapat digunakan langsung di modul transform. Untuk daftar metadata yang didukung, lihat metadata.list.

Perhatikan bahwa kolom yang ditambahkan dengan metadata.list tidak secara otomatis disertakan dalam data keluaran. Untuk menulis kolom metadata ini ke sink, Anda harus mendeklarasikannya di bagian projection modul transform.

Konfigurasi berikut menambahkan metadata __timestamp__ dan __tag__ ke data dan menuliskannya ke sink. Konfigurasi ini juga menggunakan bidang metadata tersebut di modul transform untuk memfilter data di mana timestamp lebih besar dari 1772181154 dan tag-nya adalah "test".

source:
  type: sls
  endpoint: localhost
  project: test_pj
  logstore: test_log
  accessId: access_id
  accessKey: access_key
  # Menambahkan kolom metadata.
  metadata.list: __timestamp__,__tag__
  # Mengatur strategi inferensi skema ke continuous. Strategi ini mendeteksi skema setiap pesan dan menyinkronkan perubahan skema.
  schema.inference.strategy: continuous
  # Mengabaikan error selama parsing data.
  ingestion.ignore-errors: true
  # Jika terjadi 100 error parsing data, pekerjaan gagal.
  ingestion.error-tolerance.max-count: 100
  
transform:
  - source-table: \.*.\.*
    projection: \*, __timestamp__ as timestamp_col, __tag__ as tag_col
    filter: '`__timestamp__` > 1772181154 and `__tag__` = "test"'

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Mengaktifkan pengumpulan data kotor. Data kotor ditulis ke file log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Menangani data kotor

Data log mungkin berisi data kotor dengan format salah, yang dapat menyebabkan pekerjaan gagal dan restart berulang kali. Ingesti data Flink CDC mendukung pengabaian error parsing dan pengumpulan data yang menyebabkan error tersebut. Untuk informasi lebih lanjut, lihat Pengumpulan Data Kotor.

Pekerjaan berikut menulis data kotor yang gagal diparse ke file log. Pekerjaan gagal jika terjadi lebih dari 100 error parsing.

source:
  type: sls
  endpoint: localhost
  project: test_pj
  logstore: test_log
  accessId: access_id
  accessKey: access_key
  # Mengatur strategi inferensi skema ke continuous. Strategi ini mendeteksi skema setiap pesan dan menyinkronkan perubahan skema.
  schema.inference.strategy: continuous
  # Mengabaikan error selama parsing data.
  ingestion.ignore-errors: true
  # Jika terjadi 100 error parsing data, pekerjaan gagal.
  ingestion.error-tolerance.max-count: 100

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Mengaktifkan pengumpulan data kotor. Data kotor ditulis ke file log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Menentukan TableID

Untuk menggunakan bidang tertentu dari data sebagai TableID, konfigurasikan decode.table-id.fields. Misalnya, diberikan data log {"col0":"a", "col1":"b", "col2":"c"}, konfigurasi berbeda menghasilkan TableID berikut:

Konfigurasi

TableID

col0

a

col0,col1

a.b

col0,col1,col2

a.b.c

Konfigurasi pekerjaan berikut menggabungkan kolom db dan table dari data untuk membuat TableID.

source:
  type: sls
  endpoint: localhost
  project: test_pj
  logstore: test_log
  accessId: access_id
  accessKey: access_key
  # Menggunakan bidang db dan table dari data sebagai TableID.
  decode.table-id.fields: db,table
  # Mengatur strategi inferensi skema ke continuous. Strategi ini mendeteksi skema setiap pesan dan menyinkronkan perubahan skema.
  schema.inference.strategy: continuous
  # Mengabaikan error selama parsing data.
  ingestion.ignore-errors: true
  # Jika terjadi 100 error parsing data, pekerjaan gagal.
  ingestion.error-tolerance.max-count: 100

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Mengaktifkan pengumpulan data kotor. Data kotor ditulis ke file log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Menentukan tipe bidang

Konektor ingesti data SLS secara default memperlakukan semua bidang dalam data sebagai tipe String. Jika Anda perlu menentukan tipe untuk bidang tertentu, gunakan opsi konfigurasi fixed-types.

Pekerjaan berikut menentukan tipe bidang id sebagai BIGINT dan tipe bidang name sebagai VARCHAR(10).

source:
  type: sls
  endpoint: localhost
  project: test_pj
  logstore: test_log
  accessId: access_id
  accessKey: access_key
  # Menentukan tipe bidang id sebagai BIGINT dan tipe bidang name sebagai VARCHAR(10).
  fixed-types: id BIGINT, name VARCHAR(10)
  # Mengatur strategi inferensi skema ke continuous. Strategi ini mendeteksi skema setiap pesan dan menyinkronkan perubahan skema.
  schema.inference.strategy: continuous
  # Mengabaikan error selama parsing data.
  ingestion.ignore-errors: true
  # Jika terjadi 100 error parsing data, pekerjaan gagal.
  ingestion.error-tolerance.max-count: 100

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Opsional) Mengaktifkan deletion vectors untuk meningkatkan performa baca.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Mengaktifkan pengumpulan data kotor. Data kotor ditulis ke file log.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger