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: loggerAsumsikan 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: loggerUntuk 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: loggerMembaca 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: loggerMenangani 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: loggerJika 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: loggerInferensi 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: loggerKonfigurasi 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: loggerMengurai skema dari
mysqlTypesource: 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: loggerAnda 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: loggerMempercepat 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.
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, ataucanal-json.infer-schema.primitive-as-string. Ini melewatkan inferensi tipe dan mengatur semua tipe bidang menjadi String, sehingga mempercepat parsing.Untuk data Canal JSON, Anda dapat menggunakan
canal-json.database.includedancanal-json.table.includeuntuk memfilter data dari tabel yang tidak diperlukan.Jika skema data tidak akan berubah atau jika Anda tidak memerlukan evolusi skema, Anda dapat mengubah
schema.inference.strategymenjadistatic. 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: loggerTableID: Secara default,
projectdanlogstoredigabungkan untuk membentuk TableID. Dalam pekerjaan di atas, TableID-nya adalahtest_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: loggerMenangani 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: loggerMenentukan 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: loggerMenentukan 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