すべてのプロダクト
Search
ドキュメントセンター

Realtime Compute for Apache Flink:ジョブエラーに関するよくある質問

最終更新日:Aug 27, 2026

Realtime Compute for Apache Flink の一般的なジョブのランタイムエラーと解決策。

ジョブが開始できない場合はどうすればよいですか?

  • 問題の説明

    [操作] 列で [Start] をクリックすると、ジョブのステータスが [開始中] から [失敗] に変わります。

  • 解決策

    • [イベント] タブを確認する: ジョブの詳細ページで [イベント] タブに移動します。ジョブの開始が失敗した際に発生した障害イベントを見つけ、その詳細を確認して根本原因を特定してください。

    • 起動ログを確認する: [ログ] タブに移動し、[起動ログ] サブタブを選択します。ジョブの開始が失敗した理由を説明する具体的なエラーメッセージがないか、ログを確認してください。

    • JobManager / TaskManager のログを確認する: JobManager が正常に起動しているように見えるにもかかわらずジョブが失敗する場合は、JobManager と TaskManager の詳細なログを確認してください。これらは、[ログ] タブ内の [Job Manager] または [実行中の TaskManager] サブタブで確認できます。

  • 一般的なエラーと解決策

    問題の説明

    原因

    解決策

    ERROR:exceeded quota: resourcequota

    現在のリソースキュー内のリソースが不足しています。

    リソースキューの容量を増やすか、ジョブのリソース要件を削減してください。

    ERROR:the vswitch ip is not enough

    名前空間内の IP アドレスが、必要な TaskManager に対して不足しています。

    ジョブの並列度を下げるか、スロット設定を調整するか、vSwitch の設定を変更してください。

    ERROR: pooler: ***: authentication failed

    AccessKey ペアが無効であるか、権限が不足しています。

    AccessKey ペアが有効であり、ジョブを実行および管理する権限を持つアカウントに属していることを確認してください。

データベース接続エラーの修正方法

  • 問題の説明

    failed to execute sql statement
  • 原因

    登録済みのカタログが無効、または到達不能です。

  • 解決策

    [Catalogs] ページに移動し、グレー表示されているカタログを削除してから再登録してください。

ジョブ実行後にタスクのデータが消費されない場合

  • ネットワーク接続の確認

    アップストリームおよびダウンストリームのストレージでデータが生成または消費されない場合は、[起動ログ] タブでエラーメッセージを確認してください。タイムアウトエラーが表示される場合は、ストレージシステム間のネットワーク接続をトラブルシューティングしてください。

  • タスク実行ステータスの確認

    [ステータス] タブで、ソースからデータが読み取られ、シンクに書き込まれているかどうかを確認し、エラーの発生箇所を特定してください。

    メトリクステーブルで、ソースノードの [送信バイト数] 列と [送信レコード数] 列を確認してください。値が 0 より大きい場合 (例: 19.74 GB / 76,766,861 レコード)、ソースは正常にデータを送信しています。同様に、ダウンストリームノードの [受信バイト数] 列と [受信レコード数] 列で、データが正常に受信されていることを確認してください。

  • オペレーター出力の確認

    各 オペレーター に print シンクテーブルを追加して、問題をトラブルシューティングしてください。

ジョブが予期せず再起動した場合の対処方法

エラーをトラブルシューティングするには、[ログ] タブを確認してください。

  • 例外情報の表示

    [JM 例外] サブタブで、報告されたエラーを確認し、根本原因を特定します。

  • ジョブの JobManager および TaskManager ログの表示

    [ログ]タブで[ジョブマネージャー]サブタブをクリックし、次に[ログ]サブタブを選択すると、対応するジョブログが表示されます。同様に、[実行中のタスクマネージャー]サブタブをクリックすると、TM ログが表示されます。

  • 失敗した TaskManager のログの表示

    一部の例外によって TaskManager が失敗し、ログが不完全になることがあります。 最後の無効な TaskManager のログを表示し、問題のトラブルシューティングを行います。

  • 過去のジョブインスタンスログの表示

    過去のジョブインスタンスログを確認して、失敗の原因を特定します。

