Masalah umum pada konektor dan solusinya untuk Realtime Compute for Apache Flink.
-
Kafka
-
DataHub
-
MaxCompute
-
Pembacaan data untuk sumber MaxCompute penuh dan inkremental
-
Dapatkah sumber MaxCompute membaca data yang ditambahkan setelah job dimulai?
-
Mengubah konkurensi untuk job sumber MaxCompute yang dilanjutkan
-
Mengapa sumber MaxCompute membaca partisi sebelum posisi awal?
-
Menangani partisi baru yang belum lengkap dalam sumber MaxCompute inkremental
-
Error konektor MaxCompute: ErrorMessage=Authorization Failed [4019], You have NO privilege
-
Mengonfigurasi parameter
startPartitionuntuk sumber MaxCompute inkremental -
Penundaan lama sebelum sumber MaxCompute inkremental mulai membaca
-
Job dengan sumber MaxCompute macet saat startup atau tertunda
-
Error runtime dalam tabel hasil MaxCompute: 'Invalid partition spec'
-
Error runtime dalam tabel hasil MaxCompute: 'No more available blockId'
-
MySQL
-
ApsaraDB RDS for MySQL
-
Perubahan tipe untuk primary key bigint unsigned selama sinkronisasi data
-
Apakah Flink memperbarui atau menyisipkan catatan saat menulis ke RDS?
-
Mengapa INT UNSIGNED di MySQL memerlukan tipe berbeda di Flink SQL
-
Error: Incorrect string value: '\xF0\x9F\x98\x80\xF0\x9F...' untuk kolom 'test' pada baris 1
-
Skema downstream tidak berubah setelah pembaruan skema MySQL
-
Menyelesaikan kegagalan job akibat perubahan skema yang tidak didukung selama sinkronisasi CTAS/CDAS
-
-
ClickHouse
-
Print
-
Tablestore
-
ApsaraMQ for RocketMQ
-
Hologres
-
Error: remaining connection slots are reserved for non-replication superuser connections
-
Hubungan antara interval checkpoint dan visibilitas data untuk sink Hologres
-
Pengecualian 'permission denied for database' saat penerapan
-
Masalah presisi data saat mengonsumsi data binlog dalam mode JDBC
-
Error setelah menghapus dan membuat ulang tabel dengan nama yang sama
-
Log Service
-
Paimon
-
Hudi
-
AnalyticDB for MySQL (ADB)
-
Custom connector
-
Elasticsearch
Ambil data JSON dari Kafka menggunakan Flink
-
Untuk mengambil data JSON standar, lihat JSON Format.
-
Untuk mengambil data JSON bersarang, definisikan objek JSON sebagai tipe
ROWdalam DDL untuk tabel sumber. Dalam DDL untuk tabel sink, deklarasikan kunci yang akan diambil. Kemudian, gunakan pernyataan DML untuk mengakses kunci dan mengekstrak nilainya. Kode berikut memberikan contoh:-
Data sampel
{ "a":"abc", "b":1, "c":{ "e":["1","2","3","4"], "f":{"m":"567"} } } -
DDL tabel sumber
CREATE TEMPORARY TABLE `kafka_table` ( `a` VARCHAR, b int, `c` ROW<e ARRAY<VARCHAR>,f ROW<m VARCHAR>> -- 'c' adalah objek JSON yang dipetakan ke tipe ROW di Flink. 'e' adalah array JSON yang dipetakan ke tipe ARRAY. ) WITH ( 'connector' = 'kafka', 'topic' = 'xxx', 'properties.bootstrap.servers' = 'xxx', 'properties.group.id' = 'xxx', 'format' = 'json', 'scan.startup.mode' = 'xxx' ); -
DDL tabel sink
CREATE TEMPORARY TABLE `sink` ( `a` VARCHAR, b INT, e VARCHAR, `m` varchar ) WITH ( 'connector' = 'print', 'logger' = 'true' ); -
Pernyataan DML
INSERT INTO `sink` SELECT `a`, b, c.e[1], -- Flink menggunakan pengindeksan berbasis 1 untuk array. Contoh ini menggunakan indeks 1 untuk mengambil elemen pertama. Untuk mengambil seluruh array, hilangkan [1]. c.f.m FROM `kafka_table`; -
Hasil
409 2021-04-08 10:13:11,214 INFO org.apache.flink.kafka.shaded.org.apache.kafka.common.utils.AppInfoParser [] - 410 2021-04-08 10:13:11,215 INFO org.apache.flink.kafka.shaded.org.apache.kafka.common.utils.AppInfoParser [] - Kafka commitId: cb8625948210849f 411 2021-04-08 10:13:11,215 INFO org.apache.flink.kafka.shaded.org.apache.kafka.common.utils.AppInfoParser [] - Kafka startTimeMs: 1617847991214 412 2021-04-08 10:13:11,270 INFO org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-lb_test-2, groupId=lb_test] Subscribed to partition(s): lb_test-0, lb_test-1, lb_test-2, lb_test-3, lb_test-4, lb_test-5 413 2021-04-08 10:13:11,280 INFO org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-lb_test-2, groupId=lb_test] Seeking to offset 1 for partition lb_test-0 414 2021-04-08 10:13:11,285 INFO org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-lb_test-2, groupId=lb_test] Seeking to offset 0 for partition lb_test-1 415 2021-04-08 10:13:11,285 INFO org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-lb_test-2, groupId=lb_test] Seeking to offset 0 for partition lb_test-2 416 2021-04-08 10:13:11,285 INFO org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-lb_test-2, groupId=lb_test] Seeking to offset 0 for partition lb_test-3 417 2021-04-08 10:13:11,290 INFO org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-lb_test-2, groupId=lb_test] Seeking to offset 0 for partition lb_test-4 418 2021-04-08 10:13:11,291 INFO org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-lb_test-2, groupId=lb_test] Seeking to offset 0 for partition lb_test-5 419 2021-04-08 10:13:11,302 INFO org.apache.flink.kafka.shaded.org.apache.kafka.clients.Metadata [] - [Consumer clientId=consumer-lb_test-2, groupId=lb_test] Cluster ID: -1flJPwnTvuGFSuyCtU1hw 420 2021-04-08 10:15:31,597 INFO org.apache.flink.api.common.functions.util.PrintSinkOutputWriter [] - +I(abc,1,1,567)
-
Flink tidak dapat mengonsumsi atau menulis ke Kafka
-
Penyebab
Jika terdapat mekanisme penerusan, seperti proxy atau pemetaan port, antara Flink dan Kafka, klien Kafka akan mengambil alamat jaringan internal server Kafka, bukan alamat proxy. Akibatnya, Flink dapat terhubung ke kluster Kafka tetapi tidak dapat mengonsumsi atau menulis data, meskipun jalur jaringan telah terbentuk.
Proses koneksi antara konektor Flink Kafka dan server Kafka melibatkan dua langkah:
-
Klien Kafka mengambil metadata dari broker Kafka. Metadata ini mencakup alamat jaringan semua broker dalam kluster.
-
Konektor Flink kemudian menggunakan alamat jaringan tersebut untuk mengonsumsi atau menulis data.
-
-
Pemecahan Masalah
Ikuti langkah-langkah berikut untuk menentukan apakah terdapat mekanisme penerusan, seperti proxy atau pemetaan port, antara Flink dan Kafka:
-
Gunakan alat baris perintah ZooKeeper (zkCli.sh atau zookeeper-shell.sh) untuk masuk ke kluster ZooKeeper yang digunakan oleh kluster Kafka Anda.
-
Jalankan perintah yang sesuai untuk kluster Anda guna mengambil metadata broker Kafka.
Anda biasanya dapat menggunakan perintah
get /brokers/ids/0untuk mengambil metadata broker Kafka. Alamat koneksi terletak di fieldendpoints. Misalnya, hubungkan menggunakan ZooKeeper Shell dan jalankanget /brokers/ids/0untuk melihat informasi registrasi broker. Perhatikan alamat yang dikonfigurasi di fieldendpointspada JSON yang dikembalikan:# bin/zookeeper-shell.sh localhost:2181 Connecting to localhost:2181 Welcome to ZooKeeper! JLine support is disabled WATCHER:: WatchedEvent state:SyncConnected type:None path:null get /brokers/ids/0 {"listener_security_protocol_map":{"PLAINTEXT":"PLAINTEXT"},"endpoints":["PLAINTEXT://116.62.xxx:9092"],"jmx_port":-1,"host":"116.62.xxx","timestamp":"1614840078030","port":9092,"version":4} -
Gunakan perintah seperti
pingatautelnetuntuk menguji konektivitas dari lingkungan Flink ke alamat dari field endpoints.Koneksi yang gagal menunjukkan bahwa terdapat mekanisme penerusan, seperti proxy atau pemetaan port, antara Flink dan Kafka.
-
-
Solusi
-
Jangan gunakan mekanisme penerusan. Sebaliknya, buat jalur jaringan langsung antara Flink dan Kafka. Hal ini memungkinkan Flink terhubung langsung ke endpoints yang tercantum dalam metadata Kafka.
-
Hubungi administrator Kafka Anda untuk mengonfigurasi alamat penerusan dalam properti advertised.listeners pada broker Kafka. Hal ini memastikan klien Kafka mengambil metadata yang mencakup alamat penerusan yang benar.
CatatanHanya Kafka versi 0.10.2.0 dan yang lebih baru yang mendukung penambahan alamat proxy ke listeners broker Kafka.
Untuk informasi lebih lanjut tentang cara kerja ini, lihat KIP-103: Separate Internal and External traffic dan Kafka client cannot connect to brokers.
-
Jika konektivitas jaringan antara Flink dan Kafka telah dikonfirmasi dan masalah masih berlanjut, periksa penyebab non-jaringan berikut:
Periksa 1: Strategi offset awal
Verifikasi parameter scan.startup.mode dalam klausa WITH DDL tabel sumber Kafka Anda. Jika nilainya latest-offset, Flink hanya membaca pesan yang ditulis setelah job dimulai. Jika tidak ada pesan baru yang tiba setelah startup job, job tampaknya tidak mengonsumsi data.
|
|
Behavior |
|
|
Membaca dari pesan paling awal yang tersedia di setiap partisi. |
|
|
Hanya membaca pesan yang ditulis setelah job dimulai. Data yang diproduksi sebelum startup job tidak dikonsumsi. |
|
|
Dilanjutkan dari offset terakhir yang dikomit oleh kelompok konsumen. Jika tidak ada offset yang dikomit, kembali ke |
|
|
Membaca dari timestamp yang ditentukan pengguna. Memerlukan pengaturan |
Untuk memverifikasi bahwa data baru diproduksi setelah job dimulai, gunakan klien konsumen Kafka untuk memantau topik secara real time.
Periksa 2: Ketidakcocokan format data
Verifikasi bahwa parameter format dalam klausa WITH tabel sumber Kafka Anda sesuai dengan pengkodean aktual pesan dalam topik Kafka. Ketidakcocokan format menyebabkan kegagalan deserialisasi, yang dapat mengakibatkan job diam-diam melewatkan pesan atau tidak menghasilkan output.
|
Scenario |
|
|
Pesan JSON biasa |
|
|
Pesan CDC Canal |
|
|
Pesan CDC Debezium |
|
|
Pesan CDC Maxwell |
|
Untuk memeriksa format pesan aktual, gunakan klien konsumen Kafka untuk membaca byte mentah dari topik dan memeriksa struktur muatan.
Tidak ada output data dari jendela waktu event Kafka
-
Masalah
Sebuah job tidak menghasilkan output saat menggunakan tabel sumber Kafka dengan jendela waktu event.
-
Penyebab
Partisi Kafka yang tidak aktif dapat mencegah watermark maju, sehingga menghentikan jendela waktu event menghasilkan output.
-
Solusi
-
Pastikan semua partisi menerima data.
-
Untuk mengaktifkan deteksi ketidakaktifan sumber, tambahkan kode berikut ke bagian Other Configurations dan simpan perubahan Anda. Untuk instruksi terperinci, lihat Cara mengonfigurasi parameter runtime job kustom?.
table.exec.source.idle-timeout: 5Untuk informasi lebih lanjut tentang parameter table.exec.source.idle-timeout, lihat Configuration.
-
Commit offset di Kafka
Offset commit Kafka melacak posisi data yang telah diproses, memastikan konsistensi dan keandalan dalam pemrosesan aliran dengan mencegah duplikasi atau kehilangan data. Saat checkpoint berhasil diselesaikan, Flink melakukan commit offset baca yang sesuai ke Kafka. Jika checkpointing tidak diaktifkan, atau jika interval checkpoint terlalu lama, offset yang dikomit di Kafka dapat menjadi usang, menyebabkan pemrosesan ulang data atau kehilangan data.
Mengurai JSON bersarang dengan konektor Kafka
Sebagai contoh, saat mengurai data JSON berikut secara langsung dengan format json, data tersebut diselesaikan menjadi satu field bertipe ARRAY<ROW<cola VARCHAR, colb VARCHAR>>. Field ini merupakan array dari baris, di mana setiap baris berisi dua field VARCHAR. Anda kemudian dapat mengurai array ini menggunakan user-defined table-valued function (UDTF).
{"data":[{"cola":"test1","colb":"test2"},{"cola":"test1","colb":"test2"},{"cola":"test1","colb":"test2"},{"cola":"test1","colb":"test2"},{"cola":"test1","colb":"test2"}]}
Menghubungkan ke kluster Kafka yang aman
-
Dalam klausa
WITHDDL tabel Kafka Anda, tambahkan konfigurasi keamanan untuk autentikasi dan enkripsi. Untuk daftar lengkap opsi, lihat SECURITY.PentingTambahkan prefiks semua parameter konfigurasi keamanan dengan properties.
-
Contoh ini menunjukkan cara mengonfigurasi tabel Kafka untuk menggunakan mekanisme SASL
PLAINdan menyediakan konfigurasi JAAS.CREATE TABLE KafkaTable ( `user_id` BIGINT, `item_id` BIGINT, `behavior` STRING, `ts` TIMESTAMP(3) METADATA FROM 'timestamp' ) WITH ( 'connector' = 'kafka', ... 'properties.security.protocol' = 'SASL_PLAINTEXT', 'properties.sasl.mechanism' = 'PLAIN', 'properties.sasl.jaas.config' = 'org.apache.flink.kafka.shaded.org.apache.kafka.common.security.plain.PlainLoginModule required username=\"username\" password=\"password\";' ); -
Contoh ini menunjukkan cara menggunakan protokol keamanan
SASL_SSLdengan mekanisme SASLSCRAM-SHA-256.CREATE TABLE KafkaTable ( `user_id` BIGINT, `item_id` BIGINT, `behavior` STRING, `ts` TIMESTAMP(3) METADATA FROM 'timestamp' ) WITH ( 'connector' = 'kafka', ... 'properties.security.protocol' = 'SASL_SSL', /* Konfigurasi SSL */ /* Path ke truststore (sertifikat CA) yang disediakan oleh server */ 'properties.ssl.truststore.location' = '/flink/usrlib/kafka.client.truststore.jks', 'properties.ssl.truststore.password' = 'test1234', /* Jika autentikasi sisi klien diperlukan, konfigurasikan path ke keystore (kunci privat) */ 'properties.ssl.keystore.location' = '/flink/usrlib/kafka.client.keystore.jks', 'properties.ssl.keystore.password' = 'test1234', /* Konfigurasi SASL */ /* Konfigurasikan mekanisme SASL sebagai SCRAM-SHA-256 */ 'properties.sasl.mechanism' = 'SCRAM-SHA-256', /* Konfigurasikan JAAS */ 'properties.sasl.jaas.config' = 'org.apache.flink.kafka.shaded.org.apache.kafka.common.security.scram.ScramLoginModule required username=\"username\" password=\"password\";' );Catatan-
Jika
properties.sasl.mechanismadalahSCRAM-SHA-256, gunakanorg.apache.flink.kafka.shaded.org.apache.kafka.common.security.scram.ScramLoginModuleuntukproperties.sasl.jaas.config. -
Jika
properties.sasl.mechanismadalahPLAIN, gunakanorg.apache.flink.kafka.shaded.org.apache.kafka.common.security.plain.PlainLoginModuleuntukproperties.sasl.jaas.config.
-
-
-
Di bagian Additional dependency files untuk job Anda, unggah semua file yang diperlukan, seperti sertifikat, kunci publik, dan kunci privat.
Platform menyimpan file yang diunggah di direktori
/flink/usrlib. Untuk instruksi pengunggahan, lihat penerapan job.PentingJika mekanisme autentikasi pada broker Kafka Anda adalah
SASL_SSLtetapi mekanisme sisi klien adalahSASL_PLAINTEXT, job gagal dengan pengecualianOutOfMemoryselama validasi. Untuk menyelesaikan masalah ini, pastikan mekanisme autentikasi sisi klien dan sisi server cocok.
Konflik penamaan field
-
Masalah
Sumber data Kafka membuat serial pesan menjadi dua string JSON terpisah: satu untuk kunci dan satu untuk nilai. Dalam skenario ini, baik kunci maupun nilai berisi field dengan nama yang sama, seperti field id pada contoh di bawah. Mengurai data ini langsung ke tabel Flink menyebabkan konflik penamaan field.
-
kunci
{ "id": 1 } -
nilai
{ "id": 100, "name": "flink" }
-
-
Solusi
Gunakan properti key.fields-prefix untuk menghindari masalah ini.
CREATE TABLE kafka_table ( -- Definisikan kolom untuk field kunci dan nilai key_id INT, value_id INT, name STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'test_topic', 'properties.bootstrap.servers' = 'localhost:9092', 'format' = 'json', 'json.ignore-parse-errors' = 'true', -- Tentukan field dan tipe data untuk kunci 'key.format' = 'json', 'key.fields' = 'id', 'value.format' = 'json', 'value.fields' = 'id, name', -- Tambahkan prefiks ke field dari kunci 'key.fields-prefix' = 'key_' );Menyetel properti key.fields-prefix ke key_ menginstruksikan konektor untuk menambahkan prefiks key_ ke semua field dari kunci pesan. Misalnya, field id dari kunci menjadi kolom key_id dalam tabel Flink. Hal ini mencegah konflik dengan field id dari nilai, yang dipetakan ke kolom .
Menjalankan kueri
SELECT * FROM kafka_table;menghasilkan output berikut:key_id: 1, value_id: 100, name: flink
Memecahkan masalah latensi tinggi dari sumber Kafka
-
Masalah
Saat Anda membaca dari tabel sumber Kafka, metrik
currentEmitEventTimeLagmenunjukkan nilai lebih dari 50 tahun. Misalnya, beberapa job SQL Flink berada dalam status running, tetapi kolom business latency menampilkan nilai yang sangat tinggi melebihi 19.160 hari, seperti19160d 1h 59m 28s. -
Pemecahan Masalah
-
Pertama, tentukan apakah job tersebut merupakan job JAR atau job SQL.
Untuk job JAR, verifikasi bahwa file
pom.xmlAnda menggunakan dependensi Kafka yang disediakan oleh Realtime Compute for Apache Flink. Versi open-source konektor tidak melaporkan metrik ini. -
Periksa apakah semua partisi dalam topik Kafka hulu menerima data secara real time.
-
Periksa apakah
timestampdalam metadata pesan Kafka adalah 0 atau null.Latensi sumber Kafka dihitung dengan mengurangkan timestamp pesan dari waktu saat ini. Jika sebuah pesan tidak memiliki timestamp, latensi dapat ditampilkan lebih dari 50 tahun. Anda dapat memeriksa timestamp dengan salah satu cara berikut:
-
Untuk job SQL, Anda dapat mengambil timestamp pesan dengan mendefinisikan kolom metadata. Untuk informasi lebih lanjut, lihat Tabel sumber Kafka.
CREATE TEMPORARY TABLE sk_flink_src_user_praise_rt ( `timestamp` BIGINT , `timestamp` TIMESTAMP METADATA, --Timestamp metadata. ts as to_timestamp ( from_unixtime (`timestamp`, 'yyyy-MM-dd HH:mm:ss') ), watermark for ts as ts - interval '5' second ) WITH ( 'connector' = 'kafka', 'topic' = '', 'properties.bootstrap.servers' = '', 'properties.group.id' = '', 'format' = 'json', 'scan.startup.mode' = 'latest-offset', 'json.fail-on-missing-field' = 'false', 'json.ignore-parse-errors' = 'true' ); -
Tulis program Java sederhana yang menggunakan klien
KafkaConsumeruntuk membaca pesan dan memeriksa timestamp-nya.
-
-
Error: Tabel 'upsert-kafka' memerlukan PRIMARY KEY
-
Masalah
) WITH ( 'connector' = 'upsert-kafka', 'topic' = 'flow_stay_duration', 'properties.bootstrap.servers' = 'xxx', 'key.format' = 'avro', 'value.format' = 'avro' ); insert into sink_ad_data_device_info org.apache.flink.table.api.ValidationException: SQL validation failed. Unable to create a sink for writing table 'vvp.default.sink_ad_data_device_info'. The cause is following: 'upsert-kafka' tables require to define a PRIMARY KEY constraint. The PRIMARY KEY specifies which columns should be read from or write to the Kafka message key. The PRIMARY KEY also defines records in the 'upsert-kafka' table should update or delete on which keys. Table options are: 'connector'='upsert-kafka' 'key.format'='avro' 'properties.bootstrap.servers'='xxx' 'topic'='flow_stay_duration' 'value.format'='avro' at org.apache.flink.table.sqlserver.utils.FormatValidatorExceptionUtils.newValidationException(FormatValidatorExceptionUtils.java:41) at org.apache.flink.table.sqlserver.utils.ErrorConverter.formatException(ErrorConverter.java:123) at org.apache.flink.table.sqlserver.utils.ErrorConverter.toErrorDetail(ErrorConverter.java:60) at org.apache.flink.table.sqlserver.utils.ErrorConverter.toGrpcException(ErrorConverter.java:54) at org.apache.flink.table.sqlserver.FlinkSqlServiceImpl.validateAndGeneratePlan(FlinkSqlServiceImpl.java:979) at org.apache.flink.table.sqlserver.proto.FlinkSqlServiceGrpc$MethodHandlers.invoke(FlinkSqlServiceGrpc.java:3283) at io.grpc.stub.ServerCalls$UnaryServerCallHandler$UnaryServerCallListener.onHalfClose(ServerCalls.java:172) at io.grpc.internal.ServerCallImpl$ServerStreamListenerImpl.halfClosed(ServerCallImpl.java:331) -
Penyebab
Error ini terjadi karena DDL tidak memiliki primary key. Saat digunakan sebagai tabel sink, konektor
upsert-kafkamengonsumsi aliran changelog dari logika hulu. Konektor menulis data INSERT dan UPDATE_AFTER ke Kafka. Untuk operasi DELETE, konektor menulis pesan dengan nilai null untuk menunjukkan bahwa pesan untuk kunci yang sesuai dihapus. Flink menggunakan kolom primary key untuk mempartisi data. Hal ini menjamin bahwa pesan dengan kunci yang sama diurutkan dan pesan pembaruan atau penghapusan yang sesuai berada di partisi yang sama. -
Solusi
Definisikan primary key dalam DDL.
Pulihkan job Flink setelah pemisahan atau penskalaan-in topik
Jika Anda memisahkan atau menskala-in topik DataHub yang sedang dibaca oleh job Flink, job tersebut akan masuk ke loop kegagalan dan tidak dapat pulih secara otomatis. Untuk menyelesaikannya, mulai ulang job tersebut.
Menghapus topik dengan konsumen aktif
Anda tidak dapat menghapus atau membuat ulang topik DataHub yang memiliki konsumen aktif.
Parameter endPoint dan tunnelEndpoint
Parameter endPoint dan tunnelEndpoint dijelaskan dalam endpoint. Di lingkungan VPC, salah konfigurasi parameter ini dapat menyebabkan pengecualian tugas:
-
Jika parameter
endPointsalah dikonfigurasi, penerapan tugas terhenti pada progres 91%. -
Jika parameter
tunnelEndpointsalah dikonfigurasi, tugas gagal dijalankan.
Pembuatan tabel DataHub gagal dengan NoPermissionException: dhs:ListShard
-
Gejala
Saat job Flink menerapkan tabel sumber atau sink DataHub, job tersebut gagal dengan error serupa berikut:
NoPermissionException: You have no permission to perform this action. Action: dhs:ListShard -
Penyebab
Error ini disebabkan oleh penamaan parameter
WITHyang tidak standar dalam DDL DataHub. Konektor DataHub memerlukan kredensial yang ditentukan sebagaiaccessIddanaccessKey. Jika Anda menggunakan bentuk bertitikaccess.iddanaccess.keysebagai gantinya, konektor tidak dapat mengenali field kredensial dan tidak dapat melakukan autentikasi. Akibatnya, konektor mencoba mencantumkan shard tanpa kredensial yang valid, dan DataHub mengembalikanNoPermissionException.Pesan error merujuk pada hak istimewa
dhs:ListShardyang hilang, tetapi akar penyebabnya adalah parameter kredensial yang tidak dikenali — bukan kekurangan izin IAM yang sebenarnya. -
Solusi
Dalam DDL DataHub Anda, ubah nama
access.idmenjadiaccessIddanaccess.keymenjadiaccessKey. Hapus tabel yang ada dan buat ulang dengan klausaWITHyang telah dikoreksi.Konfigurasi salah:
CREATE TABLE datahub_source (...) WITH ( 'connector' = 'datahub', 'endPoint' = 'https://dh-cn-hangzhou.aliyuncs.com', 'project' = 'your_project', 'topic' = 'your_topic', 'access.id' = 'your-access-key-id', -- Salah: bentuk bertitik tidak dikenali 'access.key' = 'your-access-key-secret' -- Salah: bentuk bertitik tidak dikenali );Konfigurasi benar:
CREATE TABLE datahub_source (...) WITH ( 'connector' = 'datahub', 'endPoint' = 'https://dh-cn-hangzhou.aliyuncs.com', 'project' = 'your_project', 'topic' = 'your_topic', 'accessId' = 'your-access-key-id', -- Benar 'accessKey' = 'your-access-key-secret' -- Benar );Untuk daftar lengkap parameter
WITHyang didukung, lihat dokumentasi konektor DataHub.
Pembacaan penuh dan inkremental dari sumber MaxCompute
Sumber MaxCompute melakukan pembacaan penuh dan inkremental melalui Tunnel MaxCompute. Throughput baca dibatasi oleh bandwidth Tunnel MaxCompute.
Dapatkah tabel sumber MaxCompute membaca data yang ditambahkan?
Tidak. Setelah job Flink dimulai, job tersebut tidak akan membaca data baru yang ditambahkan ke tabel atau partisi sumber. Hal ini berlaku baik saat sumber sedang membaca aktif maupun setelah selesai membaca. Menambahkan data dengan cara ini juga dapat menyebabkan failover job.
Baik tabel sumber MaxCompute penuh maupun inkremental menggunakan ODPS DOWNLOAD SESSION untuk membaca data tabel atau partisi. Saat Anda membuat DOWNLOAD SESSION, server membuat file Indeks. File ini merupakan snapshot data pada saat DOWNLOAD SESSION dibuat, dan pembacaan data selanjutnya didasarkan pada snapshot ini. Oleh karena itu, setelah DOWNLOAD SESSION dibuat, data yang ditambahkan ke tabel atau partisi MaxCompute tidak dibaca dalam kondisi normal. Namun, jika data baru ditulis ke tabel sumber MaxCompute, dua pengecualian dapat terjadi:
-
Kegagalan selama pembacaan: Jika data baru ditulis saat Tunnel sedang membaca aktif, operasi gagal dengan error
ErrorCode=TableModified,ErrorMessage=The specified table has been modified since the download initiated.. -
Data tidak konsisten saat failover: Jika data baru ditulis setelah Tunnel ditutup, job run saat ini tidak akan membacanya. Namun, jika job mengalami failover atau dilanjutkan dari status jeda, job tersebut dapat memproses ulang data lama dan hanya membaca sebagian data baru.
Perubahan konkurensi untuk job MaxCompute yang dijeda
Untuk tabel sumber MaxCompute dengan opsi useNewApi diaktifkan (diaktifkan secara default), job dalam mode streaming mendukung perubahan konkurensi setelah dijeda dan dilanjutkan. Tabel sumber MaxCompute membaca partisi yang cocok secara berurutan. Saat membaca partisi, data dalam partisi tersebut didistribusikan ke operator paralel. Mengubah konkurensi tidak memengaruhi distribusi data untuk partisi yang sedang diproses sebelum jeda. Paralelisme baru hanya berlaku saat job mulai memproses partisi berikutnya. Akibatnya, jika job sedang memproses satu partisi besar, meningkatkan konkurensi dan melanjutkan job dapat menyebabkan hanya beberapa operator MaxCompute yang membaca data.
Perubahan konkurensi tidak didukung untuk job batch atau untuk job di mana opsi useNewApi diatur ke false.
Mengapa MaxCompute membaca partisi sebelumnya saat posisi awal adalah 2019-10-11 00:00:00?
Pengaturan posisi awal hanya memengaruhi sumber data antrian pesan, seperti DataHub. Pengaturan ini tidak memengaruhi tabel sumber MaxCompute. Saat job Flink dimulai, job tersebut membaca data sebagai berikut:
-
Untuk tabel berpartisi: Semua partisi yang ada dibaca.
-
Untuk tabel non-partisi: Semua data yang ada dibaca.
Cegah pembacaan data tidak lengkap dari partisi baru
Saat ini, tidak ada mekanisme untuk memverifikasi apakah data dalam partisi sudah lengkap. Akibatnya, tabel sumber MaxCompute inkremental mulai membaca partisi baru segera setelah terdeteksi. Misalkan Anda menggunakan tabel sumber MaxCompute inkremental untuk membaca tabel MaxCompute berpartisi T di mana ds adalah kolom partisi. Dalam kasus ini, kami menyarankan agar Anda tidak membuat partisi terlebih dahulu. Sebagai gantinya, jalankan pernyataan INSERT OVERWRITE TABLE T PARTITION (ds='20191010') .... Saat job selesai, partisi dan datanya muncul secara bersamaan.
Jangan buat partisi terlebih dahulu (misalnya, ds=20191010) lalu tulis data ke dalamnya. Jika Anda menggunakan metode ini, tabel sumber MaxCompute inkremental mendeteksi partisi baru ds=20191010 dan segera mulai membacanya. Hal ini menyebabkan pembacaan data tidak lengkap jika operasi penulisan masih berlangsung.
Error otorisasi konektor MaxCompute
-
Detail error
Saat eksekusi job, error ditampilkan di halaman failover atau dalam file TaskManager.log:
ErrorMessage=Authorization Failed [4019], You have NO privilege'ODPS:***' -
Penyebab
Informasi identitas pengguna yang ditentukan dalam definisi DDL MaxCompute tidak memiliki izin yang diperlukan untuk mengakses MaxCompute.
-
Solusi
Lakukan autentikasi menggunakan akun Alibaba Cloud, pengguna RAM, atau peran RAM. Untuk informasi lebih lanjut, lihat otentikasi pengguna.
Mengonfigurasi parameter startPartition
|
Langkah |
Deskripsi |
Contoh |
|
1 |
Sambungkan setiap nama kolom partisi dengan nilai tetapnya menggunakan tanda sama dengan (=). |
Jika kolom partisi adalah dt dan Anda ingin membaca data mulai dari nilai partisi 20220901, hasilnya adalah dt=20220901. |
|
2 |
Urutkan hasil dari Langkah 1 berdasarkan level partisi secara ascending dan gabungkan dengan koma (,) tanpa spasi. String ini adalah nilai untuk parameter startPartition. Catatan
Anda hanya dapat menentukan beberapa level partisi pertama. |
|
Saat sistem memuat daftar partisi, sistem membandingkan setiap partisi dengan nilai startPartition secara leksikografis. Sistem kemudian memuat semua partisi yang lebih besar dari atau sama dengan nilai startPartition. Misalnya, pertimbangkan tabel berpartisi MaxCompute untuk pembacaan inkremental dengan partisi tingkat pertama ds dan partisi tingkat kedua type. Tabel tersebut berisi enam partisi berikut:
-
ds=20191201,type=a
-
ds=20191201,type=b
-
ds=20191202,type=a
-
ds=20191202,type=b
-
ds=20191202,type=c
-
ds=20191203,type=a
Jika startPartition diatur ke ds=20191202, sistem membaca empat partisi: ds=20191202,type=a, ds=20191202,type=b, ds=20191202,type=c, dan ds=20191203,type=a. Jika startPartition diatur ke ds=20191202,type=b, sistem membaca tiga partisi: ds=20191202,type=b, ds=20191202,type=c, dan ds=20191203,type=a.
Partisi yang ditentukan dalam startPartition tidak harus ada. Sistem membaca semua partisi yang secara leksikografis lebih besar dari atau sama dengan nilai startPartition.
Startup lambat untuk job MaxCompute inkremental
Job dimulai dengan lambat karena pertama-tama harus memproses metadata untuk semua partisi yang secara leksikografis lebih besar dari atau sama dengan nilai startPartition. Proses ini sangat tertunda oleh jumlah partisi yang besar atau file-file kecil. Untuk mengurangi penundaan ini, ikuti rekomendasi berikut:
-
Hindari membaca terlalu banyak data historis.
CatatanJika Anda perlu memproses data historis, jalankan job batch dengan tabel sumber MaxCompute sebagai gantinya.
-
Kurangi jumlah file kecil dalam data historis.
Atur parameter partisi
Baca dari partisi
-
Baca dari partisi statis
Saat membaca dari partisi statis tabel sumber atau tabel dimensi, atur parameter
partitionsebagai berikut.Langkah
Deskripsi
Contoh
1
-
Untuk tabel dimensi, tentukan setiap partisi sebagai 'partition_column_name=partition_value'. Nilai partisi harus berupa nilai tetap.
-
Untuk tabel sumber, tentukan setiap partisi sebagai 'partition_column_name=partition_value'. Nilai partisi dapat berupa nilai tetap atau nilai yang berisi wildcard (*). Wildcard dapat mencocokkan string apa pun, termasuk string kosong.
-
Untuk membaca data dari kolom partisi
dtdengan nilai20220901, tentukandt=20220901. -
Untuk membaca data dari partisi dalam kolom
dtyang nilainya dimulai dengan202209, tentukandt=202209*(hanya berlaku untuk tabel sumber). -
Untuk membaca data dari partisi dalam kolom
dtyang nilainya dimulai dengan2022dan diakhiri dengan01, tentukandt=2022*01(hanya berlaku untuk tabel sumber). -
Untuk membaca data dari semua partisi dalam kolom
dt, tentukandt=*(hanya berlaku untuk tabel sumber).
2
Urutkan string partisi dari Langkah 1 berdasarkan level partisi secara ascending, lalu gabungkan dengan koma (tanpa spasi). String yang dihasilkan adalah nilai parameter
partition.Anda hanya dapat menentukan beberapa level partisi pertama.
-
Tabel memiliki satu partisi tingkat pertama
dt. Untuk membaca data dari partisidt=20220901, tentukan'partition' = 'dt=20220901'. -
Tabel memiliki tiga level partisi: partisi tingkat pertama
dt, partisi tingkat keduahh, dan partisi tingkat ketigamm. Untuk membaca data daridt=20220901,hh=08, danmm=10, tentukan'partition' = 'dt=20220901,hh=08,mm=10'. -
Untuk tabel yang sama, untuk membaca data dari
dt=20220901,hh=08, dan nilai apa pun untukmm, tentukan'partition' = 'dt=20220901,hh=08' or 'partition' = 'dt=20220901,hh=08,mm=*'. -
Untuk tabel yang sama, untuk membaca data dari
dt=20220901, nilai apa pun untukhh, danmm=10, tentukan'partition' = 'dt=20220901,hh=*,mm=10'.
Jika langkah-langkah ini tidak memenuhi kebutuhan pemfilteran partisi Anda, Anda dapat menambahkan kondisi filter ke klausa
WHEREpernyataan SQL Anda. Hal ini memungkinkan pengoptimal SQL menggunakan pushdown partisi untuk pemfilteran. Misalnya, untuk membaca partisi dari tabel dengan dua level partisi (dtdanhh) di manadtberada di antara '20220901' dan '20220903', danhhberada di antara '09' dan '17', gunakan pernyataan SQL seperti berikut.CREATE TABLE maxcompute_table ( content VARCHAR, dt VARCHAR, hh VARCHAR ) PARTITIONED BY (dt, hh) WITH ( -- Anda harus menentukan kolom partisi dengan PARTITIONED BY untuk mengaktifkan -- pushdown partisi dalam pengoptimal SQL, yang meningkatkan kinerja. 'connector' = 'odps', ... -- Isi parameter yang diperlukan seperti accessId. Anda dapat menghilangkan parameter 'partition' dan membiarkan pengoptimal SQL memfilter partisi. ); SELECT content, dt, hh FROM maxcompute_table WHERE dt >= '20220901' AND dt <= '20220903' AND hh >= '09' AND hh <= '17'; -- Tentukan filter partisi dalam klausa WHERE. -
-
Baca partisi dengan urutan leksikografis terbesar
-
Untuk membaca partisi dengan urutan leksikografis terbesar dari tabel sumber atau tabel dimensi, atur parameter
partitionke'max_pt()'. -
Untuk membaca dua partisi dengan urutan leksikografis terbesar dari tabel sumber atau tabel dimensi, atur parameter
partitionke'max_two_pt()'. -
Untuk membaca partisi dengan urutan leksikografis terbesar yang juga memiliki partisi
.doneyang sesuai dari tabel sumber atau tabel dimensi, atur parameterpartitionke'max_pt_with_done()'.
Umumnya, partisi dengan urutan leksikografis terbesar adalah partisi yang paling baru dibuat. Opsi
max_pt_with_done()berguna ketika data dalam partisi terbaru mungkin belum siap, dan Anda ingin tabel dimensi sementara membaca dari partisi yang sedikit lebih lama tetapi lengkap.Saat data untuk partisi siap, Anda juga harus membuat partisi kosong yang sesuai. Namanya adalah nama partisi data dengan
.doneditambahkan. Misalnya, setelah data untuk partisidt=20220901siap, buat partisi kosong bernamadt=20220901.done. Saat Anda mengatur parameterpartitionkemax_pt_with_done(), tabel dimensi hanya membaca dari partisi yang memiliki partisi.doneyang sesuai. Partisi data tanpa partisi.donesementara diabaikan. Untuk informasi lebih lanjut, lihat Apa perbedaan antara max_pt() dan max_pt_with_done()?.CatatanTabel sumber hanya menentukan partisi dengan urutan leksikografis terbesar saat job dimulai. Tabel berhenti setelah membaca semua data dan tidak memantau partisi baru. Jika Anda perlu terus membaca partisi baru, gunakan mode tabel sumber inkremental. Tabel dimensi memeriksa dan membaca data terbaru setiap kali diperbarui.
-
Tulis ke partisi
-
Tulis ke partisi statis
Untuk menulis data ke partisi statis tabel hasil, Anda dapat mengatur parameter
partitionmenggunakan metode yang sama seperti untuk membaca dari partisi statis.PentingParameter
partitionuntuk tabel hasil tidak mendukung wildcard (*). -
Tulis ke partisi dinamis
Untuk menulis ke partisi dinamis, di mana nilai partisi berasal dari data, atur parameter partisi ke daftar nama kolom partisi yang dipisahkan koma, diurutkan berdasarkan level partisi secara ascending. Misalnya, jika tabel memiliki tiga level partisi,
dt,hh, danmm, tentukan'partition' = 'dt,hh,mm'.
Startup job lambat untuk tabel sumber MaxCompute
Penyebab yang mungkin meliputi:
-
Tabel MaxCompute berisi terlalu banyak file kecil.
-
Latensi jaringan tinggi terjadi jika kluster penyimpanan MaxCompute dan kluster komputasi Flink berada di wilayah yang berbeda. Untuk menyelesaikannya, tempatkan kedua kluster di wilayah yang sama.
-
Izin MaxCompute dikonfigurasi salah. Membaca dari tabel sumber memerlukan izin unduh untuk tabel MaxCompute.
Pilih saluran data
MaxCompute menyediakan dua saluran data: Batch Tunnel dan Streaming Tunnel. Anda dapat memilih saluran data berdasarkan persyaratan konsistensi dan efisiensi running Anda. Tabel berikut membandingkan kedua saluran data tersebut.
|
Kriteria |
Batch Tunnel |
Streaming Tunnel |
|
Konsistensi |
Batch Tunnel umumnya menulis data ke tabel MaxCompute lebih andal daripada Streaming Tunnel dan menjamin tidak ada kehilangan data (semantik at-least-once). Duplikasi data dapat terjadi di beberapa partisi, tetapi hanya jika terjadi pengecualian selama proses checkpoint saat job menulis ke beberapa partisi secara bersamaan. |
Menjamin tidak ada kehilangan data (semantik at-least-once). Namun, duplikasi data dapat terjadi jika job gagal karena alasan apa pun. |
|
Efisiensi running |
Efisiensi running keseluruhan lebih rendah daripada Streaming Tunnel karena data harus dikomit selama proses checkpoint, yang melibatkan operasi sisi server seperti pembuatan file. |
Data tidak perlu dikomit selama proses checkpoint. Jika Anda menggunakan Streaming Tunnel dan mengatur parameter |
Jika job yang menggunakan MaxCompute Batch Tunnel mengalami checkpoint lambat atau timeout, pertimbangkan untuk beralih ke Streaming Tunnel, asalkan sistem downstream Anda dapat mentolerir duplikasi data.
Duplikasi data dalam tabel hasil MaxCompute
Data duplikat dalam tabel hasil MaxCompute yang ditulis oleh job Flink dapat disebabkan oleh hal-hal berikut:
-
Periksa logika job Anda. Bahkan jika kendala primary key dideklarasikan dalam tabel hasil MaxCompute, Flink tidak melakukan pemeriksaan keunikan saat menulis ke penyimpanan eksternal. Selain itu, tabel non-transaksional di MaxCompute tidak mendukung kendala primary key. Oleh karena itu, jika logika job Flink Anda menghasilkan data duplikat, duplikat tersebut akan ditulis ke tabel MaxCompute.
-
Verifikasi apakah beberapa job Flink menulis ke tabel MaxCompute yang sama secara bersamaan. Seperti disebutkan, MaxCompute tidak menegakkan kendala primary key. Jika beberapa job Flink menghasilkan hasil yang sama, mereka akan membuat catatan duplikat dalam tabel.
-
Job Flink gagal selama checkpoint saat menggunakan Batch Tunnel. Saat kegagalan terjadi selama checkpoint, data untuk tabel hasil mungkin telah dikomit ke server. Akibatnya, saat job pulih dari checkpoint terakhir yang berhasil, job tersebut dapat menulis data duplikat untuk periode antara checkpoint terakhir yang berhasil dan kegagalan.
-
Terjadi failover job Flink saat menggunakan Stream Tunnel. Saat menulis ke MaxCompute dengan Stream Tunnel, data dikomit ke server MaxCompute antara checkpoint. Jika job mengalami failover dan pulih dari checkpoint terbaru, job tersebut dapat menulis data duplikat yang diproses setelah checkpoint selesai tetapi sebelum failover. Untuk informasi lebih lanjut, lihat Pilih saluran data. Untuk mencegah jenis duplikasi ini, Anda dapat beralih ke mode Batch Tunnel.
-
Job Flink yang menggunakan Batch Tunnel mengalami failover atau dimulai ulang setelah dibatalkan (misalnya, dipicu oleh Autopilot). Dalam versi sebelum
vvr-6.0.7-flink-1.15, job mengomit data ke tabel hasil MaxCompute saat dimatikan. Akibatnya, saat job Flink berhenti lalu pulih dari checkpoint terakhir, job tersebut dapat membuat data duplikat untuk periode antara checkpoint terakhir dan shutdown. Untuk menyelesaikan masalah ini, tingkatkan versi Flink Anda kevvr-6.0.7-flink-1.15atau yang lebih baru.
Job MaxCompute gagal dengan 'Invalid partition spec'
-
Penyebab: Error ini terjadi saat data yang ditulis ke MaxCompute berisi nilai tidak valid dalam kolom partisi. Nilai tidak valid mencakup string kosong, nilai null, atau nilai yang berisi tanda sama dengan (=), koma (,), atau garis miring (/).
-
Solusi: Verifikasi bahwa nilai dalam kolom partisi data sumber Anda valid.
Error 'No more available blockId' dalam job MaxCompute
-
Penyebab: Jumlah blok yang ditulis ke tabel hasil MaxCompute telah melebihi batas, biasanya karena terlalu sering menulis data dalam jumlah kecil.
-
Solusi: Sesuaikan parameter
batchSizedanflushIntervalMs.
Gunakan petunjuk SHUFFLE_HASH
Secara default, setiap instance paralel menyimpan cache seluruh tabel dimensi. Jika tabel dimensi besar, Anda dapat menggunakan petunjuk SHUFFLE_HASH untuk mendistribusikan data tabel secara merata di seluruh instance paralel dan mengurangi konsumsi memori heap JVM. Dalam contoh berikut, data dari tabel dimensi dim_1 dan dim_3 didistribusikan di antara instance paralel, sedangkan data dari dim_2 tetap sepenuhnya di-cache pada masing-masing instance.
-- Buat tabel sumber dan tiga tabel dimensi.
CREATE TABLE source_table (k VARCHAR, v VARCHAR) WITH ( ... );
CREATE TABLE dim_1 (k VARCHAR, v VARCHAR) WITH ('connector' = 'odps', 'cache' = 'ALL', ... );
CREATE TABLE dim_2 (k VARCHAR, v VARCHAR) WITH ('connector' = 'odps', 'cache' = 'ALL', ... );
CREATE TABLE dim_3 (k VARCHAR, v VARCHAR) WITH ('connector' = 'odps', 'cache' = 'ALL', ... );
-- Tentukan nama tabel dimensi yang akan didistribusikan dalam petunjuk SHUFFLE_HASH.
SELECT /*+ SHUFFLE_HASH(dim_1), SHUFFLE_HASH(dim_3) */
k, s.v, d1.v, d2.v, d3.v
FROM source_table AS s
INNER JOIN dim_1 FOR SYSTEM_TIME AS OF PROCTIME() AS d1 ON s.k = d1.k
LEFT JOIN dim_2 FOR SYSTEM_TIME AS OF PROCTIME() AS d2 ON s.k = d2.k
LEFT JOIN dim_3 FOR SYSTEM_TIME AS OF PROCTIME() AS d3 ON s.k = d3.k;
Mengonfigurasi CacheReloadTimeBlackList
Menentukan jendela waktu di mana pembaruan tabel dimensi dinonaktifkan.
-
tipe data:
String -
Gunakan
->antara waktu mulai dan waktu selesai. -
Pisahkan beberapa jendela waktu dengan
,. -
Format waktu:
YYYY-MM-DD HH:mm. Jika Anda hanya menentukan jam dan menit, jendela waktu berlaku setiap hari secara default.
'cacheReloadTimeBlackList' = '14:00 -> 15:00,23:00 -> 01:00'
|
Skenario |
Nilai |
|
Jendela waktu tunggal |
14:00 -> 15:00 |
|
Beberapa jendela waktu |
14:00 -> 15:00,23:00 -> 01:00 |
|
Jendela waktu khusus |
14:00 -> 15:00, 23:00 -> 01:00,2025-10-01 22:00 -> 2025-10-01 23:00 |
Error: java.io.EOFException: SSL peer shut down incorrectly
-
Detail kesalahan
Caused by: java.io.EOFException: SSL peer shut down incorrectly at sun.security.ssl.SSLSocketInputRecord.decodeInputRecord(SSLSocketInputRecord.java:239) ~[?:1.8.0_302] at sun.security.ssl.SSLSocketInputRecord.decode(SSLSocketInputRecord.java:190) ~[?:1.8.0_302] at sun.security.ssl.SSLTransport.decode(SSLTransport.java:109) ~[?:1.8.0_302] at sun.security.ssl.SSLSocketImpl.decode(SSLSocketImpl.java:1392) ~[?:1.8.0_302] at sun.security.ssl.SSLSocketImpl.readHandshakeRecord(SSLSocketImpl.java:1300) ~[?:1.8.0_302] at sun.security.ssl.SSLSocketImpl.startHandshake(SSLSocketImpl.java:435) ~[?:1.8.0_302] at com.mysql.cj.protocol.ExportControlled.performTlsHandshake(ExportControlled.java:347) ~[?:?] at com.mysql.cj.protocol.StandardSocketFactory.performTlsHandshake(StandardSocketFactory.java:194) ~[?:?] at com.mysql.cj.protocol.a.NativeSocketConnection.performTlsHandshake(NativeSocketConnection.java:101) ~[?:?] at com.mysql.cj.protocol.a.NativeProtocol.negotiateSSLConnection(NativeProtocol.java:308) ~[?:?] at com.mysql.cj.protocol.a.NativeAuthenticationProvider.connect(NativeAuthenticationProvider.java:204) ~[?:?] at com.mysql.cj.protocol.a.NativeProtocol.connect(NativeProtocol.java:1369) ~[?:?] at com.mysql.cj.NativeSession.connect(NativeSession.java:133) ~[?:?] at com.mysql.cj.jdbc.ConnectionImpl.connectOneTryOnly(ConnectionImpl.java:949) ~[?:?] at com.mysql.cj.jdbc.ConnectionImpl.createNewIO(ConnectionImpl.java:819) ~[?:?] at com.mysql.cj.jdbc.ConnectionImpl.<init>(ConnectionImpl.java:449) ~[?:?] at com.mysql.cj.jdbc.ConnectionImpl.getInstance(ConnectionImpl.java:242) ~[?:?] at com.mysql.cj.jdbc.NonRegisteringDriver.connect(NonRegisteringDriver.java:198) ~[?:?] at org.apache.flink.connector.jdbc.internal.connection.SimpleJdbcConnectionProvider.getOrEstablishConnection(SimpleJdbcConnectionProvider.java:128) ~[?:?] at org.apache.flink.connector.jdbc.internal.AbstractJdbcOutputFormat.open(AbstractJdbcOutputFormat.java:54) ~[?:?] ... 14 more -
Penyebab
Error ini biasanya terjadi saat database MySQL memiliki protokol SSL diaktifkan, tetapi koneksi SSL klien tidak dikonfigurasi dengan benar. Misalnya, dengan driver MySQL 8.0.27 dan database MySQL yang diaktifkan SSL, error ini terjadi karena metode akses default driver tidak menggunakan SSL.
-
Solusi
Tambahkan
characterEncoding=utf-8&useSSL=falseke parameter URL tabel dimensi MySQL. Misalnya:'url'='jdbc:mysql://***.***.***.***:3306/test?characterEncoding=utf-8&useSSL=false'
Perubahan tipe kunci MySQL bigint unsigned
Flink tidak mendukung tipe data bigint unsigned. Untuk mencegah potensi overflow data, Flink memetakan primary key bigint unsigned ke tipe decimal. Selama sinkronisasi ke Hologres, sistem mengonversi kolom ke tipe text karena Hologres tidak mendukung bigint unsigned dan tidak mengizinkan tipe decimal sebagai primary key.
Pertimbangkan perilaku ini selama desain dan pengembangan Anda. Untuk menjaga kolom sebagai tipe decimal, buat tabel secara manual di Hologres sebelum memulai sinkronisasi. Dalam tabel ini, Anda dapat menetapkan kolom lain sebagai primary key atau tidak menentukan primary key sama sekali. Namun, pendekatan ini dapat menyebabkan duplikasi data karena primary key asli tidak lagi menegakkan keunikan. Anda harus menangani masalah ini di tingkat aplikasi, misalnya dengan mentolerir duplikasi data atau dengan menerapkan logika deduplikasi.
Flink ke RDS: Update vs. insert
Jika primary key didefinisikan dalam DDL, konektor menggunakan pernyataan INSERT INTO tablename(field1,field2, field3, ...) VALUES(value1, value2, value3, ...) ON DUPLICATE KEY UPDATE field1=value1,field2=value2, field3=value3, ...;. Pernyataan ini menyisipkan catatan baru jika primary key tidak ada, atau memperbarui catatan yang ada jika ada. Jika tidak ada primary key yang dideklarasikan dalam DDL, konektor menyisipkan catatan baru dengan pernyataan insert into.
Menggunakan indeks unik dengan GROUP BY
-
Anda harus mendeklarasikan indeks unik dalam klausa GROUP BY job Anda.
-
Jika tabel RDS menggunakan primary key auto-increment, jangan deklarasikan sebagai PRIMARY KEY dalam job Flink.
Pemetaan INT UNSIGNED: MySQL ke Flink SQL
Driver JDBC MySQL memetakan integer tak bertanda ke tipe data Java yang lebih besar untuk mempertahankan presisi. Secara khusus, driver memetakan nilai INT UNSIGNED MySQL ke tipe Java LONG, yang kemudian diperlakukan oleh Flink SQL sebagai BIGINT. Demikian pula, driver memetakan nilai BIGINT UNSIGNED MySQL ke tipe Java BigInteger, yang diperlakukan oleh Flink SQL sebagai DECIMAL(20, 0).
Error: Incorrect string value
-
Detail error
Caused by: java.sql.BatchUpdateException: Incorrect string value: '\xF0\x9F\x98\x80\xF0\x9F...' for column 'test' at row 1 at sun.reflect.GeneratedConstructorAccessor59.newInstance(Unknown Source) at sun.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45) at java.lang.reflect.Constructor.newInstance(Constructor.java:423) at com.mysql.cj.util.Util.handleNewInstance(Util.java:192) at com.mysql.cj.util.Util.getInstance(Util.java:167) at com.mysql.cj.util.Util.getInstance(Util.java:174) at com.mysql.cj.jdbc.exceptions.SQLError.createBatchUpdateException(SQLError.java:224) at com.mysql.cj.jdbc.ClientPreparedStatement.executeBatchedInserts(ClientPreparedStatement.java:755) at com.mysql.cj.jdbc.ClientPreparedStatement.executeBatchInternal(ClientPreparedStatement.java:426) at com.mysql.cj.jdbc.StatementImpl.executeBatch(StatementImpl.java:796) at com.alibaba.druid.pool.DruidPooledPreparedStatement.executeBatch(DruidPooledPreparedStatement.java:565) at com.alibaba.ververica.connectors.rds.sink.RdsOutputFormat.executeSql(RdsOutputFormat.java:488) ... 15 more -
Penyebab
Data berisi karakter khusus atau menggunakan pengkodean karakter yang tidak didukung oleh database.
-
Solusi
Tambahkan
characterEncoding=UTF-8ke URL saat menghubungkan ke database MySQL melalui JDBC. Misalnya:jdbc:mysql://<alamat internal>/<namaDatabase>?characterEncoding=UTF-8.
Deadlock di MySQL (TDDL/RDS)
-
Masalah
Terjadi deadlock saat menulis data ke MySQL (TDDL/RDS).
PentingDi Realtime Compute for Apache Flink, jika Anda menggunakan database relasional seperti MySQL sebagai sink (melalui konektor TDDL/RDS), penulisan yang sering ke tabel atau resource dapat menyebabkan deadlock.
Asumsikan bahwa operasi
INSERTmemerlukan dua kunci, (A,B), secara berurutan. Kunci A adalah kunci rentang. Terdapat dua transaksi, (T1,T2), dan skema tabel adalah(id(primary key auto-incrementing),nid(unique key)). T1 berisi dua pernyataan,insert(null,2),(null,1), sedangkan T2 berisi satu pernyataan,insert(null,2).-
Pada waktu t, T1 mengeksekusi pernyataan
INSERTpertamanya. T1 sekarang memegang kedua kunci (A,B). -
Pada waktu t+1, T2 memulai operasi insert dan perlu menunggu kunci A untuk mengunci rentang (-inf,2]. Pada saat ini, kunci A dipegang oleh T1 dan telah mengunci rentang (-inf,2]. Karena rentang memiliki hubungan inklusi, T2 bergantung pada T1 untuk melepaskan A.
-
Pada waktu t+2, T1 mengeksekusi pernyataan INSERT keduanya, yang memerlukan kunci A pada rentang (-inf,1]. Karena rentang ini merupakan subset dari (-inf,2], T1 harus mengantri dan menunggu T2 melepaskan kunci. Akibatnya, T1 menjadi bergantung pada T2 untuk melepaskan kunci A.
Terjadi deadlock karena T1 dan T2 sekarang saling menunggu untuk melepaskan kunci masing-masing.
-
-
RDS/TDDL dan Tablestore menggunakan mekanisme penguncian yang berbeda.
-
RDS/TDDL: Kunci baris InnoDB diterapkan pada indeks, bukan pada catatan individual. Akibatnya, bahkan saat mengakses baris yang berbeda, konflik kunci dapat terjadi jika baris tersebut berbagi kunci indeks yang sama. Hal ini dapat mencegah pembaruan di seluruh rentang data.
-
Tablestore: Menggunakan kunci baris tunggal, yang tidak memengaruhi pembaruan data lainnya.
-
-
Solusi untuk deadlock
Untuk skenario QPS/TPS tinggi atau penulisan konkurensi tinggi, gunakan Tablestore sebagai tabel hasil untuk mencegah deadlock. Secara umum, tidak disarankan menggunakan TDDL atau RDS sebagai tabel hasil untuk job Flink.
Jika Anda harus menggunakan database relasional seperti MySQL sebagai node sink, pertimbangkan rekomendasi berikut:
-
Pastikan tidak ada beban kerja lain yang membaca dari atau menulis ke tabel yang sama.
-
Jika volume data job kecil, coba tulis data secara single-threaded. Namun, dalam skenario QPS/TPS tinggi dan konkurensi tinggi, pendekatan ini mengurangi kinerja penulisan.
-
Hindari menggunakan kunci unik jika memungkinkan, karena menulis ke tabel dengan kunci unik dapat menyebabkan deadlock. Jika persyaratan bisnis Anda mewajibkan kunci unik, definisikan dengan mengurutkan kolom dari yang paling selektif hingga yang paling tidak selektif. Hal ini secara signifikan mengurangi kemungkinan deadlock. Misalnya, Anda dapat menempatkan kolom dengan hash MD5 sebelum kolom
day_time(20171010)untuk memastikan kunci unik didefinisikan dengan kolom yang paling selektif terlebih dahulu. -
Gunakan sharding database dan pemisahan tabel berdasarkan karakteristik beban kerja Anda untuk mendistribusikan penulisan ke beberapa tabel. Untuk detail implementasi, hubungi administrator database Anda.
-
Kegagalan pembaruan struktur tabel downstream
Sinkronisasi struktur tabel tidak melacak pernyataan DDL. Sebaliknya, sinkronisasi mendeteksi perubahan skema dengan membandingkan catatan berurutan. Jika perubahan DDL terjadi tanpa perubahan data hulu berikutnya, struktur tabel downstream tidak akan diperbarui. Untuk detailnya, lihat kebijakan sinkronisasi untuk perubahan struktur tabel.
Kesalahan Timeout Tanggapan Penyelesaian Pemisahan
Error ini terjadi saat pemanfaatan CPU yang tinggi pada tugas mencegahnya merespons permintaan RPC dari koordinator. Untuk menyelesaikannya, tingkatkan sumber daya CPU untuk TaskManager di halaman konfigurasi sumber daya.
Dampak perubahan skema selama pemuatan penuh
Perubahan skema selama fase pemuatan penuh dapat menyebabkan job gagal atau mencegah perubahan skema disinkronkan. Untuk menyelesaikannya, hentikan job, hapus tabel downstream, lalu mulai ulang job tanpa state.
Perubahan skema yang tidak didukung selama sinkronisasi CTAS/CDAS
Sinkronkan ulang data untuk tabel tersebut. Untuk melakukannya, hentikan job, hapus tabel downstream, lalu mulai ulang job sinkronisasi dengan startup tanpa state. Hindari membuat perubahan yang tidak kompatibel seperti ini, karena job akan gagal lagi saat dimulai ulang. Untuk detail tentang perubahan skema yang didukung, lihat Pernyataan CREATE TABLE AS (CTAS).
Pembaruan retraksi di ClickHouse
Pembaruan retraksi didukung untuk tabel hasil ClickHouse jika Anda menentukan primary key dalam DDL untuk tabel hasil Flink dan mengatur parameter ignoreDelete ke false. Namun, hal ini menyebabkan penurunan kinerja yang signifikan.
ClickHouse adalah sistem manajemen basis data kolom yang dirancang untuk pemrosesan analitik online (OLAP), dan dukungannya untuk operasi UPDATE dan DELETE terbatas. Jika Anda menentukan primary key dalam DDL Flink, konektor mencoba menggunakan ALTER TABLE UPDATE dan ALTER TABLE DELETE untuk memperbarui dan menghapus data. Operasi ini sangat tidak efisien.
Visibilitas data di ClickHouse
-
Untuk tabel hasil ClickHouse dengan
exactly-once semanticsdinonaktifkan (default), data menjadi terlihat begitu buffer disiram. Sistem secara otomatis menyiram buffer ini saat jumlah catatan mencapai nilaibatchSizeatau waktu sejak penulisan terakhir melebihiflushIntervalMs. Anda tidak perlu menunggucheckpointselesai. -
Untuk tabel hasil ClickHouse dengan
exactly-once semanticsdiaktifkan, data menjadi terlihat hanya setelahcheckpointyang sesuai berhasil selesai.
Lihat hasil print
Ada dua cara untuk melihat hasil print:
-
Di Real-time Compute Development Console:
-
Dari panel navigasi kiri Real-time Compute Development Console, pilih .
-
Klik nama job target.
-
Klik tab Job Log.
-
Di tab Runtime Log, pilih job yang sedang berjalan dari daftar drop-down di samping Job.
-
Di tab Running Task Managers, klik Path, ID.
-
Klik tab Log untuk melihat hasil print.
-
-
Di UI Flink:
-
Dari panel navigasi kiri Real-time Compute Development Console, pilih .
-
Klik nama job target.
-
Di tab Status Overview, klik Flink UI.
-
Klik Task Managers.
-
Klik Path, ID.
-
Di tab logs, lihat hasil print.
-
Join tabel dimensi tidak mengembalikan data
Pastikan skema—termasuk tipe data dan nama kolom—dalam pernyataan DDL konsisten dengan tabel fisik.
max_pt() dan max_pt_with_done()
Fungsi max_pt() mengembalikan partisi dengan urutan leksikografis terbesar. Fungsi max_pt_with_done() mengembalikan partisi dengan urutan leksikografis terbesar yang memiliki partisi.done yang sesuai. Misalnya, pertimbangkan daftar partisi berikut:
-
ds=20190101
-
ds=20190101.done
-
ds=20190102
-
ds=20190102.done
-
ds=20190103
Berdasarkan daftar ini, max_pt() dan max_pt_with_done() berperilaku sebagai berikut:
-
`partition`='max_pt_with_done()'mengembalikan partisids=20190102. -
`partition`='max_pt()'mengembalikan partisids=20190103.
Error job tulis Paimon: "Heartbeat of TaskManager timed out"
Penyebab paling mungkin dari error ini adalah memori heap yang tidak mencukupi pada TaskManager. Paimon terutama menggunakan memori heap dengan cara berikut:
-
Setiap instance paralel operator penulis untuk tabel primary key Paimon memiliki buffer memori untuk pengurutan. Ukuran buffer ini dikontrol oleh properti tabel
write-buffer-size, yang default-nya 256 MB. -
Paimon menggunakan format file ORC secara default, yang memerlukan buffer memori tambahan untuk mengonversi data dalam memori ke format kolom secara batch. Ukuran buffer ini dikontrol oleh properti tabel
orc.write.batch-size, yang default-nya 1024, artinya buffer menampung 1024 baris data. -
Setiap bucket yang dimodifikasi memiliki objek penulis khusus untuk menulis datanya.
Berdasarkan pola penggunaan ini, berikut adalah penyebab potensial memori heap yang tidak mencukupi dan solusinya:
-
Nilai
write-buffer-sizeterlalu besar.Coba kurangi parameter ini. Namun, buffer yang terlalu kecil dapat menyebabkan penulisan disk yang sering dan memicu kompaksi file kecil lebih sering, yang berdampak pada kinerja penulisan.
-
Satu catatan data terlalu besar.
Misalnya, jika catatan berisi field JSON 4 MB, buffer ORC dapat tumbuh menjadi 4 MB × 1024 = 4 GB, mengonsumsi memori heap yang signifikan. Anda memiliki dua solusi:
-
Kurangi nilai
orc.write.batch-size. -
Jika Anda tidak perlu melakukan kueri ad-hoc (OLAP) pada tabel hasil Paimon dan hanya memerlukan konsumsi batch atau streaming, Anda dapat mengatur properti tabel
'file.format' = 'avro'dan'metadata.stats-mode' = 'none'saat pembuatan tabel. Hal ini mengganti tabel ke format Avro dan menonaktifkan pengumpulan statistik.CatatanParameter ini hanya dapat diatur saat pembuatan tabel. Parameter ini tidak dapat diubah dengan pernyataan ALTER TABLE atau petunjuk SQL setelah tabel dibuat.
-
-
Menulis ke terlalu banyak partisi secara bersamaan atau memiliki terlalu banyak bucket per partisi menciptakan jumlah objek penulis yang berlebihan.
Tinjau konfigurasi kolom partisi Anda untuk memastikan sesuai. Pastikan SQL yang salah tidak menyebabkan data tak terduga ditulis ke kolom partisi. Juga, verifikasi bahwa jumlah bucket masuk akal. Sebagai praktik terbaik, ukuran data total per bucket harus sekitar 2 GB dan tidak boleh melebihi 5 GB. Untuk detail tentang menyesuaikan jumlah bucket, lihat Sesuaikan jumlah bucket untuk tabel fixed-bucket.
Error: "Sink materializer must not be used with Paimon sink"
Operator sink materializer menangani data out-of-order dari join cascading dalam job streaming. Namun, dalam job yang menulis ke tabel Paimon, operator ini menimbulkan overhead dan dapat menyebabkan hasil yang salah saat agregasi digunakan. Jangan gunakan operator sink materializer dengan sink Paimon.
Anda dapat menonaktifkan operator sink materializer dengan mengatur parameter table.exec.sink.upsert-materialize ke false menggunakan pernyataan SET atau sebagai parameter runtime. Jika Anda juga perlu menangani data out-of-order, lihat penanganan data out-of-order.
Paimon: error File deletion conflicts detected atau LSM conflicts detected
Error ini dapat terjadi karena alasan-alasan berikut:
-
Beberapa job menulis ke partisi yang sama dari tabel Paimon yang sama. Dalam kasus ini, Paimon menyelesaikan konflik melalui
failover and restart. Ini adalah perilaku yang diharapkan, dan tidak diperlukan tindakan jika error tidak berulang. -
Job dipulihkan dari state yang usang, yang menyebabkan error berulang. Untuk menyelesaikannya, pulihkan job dari state terbarunya atau
mulai tanpa state. -
Paimon tidak mendukung penulisan terpisah dari beberapa pernyataan
INSERTdalam satu job. Sebagai gantinya, gunakan pernyataanUNION ALLuntuk menulis beberapa aliran data ke tabel Paimon. -
Konkurensi node Global Committer atau node Compaction Coordinator (saat menulis ke tabel Append Scalable) lebih besar dari 1.
konkurensiuntuk node-node ini harus 1 untuk memastikankonsistensi data.
Error "File xxx not found" dalam job konsumsi Paimon
Konsumsi tabel Paimon bergantung pada file snapshot. Jika periode retensi snapshot terlalu singkat atau job konsumsi tidak efisien, file snapshot dapat kedaluwarsa dan dihapus sebelum job selesai. Hal ini menyebabkan job konsumsi gagal.
Untuk menyelesaikan masalah ini, Anda dapat menyesuaikan periode retensi file snapshot, menentukan ID konsumen, atau mengoptimalkan job konsumsi. Untuk memeriksa file snapshot yang tersedia dan timestamp pembuatannya, lihat Tabel sistem Snapshots.
Error job Paimon: No space left on device
-
Jumlah file cache yang berlebihan dapat menyebabkan error ini jika job Anda menjalankan kueri Paimon, seperti menggunakan tabel Paimon sebagai tabel dimensi atau mengatur
changelog-producer='lookup'. Untuk mencegahnya, gunakan Petunjuk SQL untuk mengatur parameter berikut guna membatasi ruang disk maksimum dan periode retensi cache kueri.-
lookup.cache-max-disk-size: Ruang disk lokal maksimum yang dapat digunakan oleh cache kueri. Nilai yang direkomendasikan termasuk 256 MB, 512 MB, dan 1 GB. -
lookup.cache-file-retention: Periode retensi untuk file cache kueri. Nilai yang direkomendasikan termasuk 30 menit, 15 menit, atau interval yang lebih pendek.
-
-
Untuk job yang menulis ke tabel Paimon, gunakan Petunjuk SQL untuk mengatur parameter berikut. Pengaturan ini membatasi ukuran file sementara lokal selama proses penulisan, mencegah kekurangan ruang disk.
-
write-buffer-spillable: Mengontrol apakah buffer penulisan dapat tumpah ke disk. Mengatur ini kefalsesepenuhnya mencegah buffer menggunakan ruang disk apa pun. -
write-buffer-spill.max-disk-size: Ruang disk maksimum yang dapat digunakan oleh buffer penulisan saat tumpah. Nilai yang direkomendasikan termasuk 256 MB, 512 MB, dan 1 GB.
-
Mengelola file Paimon di OSS
-
Paimon menyimpan file data historis untuk mengakses versi sebelumnya dari tabel. Anda dapat menyesuaikan kebijakan retensi untuk file-file ini guna mengelola penyimpanan. Untuk instruksi terperinci, lihat Bersihkan data yang kedaluwarsa.
-
Konfigurasi kolom partisi yang tidak tepat atau jumlah bucket yang berlebihan juga dapat menyebabkan masalah ini. Sebagai praktik terbaik, bidik ukuran data sekitar 2 GB per bucket, dengan maksimum 5 GB. Untuk informasi lebih lanjut, lihat Bucketing.
-
Secara default, file data disimpan dalam format ORC. Untuk mengurangi ukuran total file data, Anda dapat menggunakan format kompresi ZSTD dengan mengatur parameter tabel
'file.compression' = 'zstd'saat Anda membuat tabel.CatatanParameter ini hanya dapat diatur saat pembuatan tabel dan tidak dapat diubah nanti dengan pernyataan ALTER TABLE atau petunjuk SQL.
Visibilitas data bergantung pada interval checkpoint
Ya. Paimon bergantung pada checkpoint untuk menjamin semantik exactly-once. Data dikomit dan menjadi terlihat di downstream hanya setelah checkpoint selesai. Sebelum komit ini, data dalam buffer lokal disiram ke sistem file jarak jauh, tetapi belum dapat dibaca.
Peningkatan memori perlahan dalam job Paimon jangka panjang
-
Peningkatan penggunaan memori diharapkan jika rps job juga meningkat perlahan.
-
Jika Anda menggunakan katalog filesystem Paimon untuk membaca dari atau menulis ke OSS, pastikan untuk mengonfigurasi parameter katalog
fs.oss.endpoint,fs.oss.accessKeyId, danfs.oss.accessKeySecret. Jika tidak, job Flink dapat mengalami kebocoran memori perlahan, yang merupakan masalah yang diketahui di komunitas.
IllegalArgumentException: timeout value is negative
-
Detail error
2021-02-24 15:14:58 java.lang.RuntimeException: java.lang.RuntimeException: java.lang.IllegalArgumentException: timeout value is negative at com.alibaba.ververica.connectors.common.source.reader.ParallelReader.run(ParallelReader.java:166) at com.alibaba.ververica.connectors.common.source.AbstractParallelSourceBase.run(AbstractParallelSourceBase.java:205) at com.alibaba.ververica.connectors.metaq.source.MetaQRowDataSource.run(MetaQRowDataSource.java:84) at org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:100) at org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:63) at org.apache.flink.streaming.runtime.tasks.SourceStreamTask$LegacySourceFunctionThread.run(SourceStreamTask.java:213) Caused by: java.lang.RuntimeException: java.lang.IllegalArgumentException: timeout value is negative at com.alibaba.ververica.connectors.common.source.reader.ParallelReader.runImpl(ParallelReader.java:244) -
Penyebab error
Saat tidak ada pesan MQ baru yang dikonsumsi, thread
MetaQSourcetidur selama interval yang ditentukan oleh parameterpullIntervalMs, yang default-nya -1. Job kemudian gagal denganIllegalArgumentExceptionkarena durasi tidur tidak dapat negatif. -
Solusi
Atur parameter
pullIntervalMske nilai non-negatif.
Deteksi perubahan partisi
-
Untuk versi Realtime Compute for Apache Flink sebelum 6.0.2, operator sumber mengambil jumlah partisi setiap 5 hingga 10 menit. Failover dipicu jika jumlah partisi berbeda untuk tiga pemeriksaan berturut-turut. Akibatnya, sumber memulai failover dalam 10 hingga 30 menit. Setelah job dimulai ulang, job tersebut membaca dari set partisi yang diperbarui.
-
Untuk versi Realtime Compute for Apache Flink 6.0.2 dan yang lebih baru, operator sumber mengambil jumlah partisi setiap 5 menit secara default. Saat partisi baru ditemukan, partisi tersebut langsung ditugaskan ke operator sumber di TaskManager, yang kemudian mulai membaca data. Tidak diperlukan failover job, memungkinkan sumber mendeteksi perubahan partisi dalam 1 hingga 5 menit.
Error: Backpressure exceeds reject limit
-
Detail error
26 at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:47) 27 at org.apache.flink.streaming.runtime.tasks.StreamTask.performCheckpoint(StreamTask.java:911) 28 at org.apache.flink.streaming.runtime.tasks.StreamTask.triggerCheckpointOnBarrier(StreamTask.java:879) 29 ... 13 more 30 Caused by: java.lang.RuntimeException: Rpc Exception failed errorCount=6 with RpcException: request niagara.table.proto.UpsertRecordBatchRequest@b53e2d74 failed on final try 4, maxAttempts=4, sn=11.117.xxx, errorCode=11, msg=BackPresure Exceed Reject Limit [method:UpsertRecordBatch,transaction_id:xxx,table_id:xxx,table_version:128, actor_id:74538 xxx,worker_address:11.117.xxx] 31 at com.alibaba.ververica.connectors.hologres.sink.HologresOutputFormat.sync(HologresOutputFormat.java:264) 32 at com.alibaba.ververica.connectors.common.sink.OutputFormatSinkFunction.snapshotState(OutputFormatSinkFunction.java:91) 33 at org.apache.flink.streaming.util.functions.StreamingFunctionUtils.trySnapshotFunctionState(StreamingFunctionUtils.java:128) 34 at org.apache.flink.streaming.util.functions.StreamingFunctionUtils.snapshotFunctionState(StreamingFunctionUtils.java:101) 35 at org.apache.flink.streaming.api.operators.AbstractUdfStreamOperator.snapshotState(AbstractUdfStreamOperator.java:90) 36 at org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.snapshotState(StreamOperatorStateHandler.java:186) 37 ... 23 more -
Penyebab
Tekanan tulis pada instans Hologres terlalu tinggi.
-
Solusi
Hubungi dukungan teknis Hologres dengan informasi instans Anda untuk meminta peningkatan.
Error: remaining connection slots are reserved for non-replication superuser connections
-
Detail Kesalahan
Caused by: com.alibaba.hologres.client.exception.HoloClientWithDetailsException: failed records 1, first:Record{schema=org.postgresql.model.TableSchema@188365, values=[f06b41455c694d24a18d0552b8b0****, com.chot.tpfymnq.meta, 2022-04-02 19:46:40.0, 28, 1, null], bitSet={0, 1, 2, 3, 4}},first err:[106]FATAL: remaining connection slots are reserved for non-replication superuser connections at com.alibaba.hologres.client.impl.Worker.handlePutAction(Worker.java:406) ~[?:?] at com.alibaba.hologres.client.impl.Worker.run(Worker.java:118) ~[?:?] at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) ~[?:1.8.0_302] at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) ~[?:1.8.0_302] ... 1 more Caused by: com.alibaba.hologres.org.postgresql.util.PSQLException: FATAL: remaining connection slots are reserved for non-replication superuser connections at com.alibaba.hologres.org.postgresql.core.v3.QueryExecutorImpl.receiveErrorResponse(QueryExecutorImpl.java:2553) ~[?:?] at com.alibaba.hologres.org.postgresql.core.v3.QueryExecutorImpl.readStartupMessages(QueryExecutorImpl.java:2665) ~[?:?] at com.alibaba.hologres.org.postgresql.core.v3.QueryExecutorImpl.<init>(QueryExecutorImpl.java:147) ~[?:?] at com.alibaba.hologres.org.postgresql.core.v3.ConnectionFactoryImpl.openConnectionImpl(ConnectionFactoryImpl.java:273) ~[?:?] at com.alibaba.hologres.org.postgresql.core.ConnectionFactory.openConnection(ConnectionFactory.java:51) ~[?:?] at com.alibaba.hologres.org.postgresql.jdbc.PgConnection.<init>(PgConnection.java:240) ~[?:?] at com.alibaba.hologres.org.postgresql.Driver.makeConnection(Driver.java:478) ~[?:?] at com.alibaba.hologres.org.postgresql.Driver.connect(Driver.java:277) ~[?:?] at java.sql.DriverManager.getConnection(DriverManager.java:674) ~[?:1.8.0_302] at java.sql.DriverManager.getConnection(DriverManager.java:217) ~[?:1.8.0_302] at com.alibaba.hologres.client.impl.ConnectionHolder.buildConnection(ConnectionHolder.java:122) ~[?:?] at com.alibaba.hologres.client.impl.ConnectionHolder.retryExecute(ConnectionHolder.java:195) ~[?:?] at com.alibaba.hologres.client.impl.ConnectionHolder.retryExecute(ConnectionHolder.java:184) ~[?:?] at com.alibaba.hologres.client.impl.Worker.doHandlePutAction(Worker.java:460) ~[?:?] at com.alibaba.hologres.client.impl.Worker.handlePutAction(Worker.java:389) ~[?:?] at com.alibaba.hologres.client.impl.Worker.run(Worker.java:118) ~[?:?] at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) ~[?:1.8.0_302] at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) ~[?:1.8.0_302] ... 1 more -
Penyebab
Batas koneksi untuk instans Hologres telah terlampaui.
-
Solusi
-
Periksa
app_namekoneksi ke setiap Frontend (FE) untuk menghitung koneksi klien Hologres dariflink-connector. -
Periksa job lain yang terhubung ke Hologres.
-
Lepaskan koneksi. Untuk informasi lebih lanjut, lihat manajemen koneksi.
-
No table is defined in publication
-
Detail error
Menghapus dan membuat ulang tabel dengan nama yang sama dapat menyebabkan job melaporkan
no table is defined in publication. -
Penyebab
Menghapus tabel tidak menghapus publikasi yang terkait.
-
Solusi
-
Di Hologres, jalankan perintah
select * from pg_publication where pubname not in (select pubname from pg_publication_tables);untuk menanyakan informasi publikasi yang tidak dibersihkan saat tabel dihapus. -
Jalankan pernyataan
drop publication xx;untuk menghapus publikasi yang tersisa. -
Mulai ulang job.
-
Interval checkpoint dan visibilitas data
Interval checkpoint konektor sink Flink Hologres tidak secara langsung mengontrol visibilitas data di Hologres. Peran utamanya adalah menentukan SLA untuk pemulihan kegagalan.
Konektor Hologres tidak mendukung transaksi. Konektor secara berkala menyiram buffer dalam memori ke database. Checkpoint memastikan bahwa semua data disiram pada saat selesai, tetapi konektor tidak menunggu seluruh interval berlalu sebelum menyiram. Konektor memicu penyiraman lebih awal jika kondisi buffer tertentu terpenuhi (untuk informasi lebih lanjut, lihat Hologres, Hologres, dan Hologres). Karena gudang data biasanya tidak memerlukan konsistensi transaksional, konektor menyiram data secara asinkron di latar belakang. Konektor kemudian melakukan penyiraman paksa terakhir selama setiap checkpoint untuk mempersiapkan pemulihan kegagalan.
Penerapan Job menghasilkan permission denied for database galat
-
Penyebab
Mulai dari Realtime Compute for Apache Flink VVR 8.0.4, konektor menerapkan mode JDBC untuk mengonsumsi log biner dari instans Hologres V2.0 atau yang lebih baru. Untuk instans ini, akun non-superuser memerlukan izin khusus untuk mengonsumsi log biner dalam mode JDBC.
-
Solusi
Berikan izin kepada akun non-superuser untuk mengonsumsi log biner dalam mode JDBC.
user_name mengacu pada ID akun Alibaba Cloud atau pengguna RAM. Untuk informasi lebih lanjut, lihat gambaran akun.
-- Untuk model izin ahli, berikan izin CREATE dan peran replikasi kepada pengguna. GRANT CREATE ON DATABASE <db_name> TO <user_name>; alter role <user_name> replication; -- Jika database menggunakan model izin sederhana (SPM), Anda tidak dapat menjalankan pernyataan GRANT. -- Sebagai gantinya, gunakan spm_grant untuk memberikan peran Admin kepada pengguna untuk database tersebut. Anda juga dapat memberikan izin langsung di HoloWeb. call spm_grant('<db_name>_admin', '<user_name>'); alter role <user_name> replication;
Kegagalan pemulihan job: table id parsed from checkpoint is different from the current table id
-
Penyebab
Pengecualian ini terjadi di versi Realtime Compute for Apache Flink VVR 8.0.5 hingga VVR 8.0.8. Saat job yang menggunakan tabel sumber binlog Hologres dipulihkan dari checkpoint, engine menerapkan pemeriksaan ketat pada table id. Jika table id saat ini dari tabel Hologres tidak cocok dengan yang disimpan di checkpoint, pemulihan gagal. Hal ini menunjukkan bahwa tabel sumber dipotong atau dibuat ulang saat job berjalan.
-
Solusi
Tingkatkan ke VVR 8.0.9 atau versi yang lebih baru dan mulai ulang job. VVR 8.0.9 menghapus pemeriksaan ketat pada table id untuk mengakomodasi skenario bisnis yang kompleks. Namun, hindari membuat ulang tabel sumber binlog. Saat tabel dibuat ulang, seluruh riwayat binlog-nya dihapus. Jika Flink kemudian menggunakan offset konsumen dari tabel lama untuk membaca data dari tabel baru, hal ini dapat menyebabkan ketidaksesuaian data.
Presisi data Binlog yang tidak terduga dalam mode JDBC
-
Penyebab
Di Realtime Compute for Apache Flink 8.0.10 dan yang lebih awal, presisi data yang tidak terduga terjadi jika presisi tipe
DECIMALyang dideklarasikan dalam DDL Flink untuk tabel sumber Binlog tidak cocok dengan presisi di Hologres. -
Solusi
Masalah ini diperbaiki di Realtime Compute for Apache Flink 8.0.11. Namun, pastikan presisi tipe
DECIMALkonsisten antara Flink dan Hologres untuk mencegah kehilangan presisi.
Menghapus dan membuat ulang tabel dengan nama yang sama dapat menyebabkan job melaporkan pengecualian no table is defined in publication atau The table xxx has no slot named xxx
-
Penyebab
Hal ini terjadi karena menghapus tabel tidak menghapus publikasi yang terkait.
-
Solusi
Solusi 1: Di Hologres, jalankan pernyataan
select * from pg_publication where pubname not in (select pubname from pg_publication_tables);untuk menemukan publikasi yang tersisa dari tabel yang dihapus. Kemudian, jalankan pernyataandrop publication xx;untuk menghapusnya. Terakhir, mulai ulang job.Solusi 2: Gunakan VVR 8.0.5 atau yang lebih baru. Konektor secara otomatis menangani pembersihan.
ClassCastException saat membaca dari Hologres
-
Detail error
Pesan error mirip dengan berikut:
java.lang.ClassCastException: class java.lang.Integer cannot be cast to class java.lang.Long (java.lang.Integer and java.lang.Long are in module java.base of loader 'bootstrap') -
Penyebab
Error ini terjadi saat tipe field dalam DDL Flink tidak cocok dengan tipe field yang sesuai dalam tabel fisik Hologres. Misalnya, field didefinisikan sebagai BIGINT dalam DDL Flink, tetapi field yang sesuai dalam tabel Hologres adalah INTEGER. Karena pemeriksaan tipe dilewati untuk nilai NULL, job melempar pengecualian hanya saat membaca data aktual.
-
Solusi
Lihat dokumentasi Hologres Ringkasan tipe data untuk memastikan tipe field dalam DDL Flink Anda cocok dengan tabel fisik Hologres.
LogSizeTooLargeException
-
Detail error
Disebabkan oleh: com.aliyun.openservices.aliyun.log.producer.errors.LogSizeTooLargeException: log berukuran 8785684 byte yang lebih besar dari MAX_BATCH_SIZE_IN_BYTES 8388608 at com.aliyun.openservices.aliyun.log.producer.internals.LogAccumulator.ensureValidLogSize(LogAccumulator.java:249) at com.aliyun.openservices.aliyun.log.producer.internals.LogAccumulator.doAppend(LogAccumulator.java:103) at com.aliyun.openservices.aliyun.log.producer.internals.LogAccumulator.append(LogAccumulator.java:84) at com.aliyun.openservices.aliyun.log.producer.LogProducer.send(LogProducer.java:385) at com.aliyun.openservices.aliyun.log.producer.LogProducer.send(LogProducer.java:308) at com.aliyun.openservices.aliyun.log.producer.LogProducer.send(LogProducer.java:211) at com.alibaba.ververica.connectors.sls.sink.SLSOutputFormat.writeRecord(SLSOutputFo rmat.java:100) -
Penyebab
Error ini terjadi karena log satu baris yang dikirim ke Log Service melebihi batas ukuran 8 MB.
-
Solusi
Untuk melewati entri log yang terlalu besar, ubah posisi startup. Untuk detailnya, lihat startup job.
OOM TaskManager: error Java heap space selama pemulihan
-
Penyebab
Masalah ini biasanya disebabkan oleh badan pesan SLS yang terlalu besar. Konektor SLS meminta data dalam batch. Jumlah LogGroups per batch dikontrol oleh parameter batchGetSize, yang default-nya 100. Artinya, setiap permintaan dapat mengambil hingga 100 LogGroups. Selama operasi normal, program Flink segera mengonsumsi data dan jarang mengambil batch penuh 100 LogGroups. Namun, selama failover, sejumlah besar data yang belum dikonsumsi dapat menumpuk. Jika memori yang diperlukan untuk batch penuh 100 LogGroups melebihi memori yang tersedia JVM, TaskManager mengalami OOM.
-
Solusi
Kurangi nilai parameter batchGetSize.
Atur offset konsumsi untuk tabel sumber Paimon
Untuk mengatur offset konsumsi untuk tabel sumber Paimon, gunakan parameter scan.mode. Tabel berikut menjelaskan nilai yang tersedia dan perilakunya.
|
Nilai |
Perilaku pembacaan batch |
Perilaku pembacaan stream |
|
default |
Nilai default. Perilaku aktual tergantung pada parameter lain.
Jika tidak ada parameter yang diatur, perilakunya sama dengan |
|
|
latest-full |
Membaca snapshot terbaru dari tabel. |
Saat job dimulai, job tersebut pertama-tama membaca snapshot terbaru dari tabel lalu terus membaca data inkremental. |
|
compacted-full |
Membaca snapshot terbaru dari tabel setelah kompaksi terbaru. |
Saat job dimulai, job tersebut pertama-tama membaca snapshot terbaru dari tabel setelah kompaksi terbaru, lalu terus membaca data inkremental. |
|
latest |
Sama dengan |
Saat job dimulai, job tersebut melewatkan snapshot terbaru dan malah terus membaca data inkremental. |
|
from-timestamp |
Menampilkan tabel dari snapshot terbaru pada atau sebelum |
Job tidak menghasilkan snapshot saat startup dan terus menghasilkan data inkremental mulai dari (dan termasuk) |
|
from-snapshot |
Menghasilkan snapshot dari tabel. ID snapshot ditentukan oleh |
Job tidak menghasilkan snapshot saat startup. Job tersebut kemudian terus menghasilkan data inkremental mulai dari dan termasuk |
|
from-snapshot-full |
Sama dengan |
Saat startup job, snapshot dari tabel dihasilkan. ID snapshot ditentukan oleh |
Mengonfigurasi kedaluwarsa partisi otomatis
Tabel Paimon dapat secara otomatis menghapus partisi yang masa hidup datanya melebihi waktu kedaluwarsa partisi yang ditentukan. Fitur ini membantu mengurangi biaya penyimpanan. Prosesnya sebagai berikut:
-
Masa hidup data: Waktu sistem saat ini dikurangi timestamp yang berasal dari nilai partisi. Nilai partisi dikonversi ke timestamp sebagai berikut:
-
Konversi nilai partisi ke string waktu menggunakan string format yang ditentukan oleh parameter
partition.timestamp-pattern.Dalam string format ini, kolom partisi direpresentasikan oleh tanda dolar ($) diikuti nama kolom. Misalnya, jika kolom partisi adalah year, month, day, dan hour, string format
$year-$month-$day $hour:00:00mengonversi partisiyear=2023,month=04,day=21,hour=17ke string2023-04-21 17:00:00. -
Konversi string waktu ke timestamp menggunakan string format yang ditentukan oleh parameter
partition.timestamp-formatter.Jika parameter ini tidak diatur, sistem default ke format
yyyy-MM-dd HH:mm:ssdanyyyy-MM-dd. Anda dapat menggunakan string format apa pun yang kompatibel dengan DateTimeFormatter Java.
-
-
Waktu kedaluwarsa partisi: Nilai yang Anda konfigurasikan untuk parameter
partition.expiration-time.
Pemecahan masalah data yang hilang di penyimpanan
-
Data mungkin tidak langsung terlihat di penyimpanan. Penulis Flink menyiram data ke disk dalam kondisi berikut:
-
Data yang dibuffer dalam bucket mencapai ukuran tertentu (default: 64 MB).
-
Ukuran buffer total mencapai ambang batas (default: 1 GB).
-
Checkpoint dipicu, yang menyiram semua data dalam memori.
-
-
Jika Anda menggunakan penulisan stream, pastikan checkpointing diaktifkan.
Menangani data duplikat di Hudi
-
Untuk penulisan COW, aktifkan parameter
write.insert.drop.duplicates.Secara default, penulisan COW tidak mendeduplikasi data dalam file pertama setiap bucket dan hanya menerapkan deduplikasi ke data inkremental. Untuk melakukan deduplikasi global, Anda harus mengaktifkan parameter ini. Untuk penulisan MOR, tidak diperlukan parameter tambahan. Mendefinisikan primary key mengaktifkan deduplikasi global secara default.
CatatanSejak Hudi versi 0.10.0, properti ini diganti namanya menjadi
write.precombinedan diatur ketruesecara default. -
Untuk melakukan deduplikasi di beberapa partisi, atur parameter
index.global.enabledketrue.Catatan-
Sejak Hudi versi 0.10.0, properti ini diatur ke
truesecara default. -
Saat
index.type=bucket, mengatur parameterindex.global.enabledketruetidak efektif karena indeks Bucket tidak mendukung perubahan lintas partisi. Oleh karena itu, meskipun indeks global diaktifkan, fitur deduplikasi tidak dapat bekerja di beberapa partisi.
-
-
Untuk pembaruan jendela panjang, seperti memodifikasi data dari sebulan yang lalu, tingkatkan parameter
index.state.ttl, yang diukur dalam hari.Indeks adalah struktur data inti Hudi untuk mengidentifikasi data duplikat. Parameter
index.state.ttlmengontrol berapa lama status indeks bertahan. Nilai default sebelumnya adalah 1,5 hari. Nilai 0 atau kurang menunjukkan bahwa status indeks dipertahankan secara permanen.CatatanSejak Hudi versi 0.10.0, properti ini default-nya
0.
Merge On Read hanya memiliki file log
-
Penyebab: Hudi membuat file Parquet hanya setelah kompaksi; jika tidak, Hudi hanya membuat file log. Secara default, tabel Merge On Read menggunakan kompaksi asinkron, yang memicu job kompaksi setelah setiap lima commit.
-
Solusi: Sesuaikan parameter
compaction.delta_commitsuntuk memicu job kompaksi lebih cepat.
Error: "multi-statement be found"
-
Masalah
Job Flink yang menulis data ke instans AnalyticDB for MySQL (ADB) gagal dan dimulai ulang. Log menunjukkan error serupa berikut:
Caused by: java.sql.SQLSyntaxErrorException: [13000, 2024101216171419216823505703151806929] multi-statement be found.at java.util.TimerThread.run(Timer.java:505) Caused by: java.sql.BatchUpdateException: [13000, 2024101216400819216823505703151079281] multi-statement be found. at sun.reflect.GeneratedConstructorAccessor115.newInstance(Unknown Source) at sun.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45) at java.lang.reflect.Constructor.newInstance(Constructor.java:423) at com.mysql.cj.util.Util.handleNewInstance(Util.java:192) at com.mysql.cj.util.Util.getInstance(Util.java:167) at com.mysql.cj.util.Util.getInstance(Util.java:174) at com.mysql.cj.jdbc.exceptions.SQLError.createBatchUpdateException(SQLError.java:224) at com.mysql.cj.jdbc.ClientPreparedStatement.executePreparedBatchAsMultiStatement(ClientPreparedStatement.java:584) at com.mysql.cj.jdbc.ClientPreparedStatement.executeBatchInternal(ClientPreparedStatement.java:431) at com.mysql.cj.jdbc.StatementImpl.executeBatch(StatementImpl.java:795) at com.ververica.cdc.connectors.shaded.com.zaxxer.hikari.pool.ProxyStatement.executeBatch(ProxyStatement.java:127) at com.ververica.cdc.connectors.shaded.com.zaxxer.hikari.pool.HikariProxyPreparedStatement.executeBatch(HikariProxyPreparedStatement...) at com.ververica.connectors.mysql.table.sink.MySqlOutputFormat.executeSql(MySqlOutputFormat.java:567) ... 6 more -
Penyebab
Error ini menunjukkan masalah kompatibilitas antara driver JDBC MySQL versi 8.x dan database AnalyticDB for MySQL (ADB) dengan
ALLOW_MULTI_QUERIES=truediaktifkan. -
Solusi
-
Hubungi dukungan teknis untuk mendapatkan konektor ADB 3.0 kustom yang menggunakan driver JDBC MySQL versi 5.1.46. Terapkan konektor ini ke tugas Flink Anda. Lihat Mengelola konektor kustom untuk instruksi penggunaan konektor kustom.
-
Atur parameter
allowMultiQueries=truedalam URI tabel ADB, misalnya,jdbc:mysql://xxxxx.ads.aliyuncs.com:3306/xxx?allowMultiQueries=true.
-
Error: No suitable driver found
-
Penyebab
Konektor kustom tidak dapat menemukan driver yang diperlukan.
-
Solusi
-
Muat driver di kelas factory dengan memanggil
Class.forName. -
Tambahkan driver sebagai dependensi tambahan dan atur parameter
kubernetes.application-mode.classpath.include-user-jarke true. Lihat Bagaimana cara mengonfigurasi parameter run job kustom? untuk instruksi.
-
Kehilangan data atau penimpaan saat Flink menulis ke Elasticsearch
-
Penyebab 1: Konflik antara
doc_as_upsertdan Pipeline Ingest ElasticsearchSaat keduanya dikonfigurasi secara bersamaan, Elasticsearch dapat memproses pembaruan parsial dan transformasi pipeline dalam urutan yang tidak kompatibel, menyebabkan dokumen ditimpa atau hilang secara tak terduga:
-
sink.bulk-flush.update.doc_as_upsert = 'true'dalam klausaWITHDDL Flink -
Pipeline Ingest Elasticsearch yang ditetapkan ke indeks target
Dalam konfigurasi ini, Elasticsearch menerapkan Pipeline Ingest sebelum memproses pembaruan parsial. Bergantung pada logika pipeline dan versi Elasticsearch, interaksi ini dapat menyebabkan field dokumen ditimpa dengan nilai yang salah atau dokumen dijatuhkan diam-diam.
-
-
Penyebab 2: Beberapa tugas sink Flink menulis ke indeks yang sama
Jika dua atau lebih job Flink — atau beberapa instance sink paralel dengan rentang kunci yang tidak tumpang tindih — menulis ke indeks Elasticsearch yang sama tanpa penanganan kunci yang terkoordinasi, dokumen dapat saling menimpa. Hal ini mengakibatkan kehilangan data untuk sink mana pun yang menulis terakhir ke ID dokumen tertentu.
-
Solusi
-
Untuk konflik
doc_as_upsertdan Pipeline Ingest: Hapus parametersink.bulk-flush.update.doc_as_upsert = 'true'dari DDL Flink, dan hapus atau tetapkan ulang Pipeline Ingest Elasticsearch dari indeks target. Pindahkan logika transformasi data yang ditangani oleh pipeline ke dalam job Flink itu sendiri — misalnya, menggunakanProcessFunctionatau kolom terhitung sebelum operator sink. -
Untuk beberapa tugas sink yang saling menimpa: Pastikan setiap job Flink menulis ke indeks Elasticsearch yang terpisah, atau konsolidasikan penulisan ke satu job Flink yang mengelola kunci dokumen secara deterministik.
-