Error runtime pekerjaan umum dan solusinya di Realtime Compute for Apache Flink.
-
Apa yang harus saya lakukan jika pekerjaan tidak dapat dimulai?
-
Apa yang harus saya lakukan jika pekerjaan restart secara tak terduga?
-
Mengapa output data terhenti pada operator LocalGroupAggregate?
-
Apa yang harus saya lakukan jika ketidakaktifan partisi Kafka menunda output window?
-
Bagaimana cara menemukan error jika JobManager tidak berjalan?
-
Apa yang harus saya lakukan jika muncul pesan error "akka.pattern.AskTimeoutException"?
-
Bagaimana cara memperbaiki pesan error "The GRPC call timed out in sqlserver"?
-
Apa yang harus saya lakukan jika muncul pesan error "Caused by: java.lang.NoSuchMethodError"?
Apa yang harus saya lakukan jika pekerjaan tidak dapat dimulai?
-
Deskripsi masalah
Saat Anda mengklik Start di kolom Actions, status pekerjaan berubah dari STARTING menjadi FAILED.
-
Solusi
-
Periksa tab Events: Buka tab Events di halaman detail pekerjaan. Temukan event kegagalan yang terjadi saat pekerjaan gagal dimulai dan tinjau detailnya untuk mengidentifikasi akar permasalahan.
-
Tinjau log startup: Buka tab Logs dan pilih sub-tab Startup Logs. Periksa log tersebut untuk menemukan pesan error spesifik yang menjelaskan mengapa pekerjaan gagal dimulai.
-
Periksa log JobManager/TaskManager: Jika JobManager tampaknya berhasil dimulai tetapi pekerjaan tetap gagal, periksa log detail untuk JobManager dan TaskManager. Log tersebut dapat ditemukan di sub-tab Job Manager atau Running Task Managers dalam tab Logs.
-
-
Error umum dan solusinya
Deskripsi masalah
Penyebab
Solusi
ERROR:exceeded quota: resourcequotaSumber daya tidak mencukupi di Resource Queue saat ini.
Tingkatkan kapasitas Resource Queue atau kurangi kebutuhan sumber daya pekerjaan.
ERROR:the vswitch ip is not enoughAlamat IP tidak mencukupi di namespace untuk TaskManager yang diperlukan.
Kurangi paralelisme pekerjaan, sesuaikan konfigurasi slot, atau ubah pengaturan vSwitch.
ERROR: pooler: ***: authentication failedPasangan Kunci Akses tidak valid atau izin tidak mencukupi.
Verifikasi bahwa Pasangan Kunci Akses valid dan dimiliki oleh akun yang memiliki izin untuk mengeksekusi dan mengelola pekerjaan.
Bagaimana cara memperbaiki error koneksi database?
-
Deskripsi masalah
failed to execute sql statement -
Penyebab
Katalog yang terdaftar tidak valid atau tidak dapat dijangkau.
-
Solusi
Buka halaman Catalogs, hapus semua katalog yang berwarna abu-abu, lalu daftarkan ulang.
Apa yang harus saya lakukan jika data dalam task pekerjaan tidak dikonsumsi setelah pekerjaan dijalankan?
-
Periksa konektivitas jaringan
Jika data tidak dihasilkan atau dikonsumsi di penyimpanan hulu dan hilir, periksa tab Startup Logs untuk pesan error. Jika Anda melihat error timeout, pecahkan masalah konektivitas jaringan antara sistem penyimpanan.
-
Periksa status eksekusi task
Di tab Status, verifikasi apakah data sedang dibaca dari source dan ditulis ke sink untuk menentukan lokasi error.
Di tabel metrik, periksa kolom Bytes Sent dan Records sent untuk node Source. Jika nilainya lebih besar dari 0 (misalnya, 19,74 GB / 76.766.861 catatan), berarti source mengirim data secara normal. Demikian pula, periksa kolom Bytes Received dan Records Received untuk node hilir guna memastikan mereka menerima data secara normal.
-
Periksa output operator
Tambahkan tabel sink print ke setiap operator untuk memecahkan masalah ini.
Apa yang harus saya lakukan jika pekerjaan restart secara tak terduga?
Untuk memecahkan masalah error tersebut, periksa tab Logs.
-
Lihat informasi exception.
Di sub-tab JM Exceptions, tinjau error yang dilaporkan dan identifikasi akar permasalahannya.
-
Lihat log JobManager dan TaskManager untuk pekerjaan tersebut.
Di tab Logs, klik sub-tab Job Manager, lalu pilih sub-tab Logs untuk melihat log pekerjaan yang sesuai. Demikian pula, klik sub-tab Running Task Managers untuk melihat log TM.
-
Lihat log TaskManager yang gagal untuk pekerjaan tersebut.
Beberapa exception dapat menyebabkan TaskManager gagal, sehingga log menjadi tidak lengkap. Lihat log TaskManager tidak valid terakhir untuk memecahkan masalah ini.
-
Lihat log instans pekerjaan historis.
Tinjau log instans pekerjaan historis untuk mengidentifikasi penyebab kegagalan.
Mengapa output data terhenti pada operator LocalGroupAggregate?
-
Kode
CREATE TEMPORARY TABLE s1 ( a INT, b INT, ts as PROCTIME(), PRIMARY KEY (a) NOT ENFORCED ) WITH ( 'connector'='datagen', 'rows-per-second'='1', 'fields.b.kind'='random', 'fields.b.min'='0', 'fields.b.max'='10' ); CREATE TEMPORARY TABLE sink ( a BIGINT, b BIGINT ) WITH ( 'connector'='print' ); CREATE TEMPORARY VIEW window_view AS SELECT window_start, window_end, a, sum(b) as b_sum FROM TABLE(TUMBLE(TABLE s1, DESCRIPTOR(ts), INTERVAL '2' SECONDS)) GROUP BY window_start, window_end, a; INSERT INTO sink SELECT count(distinct a), b_sum FROM window_view GROUP BY b_sum; -
Deskripsi masalah
Output data terhenti pada operator LocalGroupAggregate dalam waktu yang lama, dan operator MiniBatchAssigner tidak muncul dalam topologi pekerjaan.
Di antarmuka Operator Analysis (Beta), topologi pekerjaan menunjukkan bahwa LocalGroupAggregate[201] memiliki RecordsIn sebesar 1.853 dan RecordsOut hanya 116. Vertex2 hilir (berisi operator GlobalGroupAggregate dan Calc) memiliki Out (sum) sebesar 0 dan Backpressured (max) sebesar 0%. Tidak ada node MiniBatchAssigner yang muncul di seluruh topologi.
-
Penyebab
Pekerjaan mencakup operator WindowAggregate dan GroupAggregate. Operator WindowAggregate menggunakan proctime sebagai kolom waktu. Memori terkelola digunakan untuk menyimpan cache data dalam mode pemrosesan miniBatch jika parameter
table.exec.mini-batch.sizetidak dikonfigurasi atau diatur ke nilai negatif.Operator MiniBatchAssigner gagal dibuat dan tidak dapat mengirim pesan watermark ke operator komputasi untuk memicu perhitungan akhir dan output data. Perhitungan akhir dan output data hanya dipicu ketika salah satu kondisi berikut terpenuhi: memori terkelola penuh, perintah CHECKPOINT diterima dan checkpoint belum dilakukan, atau pekerjaan dibatalkan. Untuk informasi lebih lanjut, lihat table.exec.mini-batch.size. Jika interval checkpoint diatur ke nilai yang terlalu besar, operator LocalGroupAggregate tidak memicu output data dalam waktu yang lama.
-
Solusi
-
Kurangi interval checkpoint agar operator LocalGroupAggregate memicu output data sebelum checkpointing. Tuning Checkpointing.
-
Gunakan memori heap untuk menyimpan cache data. Output dipicu secara otomatis ketika data cache mencapai nilai
table.exec.mini-batch.size. Atur parameter ini ke nilai positif N. Bagaimana cara mengonfigurasi parameter runtime kustom untuk pekerjaan?
-
Apa yang harus saya lakukan jika ketidakaktifan partisi Kafka menunda output window?
Jika connector Kafka hulu memiliki beberapa partisi tetapi hanya sebagian yang menerima data, partisi yang tidak aktif mencegah watermark maju. Window tidak dapat ditutup tepat waktu, sehingga menunda output real-time.
Konfigurasikan timeout untuk menandai partisi yang tidak aktif. Partisi yang tidak aktif akan dikecualikan dari perhitungan watermark hingga menerima data lagi. Configuration.
Tambahkan konfigurasi berikut ke bidang Other Configuration di bagian Parameters pada tab Configuration. Bagaimana cara mengonfigurasi parameter runtime kustom untuk pekerjaan?
table.exec.source.idle-timeout: 1s
Bagaimana cara menemukan error jika JobManager tidak berjalan?
Halaman Flink UI tidak muncul karena JobManager tidak berjalan. Untuk mengidentifikasi penyebabnya, lakukan langkah-langkah berikut:
-
Di bilah navigasi kiri Development Console, pilih . Di halaman Deployments, temukan deployment pekerjaan target dan klik namanya.
-
Klik tab Events.
-
Cari error menggunakan pintasan keyboard OS Anda:
-
Windows: Ctrl+F
-
macOS: Command+F
Di halaman detail pekerjaan, klik tab Events untuk melihat daftar event siklus hidup pekerjaan dan pesan error. Di penampil log, temukan baris exception yang sesuai. Misalnya, log WARN baris 52 dapat menunjukkan error parsing pasangan kunci-nilai untuk kunci
$internal.application.program-argsdi baris 43 dalam/flink/conf/flink-conf.yaml. -
Apa yang harus saya lakukan saat muncul pesan "INFO: org.apache.flink.fs.osshadoop.shaded.com.aliyun.oss"?
-
Deskripsi masalah
2020-08-09 10:18:06,010 INFO org.apache.flink.runtime.jobmaster.JobMaster [] - Configuring application-defined state backend with job/cluster config 2020-08-09 10:18:06,249 INFO org.apache.flink.fs.osshadoop.shaded.com.aliyun.oss [] - [Server]Unable to execute HTTP request: Not Found [ErrorCode]: NoSuchKey [RequestId]: 5F2F5CDE79BCB53539AF510C [HostId]: null 2020-08-09 10:18:06,262 INFO org.apache.flink.fs.osshadoop.shaded.com.aliyun.oss [] - [Server]Unable to execute HTTP request: Not Found [ErrorCode]: NoSuchKey [RequestId]: 5F2F5CDE79BCB53539CF510C [HostId]: null 2020-08-09 10:18:06,349 INFO org.apache.flink.fs.osshadoop.shaded.com.aliyun.oss [] - [Server]Unable to execute HTTP request: Not Found [ErrorCode]: NoSuchKey [RequestId]: 5F2F5CDE79BCB535391452OC [HostId]: null -
Penyebab
Data disimpan di Bucket OSS. Saat OSS membuat direktori, sistem memeriksa apakah direktori tersebut sudah ada. Jika tidak, pesan INFO ini dicetak. Hal ini tidak memengaruhi pekerjaan Anda.
-
Solusi
Tambahkan konfigurasi logger berikut ke templat log Anda untuk menekan pesan ini:
<Logger level="ERROR" name="org.apache.flink.fs.osshadoop.shaded.com.aliyun.oss"/>. topik konfigurasi log.
Apa yang harus saya lakukan jika muncul pesan error "akka.pattern.AskTimeoutException"?
-
Penyebab
-
Penyebab 1: Pengumpulan sampah (GC) yang sering terjadi. Memori JobManager atau TaskManager tidak mencukupi menyebabkan GC sering terjadi, sehingga menyebabkan timeout heartbeat dan RPC antara JobManager dan TaskManager.
-
Penyebab 2: Volume permintaan RPC yang tinggi. Terlalu banyak permintaan RPC membebani JobManager, menyebabkan backlog RPC dan timeout heartbeat serta RPC.
-
Penyebab 3: Nilai timeout terlalu kecil. Pengaturan timeout terlalu rendah. Saat Realtime Compute for Apache Flink mencoba koneksi ulang ke layanan pihak ketiga, timeout habis sebelum kegagalan dilaporkan.
-
-
Solusi
-
Solusi 1: Periksa frekuensi dan durasi GC dari penggunaan memori pekerjaan dan log GC. Jika GC sering terjadi atau berlangsung lama, tingkatkan memori JobManager dan TaskManager.
-
Solusi 2: Untuk menangani volume permintaan RPC yang tinggi, tingkatkan jumlah core CPU dan ukuran memori JobManager, serta atur parameter
akka.ask.timeoutdanheartbeat.timeoutke nilai yang lebih besar.Penting-
Sesuaikan
akka.ask.timeoutdanheartbeat.timeouthanya saat terdapat banyak permintaan RPC. Untuk pekerjaan dengan sedikit permintaan RPC, nilai kecil biasanya tidak menyebabkan masalah ini. -
Atur nilai berdasarkan kebutuhan bisnis Anda. Nilai yang terlalu besar akan meningkatkan waktu pemulihan saat TaskManager keluar secara tak terduga.
-
-
Solusi 3: Untuk menangani kegagalan koneksi layanan pihak ketiga, tingkatkan parameter berikut agar kegagalan koneksi dilaporkan segera:
-
client.timeout: Nilai default: 60. Nilai yang direkomendasikan: 600. Satuan: detik. -
akka.ask.timeout: Nilai default: 10. Nilai yang direkomendasikan: 600. Satuan: detik. -
client.heartbeat.timeout: Nilai default: 180.000. Nilai yang direkomendasikan: 600.000. Satuan: detik.CatatanUntuk mencegah error, jangan sertakan satuan dalam nilainya.
-
heartbeat.timeout: Nilai default: 50.000. Nilai yang direkomendasikan: 600.000. Satuan: milidetik.CatatanUntuk mencegah error, jangan sertakan satuan dalam nilainya.
Sebagai contoh, jika muncul pesan error
"Caused by: java.sql.SQLTransientConnectionException: connection-pool-xxx.mysql.rds.aliyuncs.com:3306 - Connection is not available, request timed out after 30000ms", berarti kolam koneksi MySQL penuh. Dalam kasus ini, Anda harus meningkatkan nilai parameterconnection.pool.sizeyang dijelaskan dalam parameter WITH MySQL. Nilai default: 20.CatatanTentukan nilai minimum berdasarkan pesan error timeout. Nilai yang ditampilkan dalam error menunjukkan pengaturan saat ini. Misalnya, "60.000 ms" dalam
"pattern.AskTimeoutException: Ask timed out on [Actor[akka://flink/user/rpc/dispatcher_1#1064915964]] after [60000 ms]."adalah nilaiclient.timeout. -
-
Apa yang harus saya lakukan jika muncul pesan error "Task did not exit gracefully within 180 + seconds."?
-
Deskripsi masalah
Task did not exit gracefully within 180 + seconds. 2022-04-22T17:32:25.852861506+08:00 stdout F org.apache.flink.util.FlinkRuntimeException: Task did not exit gracefully within 180 + seconds. 2022-04-22T17:32:25.852865065+08:00 stdout F at org.apache.flink.runtime.taskmanager.Task$TaskCancelerWatchDog.run(Task.java:1709) [flink-dist_2.11-1.12-vvr-3.0.4-SNAPSHOT.jar:1.12-vvr-3.0.4-SNAPSHOT] 2022-04-22T17:32:25.852867996+08:00 stdout F at java.lang.Thread.run(Thread.java:834) [?:1.8.0_102] log_level:ERROR -
Penyebab
Error ini tidak menunjukkan akar permasalahan. Artinya, proses keluar task terhambat selama failover atau pembatalan lebih dari
task.cancellation.timeoutdefault yaitu 180 detik. Realtime Compute for Apache Flink menganggap task tersebut tidak dapat dipulihkan, menghentikan TaskManager yang terpengaruh, dan memungkinkan failover atau pembatalan dilanjutkan.Hal ini sering disebabkan oleh user-defined function (UDF). Misalnya, jika metode close dalam UDF terblokir atau tidak mengembalikan nilai, task tidak dapat keluar.
-
Solusi
Untuk debugging, atur
task.cancellation.timeoutke 0. Bagaimana cara mengonfigurasi parameter runtime kustom untuk pekerjaan? Saat diatur ke 0, task yang terblokir akan menunggu tanpa batas untuk keluar tanpa memicu timeout. Jika failover terpicu lagi atau task tetap terhambat setelah restart, temukan task dalam status CANCELLING, periksa jejak stack-nya, dan perbaiki akar permasalahannya.PentingParameter
task.cancellation.timeouthanya untuk debugging. Jangan atur ke 0 di lingkungan produksi. Gunakan timeout yang sesuai dan perbaiki masalah mendasar pada UDF atau logika bisnis.
Apa yang harus saya lakukan saat muncul pesan error "Can not retract a non-existent record. This should never happen."?
-
Deskripsi masalah
java.lang.RuntimeException: Can not retract a non-existent record. This should never happen. at org.apache.flink.table.runtime.operators.rank.RetractableTopNFunction.processElement(RetractableTopNFunction.java:196) at org.apache.flink.table.runtime.operators.rank.RetractableTopNFunction.processElement(RetractableTopNFunction.java:55) at org.apache.flink.streaming.api.operators.KeyedProcessOperator.processElement(KeyedProcessOperator.java:83) at org.apache.flink.streaming.runtime.tasks.OneInputStreamTask$StreamTaskNetworkOutput.emitRecord(OneInputStreamTask.java:205) at org.apache.flink.streaming.runtime.io.AbstractStreamTaskNetworkInput.processElement(AbstractStreamTaskNetworkInput.java:135) at org.apache.flink.streaming.runtime.io.AbstractStreamTaskNetworkInput.emitNext(AbstractStreamTaskNetworkInput.java:106) at org.apache.flink.streaming.runtime.io.StreamOneInputProcessor.processInput(StreamOneInputProcessor.java:66) at org.apache.flink.streaming.runtime.tasks.StreamTask.processInput(StreamTask.java:424) at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:204) at org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:685) at org.apache.flink.streaming.runtime.tasks.StreamTask.executeInvoke(StreamTask.java:640) at org.apache.flink.streaming.runtime.tasks.StreamTask.runWithCleanUpOnFail(StreamTask.java:651) at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:624) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:799) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:586) at java.lang.Thread.run(Thread.java:877) -
Penyebab dan solusi
Skenario
Penyebab
Solusi
Skenario 1
Masalah ini disebabkan oleh fungsi
now()dalam kode.Algoritma TopN tidak mengizinkan bidang non-deterministik digunakan dalam klausa ORDER BY atau PARTITION BY. Jika bidang non-deterministik digunakan, nilai yang dikembalikan oleh fungsi
now()berbeda untuk setiap catatan, sehingga nilai sebelumnya tidak dapat ditemukan dalam status.Gunakan bidang deterministik dalam klausa ORDER BY atau PARTITION BY.
Skenario 2
Parameter
table.exec.state.ttldiatur ke nilai yang terlalu kecil. Akibatnya, entri status kedaluwarsa dan dihapus, sehingga status kunci yang diperlukan tidak dapat ditemukan dalam status.Tingkatkan nilai
table.exec.state.ttl. Bagaimana cara mengonfigurasi parameter runtime kustom untuk pekerjaan?
Bagaimana cara memperbaiki pesan error "The GRPC call timed out in sqlserver"?
-
Deskripsi masalah
org.apache.flink.table.sqlserver.utils.ExecutionTimeoutException: The GRPC call timed out in sqlserver, please check the thread stacktrace for root cause: Thread name: sqlserver-operation-pool-thread-4, thread state: TIMED_WAITING, thread stacktrace: at java.lang.Thread.sleep0(Native Method) at java.lang.Thread.sleep(Thread.java:360) at org.apache.hadoop.io.retry.RetryInvocationHandler$Call.processWaitTimeAndRetryInfo(RetryInvocationHandler.java:130) at org.apache.hadoop.io.retry.RetryInvocationHandler$Call.invokeOnce(RetryInvocationHandler.java:107) at org.apache.hadoop.io.retry.RetryInvocationHandler.invoke(RetryInvocationHandler.java:359) at com.sun.proxy.$Proxy195.getFileInfo(Unknown Source) at org.apache.hadoop.hdfs.DFSClient.getFileInfo(DFSClient.java:1661) at org.apache.hadoop.hdfs.DistributedFileSystem$29.doCall(DistributedFileSystem.java:1577) at org.apache.hadoop.hdfs.DistributedFileSystem$29.doCall(DistributedFileSystem.java:1574) at org.apache.hadoop.fs.FileSystemLinkResolver.resolve(FileSystemLinkResolver.java:81) at org.apache.hadoop.hdfs.DistributedFileSystem.getFileStatus(DistributedFileSystem.java:1589) at org.apache.hadoop.fs.FileSystem.exists(FileSystem.java:1683) at org.apache.flink.connectors.hive.HiveSourceFileEnumerator.getNumFiles(HiveSourceFileEnumerator.java:118) at org.apache.flink.connectors.hive.HiveTableSource.lambda$getDataStream$0(HiveTableSource.java:209) at org.apache.flink.connectors.hive.HiveTableSource$$Lambda$972/1139330351.get(Unknown Source) at org.apache.flink.connectors.hive.HiveParallelismInference.logRunningTime(HiveParallelismInference.java:118) at org.apache.flink.connectors.hive.HiveParallelismInference.infer(HiveParallelismInference.java:100) at org.apache.flink.connectors.hive.HiveTableSource.getDataStream(HiveTableSource.java:207) at org.apache.flink.connectors.hive.HiveTableSource$1.produceDataStream(HiveTableSource.java:123) at org.apache.flink.table.planner.plan.nodes.exec.common.CommonExecTableSourceScan.translateToPlanInternal(CommonExecTableSourceScan.java:127) at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.translateToPlan(ExecNodeBase.java:226) at org.apache.flink.table.planner.plan.nodes.exec.ExecEdge.translateToPlan(ExecEdge.java:290) at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.lambda$translateInputToPlan$5(ExecNodeBase.java:267) at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase$$Lambda$949/77002396.apply(Unknown Source) at java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:193) at java.util.stream.ReferencePipeline$2$1.accept(ReferencePipeline.java:175) at java.util.ArrayList$ArrayListSpliterator.forEachRemaining(ArrayList.java:1374) at java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:481) at java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:471) at java.util.stream.ReduceOps$ReduceOp.evaluateSequential(ReduceOps.java:708) at java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234) at java.util.stream.ReferencePipeline.collect(ReferencePipeline.java:499) at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.translateInputToPlan(ExecNodeBase.java:268) at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.translateInputToPlan(ExecNodeBase.java:241) at org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecExchange.translateToPlanInternal(StreamExecExchange.java:87) at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.translateToPlan(ExecNodeBase.java:226) at org.apache.flink.table.planner.plan.nodes.exec.ExecEdge.translateToPlan(ExecEdge.java:290) at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.lambda$translateInputToPlan$5(ExecNodeBase.java:267) at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase$$Lambda$949/77002396.apply(Unknown Source) at java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:193) at java.util.stream.ReferencePipeline$2$1.accept(ReferencePipeline.java:175) at java.util.ArrayList$ArrayListSpliterator.forEachRemaining(ArrayList.java:1374) at java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:481) at java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:471) at java.util.stream.ReduceOps$ReduceOp.evaluateSequential(ReduceOps.java:708) at java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234) at java.util.stream.ReferencePipeline.collect(ReferencePipeline.java:499) at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.translateInputToPlan(ExecNodeBase.java:268) at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.translateInputToPlan(ExecNodeBase.java:241) at org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecGroupAggregate.translateToPlanInternal(StreamExecGroupAggregate.java:148) at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.translateToPlan(ExecNodeBase.java:226) at org.apache.flink.table.planner.plan.nodes.exec.ExecEdge.translateToPlan(ExecEdge.java:290) at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.lambda$translateInputToPlan$5(ExecNodeBase.java:267) at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase$$Lambda$949/77002396.apply(Unknown Source) at java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:193) at java.util.stream.ReferencePipeline$2$1.accept(ReferencePipeline.java:175) at java.util.ArrayList$ArrayListSpliterator.forEachRemaining(ArrayList.java:1374) at java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:481) at java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:471) at java.util.stream.ReduceOps$ReduceOp.evaluateSequential(ReduceOps.java:708) at java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234) at java.util.stream.ReferencePipeline.collect(ReferencePipeline.java:499) at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.translateInputToPlan(ExecNodeBase.java:268) at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.translateInputToPlan(ExecNodeBase.java:241) at org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecSink.translateToPlanInternal(StreamExecSink.java:108) at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.translateToPlan(ExecNodeBase.java:226) at org.apache.flink.table.planner.delegation.StreamPlanner$$anonfun$1.apply(StreamPlanner.scala:74) at org.apache.flink.table.planner.delegation.StreamPlanner$$anonfun$1.apply(StreamPlanner.scala:73) at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234) at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234) at scala.collection.Iterator$class.foreach(Iterator.scala:891) at scala.collection.AbstractIterator.foreach(Iterator.scala:1334) at scala.collection.IterableLike$class.foreach(IterableLike.scala:72) at scala.collection.AbstractIterable.foreach(Iterable.scala:54) at scala.collection.TraversableLike$class.map(TraversableLike.scala:234) at scala.collection.AbstractTraversable.map(Traversable.scala:104) at org.apache.flink.table.planner.delegation.StreamPlanner.translateToPlan(StreamPlanner.scala:73) at org.apache.flink.table.planner.delegation.StreamExecutor.createStreamGraph(StreamExecutor.scala:52) at org.apache.flink.table.planner.delegation.PlannerBase.createStreamGraph(PlannerBase.scala:610) at org.apache.flink.table.planner.delegation.StreamPlanner.explainExecNodeGraphInternal(StreamPlanner.scala:166) at org.apache.flink.table.planner.delegation.StreamPlanner.explainExecNodeGraph(StreamPlanner.scala:159) at org.apache.flink.table.sqlserver.execution.OperationExecutorImpl.validate(OperationExecutorImpl.java:304) at org.apache.flink.table.sqlserver.execution.OperationExecutorImpl.validate(OperationExecutorImpl.java:288) at org.apache.flink.table.sqlserver.execution.DelegateOperationExecutor.lambda$validate$22(DelegateOperationExecutor.java:211) at org.apache.flink.table.sqlserver.execution.DelegateOperationExecutor$$Lambda$394/1626790418.run(Unknown Source) at org.apache.flink.table.sqlserver.execution.DelegateOperationExecutor.wrapClassLoader(DelegateOperationExecutor.java:250) at org.apache.flink.table.sqlserver.execution.DelegateOperationExecutor.lambda$wrapExecutor$26(DelegateOperationExecutor.java:275) at org.apache.flink.table.sqlserver.execution.DelegateOperationExecutor$$Lambda$395/1157752141.run(Unknown Source) at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) at java.util.concurrent.FutureTask.run(FutureTask.java:266) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1147) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:622) at java.lang.Thread.run(Thread.java:834) at org.apache.flink.table.sqlserver.execution.DelegateOperationExecutor.wrapExecutor(DelegateOperationExecutor.java:281) at org.apache.flink.table.sqlserver.execution.DelegateOperationExecutor.validate(DelegateOperationExecutor.java:211) at org.apache.flink.table.sqlserver.FlinkSqlServiceImpl.validate(FlinkSqlServiceImpl.java:786) at org.apache.flink.table.sqlserver.proto.FlinkSqlServiceGrpc$MethodHandlers.invoke(FlinkSqlServiceGrpc.java:2522) at io.grpc.stub.ServerCalls$UnaryServerCallHandler$UnaryServerCallListener.onHalfClose(ServerCalls.java:172) at io.grpc.internal.ServerCallImpl$ServerStreamListenerImpl.halfClosed(ServerCallImpl.java:331) at io.grpc.internal.ServerImpl$JumpToApplicationThreadServerStreamListener$1HalfClosed.runInContext(ServerImpl.java:820) at io.grpc.internal.ContextRunnable.run(ContextRunnable.java:37) at io.grpc.internal.SerializingExecutor.run(SerializingExecutor.java:123) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1147) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:622) at java.lang.Thread.run(Thread.java:834) Caused by: java.util.concurrent.TimeoutException at java.util.concurrent.FutureTask.get(FutureTask.java:205) at org.apache.flink.table.sqlserver.execution.DelegateOperationExecutor.wrapExecutor(DelegateOperationExecutor.java:277) ... 11 more -
Penyebab
SQL kompleks dalam draft menyebabkan timeout eksekusi RPC.
-
Solusi
Tambahkan kode berikut ke bidang Other Configuration di bagian Parameters pada tab Configuration untuk meningkatkan waktu tunggu RPC. Nilai default-nya adalah 120 detik. Untuk informasi selengkapnya, lihat Mengonfigurasi parameter berjalan kustom.
flink.sqlserver.rpc.execution.timeout: 600s
Bagaimana cara mengatasi pesan error "RESOURCE_EXHAUSTED: gRPC message exceeds maximum size 41943040: 58384051"?
-
Deskripsi masalah
Caused by: io.grpc.StatusRuntimeException: RESOURCE_EXHAUSTED: gRPC message exceeds maximum size 41943040: 58384051 at io.grpc.stub.ClientCalls.toStatusRuntimeException(ClientCalls.java:244) at io.grpc.stub.ClientCalls.getUnchecked(ClientCalls.java:225) at io.grpc.stub.ClientCalls.blockingUnaryCall(ClientCalls.java:142) at org.apache.flink.table.sqlserver.proto.FlinkSqlServiceGrpc$FlinkSqlServiceBlockingStub.generateJobGraph(FlinkSqlServiceGrpc.java:2478) at org.apache.flink.table.sqlserver.api.client.FlinkSqlServerProtoClientImpl.generateJobGraph(FlinkSqlServerProtoClientImpl.java:456) at org.apache.flink.table.sqlserver.api.client.ErrorHandlingProtoClient.lambda$generateJobGraph$25(ErrorHandlingProtoClient.java:251) at org.apache.flink.table.sqlserver.api.client.ErrorHandlingProtoClient.invokeRequest(ErrorHandlingProtoClient.java:335) ... 6 more Cause: RESOURCE_EXHAUSTED: gRPC message exceeds maximum size 41943040: 58384051) -
Penyebab
JobGraph terlalu besar karena logika draft yang kompleks. Hal ini menyebabkan error verifikasi atau mencegah pekerjaan draft dimulai atau dibatalkan.
-
Solusi
Tambahkan kode berikut ke bidang Other Configuration di bagian Parameters pada tab Configuration. Bagaimana cara mengonfigurasi parameter runtime kustom untuk pekerjaan?
table.exec.operator-name.max-length: 1000
Apa yang harus saya lakukan jika muncul pesan error "Caused by: java.lang.NoSuchMethodError"?
-
Deskripsi masalah
Error message: Caused by: java.lang.NoSuchMethodError: org.apache.flink.table.planner.plan.metadata.FlinkRelMetadataQuery.getUpsertKeysInKeyGroupRange(Lorg/apache/calcite/rel/RelNode;[I)Ljava/util/Set; -
Penyebab
Jika Anda memanggil API Apache Flink dan Realtime Compute for Apache Flink menyediakan versi yang dioptimalkan, exception seperti konflik paket dapat terjadi.
-
Solusi
Batasi pemanggilan metode Anda hanya pada metode yang secara eksplisit ditandai dengan @Public atau @PublicEvolving dalam kode sumber Apache Flink. Realtime Compute for Apache Flink menjamin kompatibilitas dengan metode-metode tersebut.
Apa yang harus saya lakukan jika muncul pesan error "java.lang.ClassCastException: org.codehaus.janino.CompilerFactory cannot be cast to org.codehaus.commons.compiler.ICompilerFactory"?
-
Deskripsi masalah
Causedby:java.lang.ClassCastException:org.codehaus.janino.CompilerFactorycannotbecasttoorg.codehaus.commons.compiler.ICompilerFactory atorg.codehaus.commons.compiler.CompilerFactoryFactory.getCompilerFactory(CompilerFactoryFactory.java:129) atorg.codehaus.commons.compiler.CompilerFactoryFactory.getDefaultCompilerFactory(CompilerFactoryFactory.java:79) atorg.apache.calcite.rel.metadata.JaninoRelMetadataProvider.compile(JaninoRelMetadataProvider.java:426) ...66more -
Penyebab
-
Paket JAR berisi dependensi Janino yang menyebabkan konflik.
-
Paket JAR tertentu yang diawali dengan
Flink-sepertiflink-table-plannerdanflink-table-runtimeditambahkan ke paket JAR UDF atau connector.
-
-
Solusi
-
Periksa apakah paket JAR berisi
org.codehaus.janino.CompilerFactory. Konflik kelas dapat terjadi karena urutan pemuatan kelas berbeda di mesin yang berbeda. Untuk mengatasi masalah ini, lakukan langkah-langkah berikut:-
Di bilah navigasi kiri Development Console, pilih . Di halaman Deployments, temukan pekerjaan target dan klik namanya.
-
Di tab Configuration di halaman detail pekerjaan, klik Edit di pojok kanan atas bagian Parameters.
-
Tambahkan kode berikut ke bidang Other Configuration dan klik Save.
classloader.parent-first-patterns.additional: org.codehaus.janinoGanti nilai parameter
classloader.parent-first-patterns.additionaldengan kelas yang konflik.
-
-
Tentukan
<scope>provided</scope>untuk dependensi Apache Flink, seperti dependensi non-connector yang namanya diawali denganflink-dalam gruporg.apache.flink.
-