LocalGroupAggregate オペレーターでのデータ出力の中断

  • コード

    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;
  • 問題の説明

    LocalGroupAggregate オペレーターでデータ出力が長時間中断され、ジョブトポロジに MiniBatchAssigner オペレーターが存在しません。

    [Operator Analysis (Beta)] 画面では、ジョブトポロジで LocalGroupAggregate[201] の RecordsIn が 1853、RecordsOut がわずか 116 となっています。ダウンストリームの Vertex2 (GlobalGroupAggregate および Calc オペレーターを含む) の Out (sum) は 0、Backpressured (max) は 0% です。トポロジ全体に MiniBatchAssigner ノードは表示されません。

  • 原因

    ジョブには WindowAggregate オペレーターと GroupAggregate オペレーターの両方が含まれています。WindowAggregate オペレーターは時間列として proctime を使用しています。table.exec.mini-batch.size パラメーターが設定されていないか、負の値に設定されている場合、ミニバッチ処理モードでデータをキャッシュするために管理メモリが使用されます。

    MiniBatchAssigner オペレーターの生成に失敗し、計算オペレーターにウォーターマークメッセージを送信して最終計算とデータ出力をトリガーできません。最終計算とデータ出力は、次のいずれかの条件が満たされた場合にのみトリガーされます:管理メモリが満杯になる、CHECKPOINT コマンドを受信したにもかかわらずチェックポイント処理が実行されていない、またはジョブがキャンセルされる。詳細については、table.exec.mini-batch.size をご参照ください。チェックポイント間隔が非常に大きな値に設定されている場合、LocalGroupAggregate オペレーターは長時間データ出力をトリガーしません。

  • 解決策

    • チェックポイント間隔を短くして、LocalGroupAggregate オペレーターがチェックポイント処理の前にデータ出力をトリガーするようにします。詳細については、「チェックポイントのチューニング」をご参照ください。

    • ヒープメモリを使用してデータをキャッシュします。キャッシュされたデータが table.exec.mini-batch.size の値に達すると、出力が自動的にトリガーされます。このパラメーターを正の値 N に設定してください。詳細については、「ジョブのカスタムランタイムパラメーターの設定」をご参照ください。

Kafka パーティションのアイドル状態によるウィンドウ出力の遅延

アップストリームの Kafka コネクタに複数のパーティションがあり、一部のパーティションのみがデータを受信している場合、アイドル状態のパーティションによってウォーターマークが進まなくなります。その結果、ウィンドウを適時にクローズできず、リアルタイム出力が遅延します。

タイムアウトを設定して、アイドル状態のパーティションをマークします。アイドル状態のパーティションは、再度データを受信するまでウォーターマークの計算から除外されます。詳細については、「設定」をご参照ください。

[設定] タブの [パラメーター] セクションにある [その他の設定] フィールドに、次の設定を追加します。詳細については、「ジョブのカスタムランタイムパラメーターを設定するにはどうすればよいですか?」をご参照ください。

table.exec.source.idle-timeout: 1s

JobManager が実行されていない場合のエラー特定方法

JobManager が実行されていないため、[Flink UI] ページが表示されません。原因を特定するには、次の手順を実行してください。

  1. 開発コンソールの左側メニューで、[O&M] > [デプロイメント] を選択します。[デプロイメント] ページで、対象のジョブのデプロイメントを見つけて、その名前をクリックします。

  2. [イベント] タブをクリックします。

  3. OS のキーボードショートカットを使用してエラーを検索します。

    • Windows: Ctrl+F

    • macOS: Command+F

    ジョブの詳細ページで [イベント] タブをクリックして、ジョブのライフサイクルイベントリストとエラーメッセージを表示します。ログビューアーで、一致する例外の行を見つけます。たとえば、52 行目の WARN ログは、/flink/conf/flink-conf.yaml の 43 行目における、キー $internal.application.program-args のキーと値のペアの解析エラーを示している可能性があります。

「INFO: org.apache.flink.fs.osshadoop.shaded.com.aliyun.oss」というメッセージが表示された場合の対処法

  • 問題の説明

    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
  • 原因

    データは OSS バケットに保存されます。OSS がディレクトリを作成する際、そのディレクトリが存在するかどうかを確認します。存在しない場合、この INFO メッセージが出力されます。これはジョブに影響を与えません。

  • 解決策

    このメッセージを抑制するには、次のロガー設定をログテンプレートに追加してください:<Logger level="ERROR" name="org.apache.flink.fs.osshadoop.shaded.com.aliyun.oss"/>。詳細については、「ログ設定のトピック」をご参照ください。

