All Products
Search
Document Center

Realtime Compute for Apache Flink:FAQ error pekerjaan

Last Updated:Aug 27, 2026

Error runtime pekerjaan umum dan solusinya di Realtime Compute for Apache Flink.

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: resourcequota

    Sumber daya tidak mencukupi di Resource Queue saat ini.

    Tingkatkan kapasitas Resource Queue atau kurangi kebutuhan sumber daya pekerjaan.

    ERROR:the vswitch ip is not enough

    Alamat IP tidak mencukupi di namespace untuk TaskManager yang diperlukan.

    Kurangi paralelisme pekerjaan, sesuaikan konfigurasi slot, atau ubah pengaturan vSwitch.

    ERROR: pooler: ***: authentication failed

    Pasangan 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.size tidak 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

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:

  1. Di bilah navigasi kiri Development Console, pilih O&M > Deployments. Di halaman Deployments, temukan deployment pekerjaan target dan klik namanya.

  2. Klik tab Events.

  3. 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-args di 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.timeout dan heartbeat.timeout ke nilai yang lebih besar.

      Penting
      • Sesuaikan akka.ask.timeout dan heartbeat.timeout hanya 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.

        Catatan

        Untuk mencegah error, jangan sertakan satuan dalam nilainya.

      • heartbeat.timeout: Nilai default: 50.000. Nilai yang direkomendasikan: 600.000. Satuan: milidetik.

        Catatan

        Untuk 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 parameter connection.pool.size yang dijelaskan dalam parameter WITH MySQL. Nilai default: 20.

      Catatan

      Tentukan 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 nilai client.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.timeout default 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.timeout ke 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.

    Penting

    Parameter task.cancellation.timeout hanya 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.ttl diatur 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- seperti flink-table-planner dan flink-table-runtime ditambahkan 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:

      1. Di bilah navigasi kiri Development Console, pilih O&M > Deployments. Di halaman Deployments, temukan pekerjaan target dan klik namanya.

      2. Di tab Configuration di halaman detail pekerjaan, klik Edit di pojok kanan atas bagian Parameters.

      3. Tambahkan kode berikut ke bidang Other Configuration dan klik Save.

        classloader.parent-first-patterns.additional: org.codehaus.janino

        Ganti nilai parameter classloader.parent-first-patterns.additional dengan kelas yang konflik.

    • Tentukan <scope>provided</scope> untuk dependensi Apache Flink, seperti dependensi non-connector yang namanya diawali dengan flink- dalam grup org.apache.flink.