「akka.pattern.AskTimeoutException」エラーの対処法

  • 原因

    • 原因1: ガベージコレクション (GC) の頻発 JobManager または TaskManager のメモリが不足すると、GC が頻繁に発生し、JobManager と TaskManager 間のハートビートおよび RPC タイムアウトが発生します。

    • 原因2: 大量の RPC リクエスト RPC リクエストが多すぎると JobManager が過負荷になり、RPC のバックログが発生し、ハートビートおよび RPC タイムアウトが発生します。

    • 原因3: 小さすぎるタイムアウト値 タイムアウト値の設定が小さすぎるため、Realtime Compute for Apache Flink がサードパーティサービスへの接続を再試行する際に、障害が報告される前にタイムアウトが発生します。

  • 解決策

    • 解決策1: ジョブのメモリ使用量と GC ログから GC の頻度と期間を確認してください。GC が頻繁に発生する場合、または長時間続く場合は、JobManager と TaskManager のメモリを増やしてください。

    • 解決策2: 大量の RPC リクエストを処理するには、JobManager の CPU コア数とメモリサイズを増やし、akka.ask.timeout および heartbeat.timeout パラメーターをより大きな値に設定してください。

      重要
      • akka.ask.timeout と heartbeat.timeout は、大量の RPC リクエストがある場合にのみ調整してください。RPC リクエストが少ないジョブでは、通常、小さい値でもこの問題は発生しません。

      • ビジネス要件に基づいて値を設定してください。値を大きく設定しすぎると、TaskManager が予期せず終了した際の復旧時間が長くなります。

    • 解決策3: サードパーティサービスの接続障害に対処するには、接続障害が速やかに報告されるように、次のパラメーターを増やしてください。

      • client.timeout:デフォルト値:60。推奨値:600。単位:秒。

      • akka.ask.timeout:デフォルト値:10。推奨値:600。単位:秒。

      • client.heartbeat.timeout:デフォルト値:180000。推奨値:600000。単位:ミリ秒。

        説明

        エラーを防ぐため、値に単位を含めないでください。

      • heartbeat.timeout:デフォルト値:50000。推奨値:600000。単位:ミリ秒。

        説明

        エラーを防ぐため、値に単位を含めないでください。

      たとえば、"Caused by: java.sql.SQLTransientConnectionException: connection-pool-xxx.mysql.rds.aliyuncs.com:3306 - Connection is not available, request timed out after 30000ms" というエラーメッセージが表示される場合、MySQL コネクションプールがいっぱいです。この場合、MySQL の WITH パラメーターにある connection.pool.size パラメーターの値を増やす必要があります。デフォルト値:20。

      説明

      タイムアウトエラーメッセージから最小値を特定してください。エラーに表示される値は、現在の設定を示しています。たとえば、"pattern.AskTimeoutException: Ask timed out on [Actor[akka://flink/user/rpc/dispatcher_1#1064915964]] after [60000 ms]." の "60000 ms" は、akka.ask.timeout の値です。

「Task did not exit gracefully within 180 + seconds.」というエラーメッセージが表示された場合はどうすればよいですか?

  • 問題の説明

    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
  • 原因

    このエラーは根本原因を示すものではありません。これは、フェイルオーバーまたはキャンセル中にタスクの終了が、デフォルトの task.cancellation.timeout である 180 秒を超えてスタックしたことを意味します。リアルタイムコンピューティング for Apache Flink は、タスクを回復不能と見なし、影響を受ける TaskManager を停止して、フェイルオーバーまたはキャンセルの続行を許可します。

    これは多くの場合、ユーザー定義関数 (UDF) が原因で発生します。例えば、UDF の close メソッドがブロックされたり、戻らなかったりすると、タスクは終了できません。

  • 解決策

    デバッグでは、task.cancellation.timeout を 0 に設定します。ジョブのカスタムランタイムパラメーターを設定する方法 0 に設定すると、ブロックされたタスクはタイムアウトすることなく、無期限に終了を待機します。フェイルオーバーが再度トリガーされた場合、または再起動後にタスクがスタックしたままの場合は、CANCELLING 状態のタスクを特定し、そのスタックトレースを検査して、根本原因を修正します。

    重要

    task.cancellation.timeout パラメーターはデバッグ専用です。本番環境では 0 に設定しないでください。適切なタイムアウトを使用し、根本的な UDF またはビジネスロジックの問題を修正してください。

「Can not retract a non-existent record. This should never happen.」エラーへの対処法

  • 問題の説明

    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)
                        
  • 原因と解決策

    シナリオ

    原因

    解決策

    シナリオ 1

    この問題は、コード内の now() 関数が原因です。

    TopN アルゴリズムでは、ORDER BY 句または PARTITION BY 句で非決定性フィールドを使用できません。非決定性フィールドが使用されると、now() 関数によって返される値はレコードごとに異なり、状態で前の値を見つけることができなくなります。

    ORDER BY 句または PARTITION BY 句で決定的フィールドを使用してください。

    シナリオ 2

    table.exec.state.ttl パラメーターが小さすぎる値に設定されているため、状態エントリの期限が切れて削除され、必要なキーステートがステートで見つからなくなります。

    table.exec.state.ttl の値を大きくします。ジョブのカスタムランタイムパラメーターを設定する方法

エラー「The GRPC call timed out in sqlserver」の対処法

  • 問題の説明

    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.java: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
                        
  • 原因

    下書き内の複雑な SQL が、RPC 実行のタイムアウトを引き起こします。

  • 解決策

    RPC タイムアウトの値を大きくするには、[設定] タブの [パラメーター] セクションにある [その他の設定] フィールドに次のコードを追加します。デフォルト値は 120 秒です。詳細については、カスタム実行パラメーターの設定をご参照ください。

    flink.sqlserver.rpc.execution.timeout: 600s

「RESOURCE_EXHAUSTED: gRPC message exceeds maximum size 41943040: 58384051」エラーの解決策

  • 問題の説明

    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)
  • 原因

    複雑なドラフトロジックが原因で JobGraph が大きすぎます。このため、検証エラーが発生したり、下書きジョブの開始やキャンセルができなくなったりします。

  • 解決策

    [設定] タブの [パラメーター] セクションにある [その他の設定] フィールドに、次のコードを追加してください。詳細については、「ジョブのカスタムランタイムパラメーターを設定するにはどうすればよいですか?」をご参照ください。

     table.exec.operator-name.max-length: 1000

エラーメッセージ「Caused by: java.lang.NoSuchMethodError」への対処法

  • 問題の説明

    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;
  • 原因

    Apache Flink の API を呼び出す際に、Realtime Compute for Apache Flink が最適化バージョンを提供している場合、パッケージの競合などの例外が発生することがあります。

  • 解決策

    Apache Flink のソースコードで @Public または @PublicEvolving と明示的にマークされているメソッドのみを呼び出してください。Realtime Compute for Apache Flink は、これらのメソッドとの互換性を確保しています。

"java.lang.ClassCastException: org.codehaus.janino.CompilerFactory cannot be cast to org.codehaus.commons.compiler.ICompilerFactory" というエラーメッセージが表示される場合はどうすればよいですか?

  • 問題の説明

    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
  • 原因

    • JAR パッケージに、競合を引き起こす Janino の依存関係が含まれています。

    • flink- で始まる特定の JAR パッケージ (flink-table-planner や flink-table-runtime など) が、UDF またはコネクタの JAR パッケージに追加されています。

  • 解決策

    • JAR パッケージに org.codehaus.janino.CompilerFactory が含まれているかどうかを確認してください。マシンによってクラスの読み込み順序が異なるため、クラスの競合が発生する可能性があります。この問題を解決するには、次の手順を実行してください。

      1. 開発コンソールの左側メニューで、[O&M] > [Deployments] を選択します。[Deployments] ページで、対象のジョブを見つけて、その名前をクリックします。

      2. ジョブの詳細ページの [Configuration] タブで、[Parameters] セクションの右上隅にある [Edit] をクリックします。

      3. [Other Configuration] フィールドに次のコードを追加し、[Save] をクリックします。

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

        classloader.parent-first-patterns.additional パラメーターの値を、競合する依存関係のパッケージプレフィックスに置き換えてください。

    • org.apache.flink グループ内で、flink- で始まる名前を持つ非コネクタ依存関係など、Apache Flink の依存関係には <scope>provided</scope> を指定してください。