Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Job errors FAQ

Dernière mise à jour :Aug 26, 2026

Erreurs d'exécution courantes des jobs et solutions dans Realtime Compute for Apache Flink.

Que faire si un job ne démarre pas ?

  • Description du problème

    Lorsque vous cliquez sur Start dans la colonne Actions, l'état du job passe de STARTING à FAILED.

  • Solutions

    • Vérifiez l'onglet Events : Accédez à l'onglet Events sur la page des détails du job. Localisez l'événement d'échec survenu lors du démarrage du job et examinez ses détails pour identifier la cause racine.

    • Examinez les journaux de démarrage : Accédez à l'onglet Logs et sélectionnez le sous-onglet Startup Logs. Recherchez dans les journaux des messages d'erreur spécifiques expliquant l'échec du démarrage du job.

    • Inspectez les journaux du JobManager et des TaskManagers : Si le JobManager semble démarrer correctement mais que le job échoue toujours, vérifiez les journaux détaillés du JobManager et des TaskManagers. Vous pouvez les trouver dans les sous-onglets Job Manager ou Running Task Managers de l'onglet Logs.

  • Erreurs courantes et solutions

    Description du problème

    Cause

    Solution

    ERROR:exceeded quota: resourcequota

    Ressources insuffisantes dans la file d'attente de ressources actuelle.

    Augmentez la capacité de la file d'attente de ressources ou réduisez les exigences en ressources du job.

    ERROR:the vswitch ip is not enough

    Adresses IP insuffisantes dans le namespace pour les TaskManagers requis.

    Réduisez le parallélisme du job, ajustez la configuration des slots ou modifiez les paramètres du vSwitch.

    ERROR: pooler: *: authentication failed**

    Paire AccessKey invalide ou autorisations insuffisantes.

    Vérifiez que la paire AccessKey est valide et appartient à un compte disposant des autorisations nécessaires pour exécuter et gérer les jobs.

Comment résoudre une erreur de connexion à la base de données ?

  • Description du problème

    failed to execute sql statement
  • Cause

    Le catalogue enregistré est invalide ou inaccessible.

  • Solution

    Accédez à la page Catalogs, supprimez tous les catalogues grisés et enregistrez-les à nouveau.

Que faire si les données des tâches d'un job ne sont pas consommées après l'exécution du job ?

  • Vérifiez la connectivité réseau

    Si aucune donnée n'est générée ou consommée dans les systèmes de stockage en amont et en aval, consultez l'onglet Startup Logs pour rechercher des messages d'erreur. En cas d'erreurs de délai d'expiration, résolvez les problèmes de connectivité réseau entre les systèmes de stockage.

  • Vérifiez l'état d'exécution des tâches

    Dans l'onglet Status, vérifiez si les données sont lues depuis la source et écrites dans le puits (sink) afin d'identifier l'origine de l'erreur.

    Dans le tableau des métriques, examinez les colonnes Bytes Sent et Records sent pour le nœud Source. Si les valeurs sont supérieures à 0 (par exemple, 19,74 Go / 76 766 861 enregistrements), la source envoie normalement les données. De même, vérifiez les colonnes Bytes Received et Records Received pour les nœuds en aval afin de confirmer qu'ils reçoivent bien les données.

  • Vérifiez la sortie de l'opérateur

    Ajoutez une table de type print sink à chaque opérateur pour diagnostiquer le problème.

Que faire en cas de redémarrage inattendu d'un job ?

Pour diagnostiquer l'erreur, consultez l'onglet Logs.

  • Consultez les informations sur les exceptions.

    Dans le sous-onglet JM Exceptions, examinez l'erreur signalée et identifiez la cause racine.

  • Consultez les journaux du JobManager et des TaskManagers pour le job.

    Dans l'onglet Logs, cliquez sur le sous-onglet Job Manager, puis sélectionnez le sous-onglet Logs pour afficher les journaux correspondants du job. De même, cliquez sur le sous-onglet Running Task Managers pour consulter les journaux des TM.

  • Consultez les journaux des TaskManagers ayant échoué pour le job.

    Certaines exceptions peuvent provoquer l'échec des TaskManagers, entraînant des journaux incomplets. Examinez les derniers journaux du TaskManager invalide pour résoudre le problème.

  • Consultez les journaux des instances de job historiques.

    Examinez les journaux des instances de job historiques pour identifier la cause de l'échec.

Pourquoi la sortie des données est-elle suspendue sur l'opérateur LocalGroupAggregate ?

  • Code

    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;
  • Description du problème

    La sortie des données est suspendue sur l'opérateur LocalGroupAggregate pendant une période prolongée, et l'opérateur MiniBatchAssigner est absent de la topologie du job.

    Dans l'interface Operator Analysis (Beta), la topologie du job indique que LocalGroupAggregate[201] affiche un RecordsIn de 1853 et un RecordsOut de seulement 116. Le Vertex2 en aval (contenant les opérateurs GlobalGroupAggregate et Calc) présente un Out (sum) de 0 et un Backpressured (max) de 0 %. Aucun nœud MiniBatchAssigner n'apparaît dans l'ensemble de la topologie.

  • Cause

    Le job inclut à la fois des opérateurs WindowAggregate et GroupAggregate. L'opérateur WindowAggregate utilise proctime comme colonne temporelle. La mémoire managée est utilisée pour mettre en cache les données en mode de traitement miniBatch si le paramètre table.exec.mini-batch.size n'est pas configuré ou est défini sur une valeur négative.

    L'opérateur MiniBatchAssigner ne parvient pas à générer et à envoyer des messages watermark aux opérateurs de calcul pour déclencher le calcul final et la sortie des données. Le calcul final et la sortie des données ne sont déclenchés que lorsqu'une des conditions suivantes est remplie : la mémoire managée est pleine, une commande CHECKPOINT est reçue et aucun point de contrôle n'a été effectué, ou le job est annulé. Pour plus d'informations, consultez table.exec.mini-batch.size. Si l'intervalle de point de contrôle est défini sur une valeur excessivement grande, l'opérateur LocalGroupAggregate ne déclenche pas la sortie des données pendant une période prolongée.

  • Solutions

    • Réduisez l'intervalle de point de contrôle afin que l'opérateur LocalGroupAggregate déclenche la sortie des données avant le point de contrôle. Tuning Checkpointing.

    • Utilisez la mémoire heap pour mettre les données en cache. La sortie se déclenche automatiquement lorsque les données mises en cache atteignent la valeur table.exec.mini-batch.size. Définissez ce paramètre sur une valeur positive N. Comment configurer des paramètres d'exécution personnalisés pour un job ?

Que faire si l'inactivité d'une partition Kafka retarde la sortie de la fenêtre ?

Si le connecteur Kafka en amont comporte plusieurs partitions mais que seules certaines reçoivent des données, les partitions inactives empêchent l'avancement du watermark. Les fenêtres ne peuvent pas se fermer rapidement, ce qui retarde la sortie en temps réel.

Configurez un délai d'expiration pour marquer les partitions inactives. Les partitions inactives sont exclues des calculs de watermark jusqu'à ce qu'elles reçoivent à nouveau des données. Configuration.

Ajoutez la configuration suivante au champ Other Configuration dans la section Parameters de l'onglet Configuration. Comment configurer des paramètres d'exécution personnalisés pour un job ?

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

Comment localiser l'erreur si le JobManager n'est pas en cours d'exécution ?

La page Flink UI n'apparaît pas car le JobManager n'est pas en cours d'exécution. Pour identifier la cause, procédez comme suit :

  1. Dans le menu de navigation de gauche de la Development Console, choisissez O&M > Deployments. Sur la page Deployments, recherchez le déploiement du job cible et cliquez sur son nom.

  2. Cliquez sur l'onglet Events.

  3. Recherchez les erreurs à l'aide du raccourci clavier de votre système d'exploitation :

    • Windows : Ctrl+F

    • macOS : Command+F

    Sur la page des détails du job, cliquez sur l'onglet Events pour afficher la liste des événements du cycle de vie du job et les messages d'erreur. Dans la visionneuse de journaux, localisez la ligne d'exception correspondante. Par exemple, un journal WARN à la ligne 52 peut indiquer une erreur d'analyse de la paire clé-valeur pour la clé $internal.application.program-args à la ligne 43 dans /flink/conf/flink-conf.yaml.

Que faire lorsque le message « INFO: org.apache.flink.fs.osshadoop.shaded.com.aliyun.oss » apparaît ?

  • Description du problème

    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
  • Cause

    Les données sont stockées dans un compartiment OSS. Lorsqu'OSS crée un répertoire, il vérifie si celui-ci existe. Si ce n'est pas le cas, ce message INFO est imprimé. Cela n'affecte pas vos jobs.

  • Solution

    Ajoutez la configuration de journaliseur suivante à votre modèle de journal pour supprimer ce message : <Logger level="ERROR" name="org.apache.flink.fs.osshadoop.shaded.com.aliyun.oss"/>. la rubrique sur la configuration des journaux.

Que faire si le message d'erreur « akka.pattern.AskTimeoutException » apparaît ?

  • Causes

    • Cause 1 : Collecte des ordures (GC) fréquente. Une mémoire insuffisante pour le JobManager ou les TaskManagers provoque des GC fréquentes, entraînant des délais d'expiration des signaux heartbeat et des appels RPC entre le JobManager et les TaskManagers.

    • Cause 2 : Volume élevé de requêtes RPC. Un trop grand nombre de requêtes RPC sature le JobManager, provoquant un engorgement des RPC ainsi que des délais d'expiration des signaux heartbeat et des appels RPC.

    • Cause 3 : Valeurs de délai d'expiration trop faibles. Les paramètres de délai d'expiration sont trop bas. Lorsque Realtime Compute for Apache Flink tente de se reconnecter à un service tiers, le délai d'expiration expire avant que l'échec ne soit signalé.

  • Solutions

    • Solution 1 : Vérifiez la fréquence et la durée des GC à partir de l'utilisation de la mémoire du job et des journaux GC. Si les GC sont fréquentes ou longues, augmentez la mémoire du JobManager et des TaskManagers.

    • Solution 2 : Pour gérer un volume élevé de requêtes RPC, augmentez le nombre de cœurs CPU et la taille de la mémoire du JobManager, et définissez les paramètres akka.ask.timeout et heartbeat.timeout sur des valeurs plus élevées.

      Important
      • Ajustez akka.ask.timeout et heartbeat.timeout uniquement en présence d'un grand nombre de requêtes RPC. Pour les jobs avec peu de requêtes RPC, des valeurs plus faisses ne provoquent généralement pas ce problème.

      • Définissez les valeurs en fonction de vos besoins métier. Des valeurs excessivement grandes augmentent le temps de récupération lorsqu'un TaskManager quitte inopinément.

    • Solution 3 : Pour gérer les échecs de connexion aux services tiers, augmentez les paramètres suivants afin que les échecs de connexion soient signalés rapidement :

      • client.timeout : Valeur par défaut : 60. Valeur recommandée : 600. Unité : secondes.

      • akka.ask.timeout : Valeur par défaut : 10. Valeur recommandée : 600. Unité : secondes.

      • client.heartbeat.timeout : Valeur par défaut : 180000. Valeur recommandée : 600000. Unité : secondes.

        Remarque

        Pour éviter les erreurs, n'incluez pas l'unité dans la valeur.

      • heartbeat.timeout : Valeur par défaut : 50000. Valeur recommandée : 600000. Unité : millisecondes.

        Remarque

        Pour éviter les erreurs, n'incluez pas l'unité dans la valeur.

      Par exemple, si le message d'erreur "Caused by: java.sql.SQLTransientConnectionException: connection-pool-xxx.mysql.rds.aliyuncs.com:3306 - Connection is not available, request timed out after 30000ms" apparaît, le pool de connexions MySQL est saturé. Dans ce cas, vous devez augmenter la valeur du paramètre connection.pool.size décrit dans les paramètres WITH de MySQL. Valeur par défaut : 20.

      Remarque

      Déterminez les valeurs minimales à partir du message d'erreur de délai d'expiration. La valeur indiquée dans l'erreur correspond au paramétrage actuel. Par exemple, « 60000 ms » dans "pattern.AskTimeoutException: Ask timed out on [Actor[akka://flink/user/rpc/dispatcher_1#1064915964]] after [60000 ms]." correspond à la valeur client.timeout.

Que faire si le message d'erreur « Task did not exit gracefully within 180 + seconds. » s'affiche ?

  • Description du problème

    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
  • Cause

    Cette erreur n'indique pas la cause racine. Elle signifie que la sortie de la tâche a été bloquée pendant un basculement ou une annulation plus longtemps que la valeur par défaut de task.cancellation.timeout (180 secondes). Realtime Compute for Apache Flink considère la tâche comme irrécupérable, arrête le TaskManager concerné et permet au basculement ou à l'annulation de se poursuivre.

    Ce problème est souvent causé par des fonctions définies par l'utilisateur (UDF). Par exemple, si la méthode close d'une UDF se bloque ou ne renvoie pas de résultat, la tâche ne peut pas se terminer correctement.

  • Solution

    À des fins de débogage, définissez task.cancellation.timeout sur 0. Comment configurer les paramètres d'exécution personnalisés pour un job ? Lorsque ce paramètre est défini sur 0, une tâche bloquée attend indéfiniment sa sortie sans déclencher de délai d'expiration. Si un basculement se déclenche à nouveau ou si une tâche reste bloquée après le redémarrage, localisez la tâche à l'état CANCELLING, inspectez sa trace de pile et corrigez la cause racine.

    Important

    Le paramètre task.cancellation.timeout est réservé au débogage. Ne le définissez pas sur 0 en production. Utilisez un délai d'expiration approprié et corrigez le problème sous-jacent lié à l'UDF ou à la logique métier.

Que faire lorsque le message d'erreur « Can not retract a non-existent record. This should never happen. » s'affiche ?

  • Description du problème

    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)
                        
  • Causes et solutions

    Scénario

    Cause

    Solution

    Scénario 1

    Le problème est causé par la fonction now() dans le code.

    L'algorithme TopN n'autorise pas l'utilisation d'un champ non déterministe dans la clause ORDER BY ou PARTITION BY. Si un champ non déterministe est utilisé, les valeurs renvoyées par la fonction now() diffèrent pour chaque enregistrement et la valeur précédente est introuvable dans l'état.

    Utilisez un champ déterministe dans la clause ORDER BY ou PARTITION BY.

    Scénario 2

    Le paramètre table.exec.state.ttl est défini sur une valeur excessivement faible. Par conséquent, les entrées d'état expirent et sont supprimées, et l'état de clé requis est introuvable.

    Augmentez la valeur de table.exec.state.ttl. Comment configurer les paramètres d'exécution personnalisés pour un job ?

Comment résoudre le message d'erreur « The GRPC call timed out in sqlserver » ?

  • Description du problème

    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
                        
  • Cause

    Un code SQL complexe dans le brouillon provoque un délai d'expiration de l'exécution RPC.

  • Solution

    Ajoutez le code suivant au champ Other Configuration de la section Parameters dans l'onglet Configuration afin d'augmenter le délai d'expiration RPC. La valeur par défaut est de 120 secondes. Pour plus d'informations, consultez la rubrique Configurer les paramètres d'exécution personnalisés

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

Comment résoudre le message d'erreur « RESOURCE_EXHAUSTED: gRPC message exceeds maximum size 41943040: 58384051 » ?

  • Description du problème

    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)
  • Cause

    Le JobGraph est trop volumineux en raison d'une logique de brouillon complexe. Cela entraîne des erreurs de validation ou empêche le démarrage ou l'annulation du job de brouillon.

  • Solution

    Ajoutez le code suivant au champ Other Configuration de la section Parameters dans l'onglet Configuration. Comment configurer les paramètres d'exécution personnalisés pour un job ?

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

Que faire si le message d'erreur « Caused by: java.lang.NoSuchMethodError » s'affiche ?

  • Description du problème

    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;
  • Cause

    Si vous appelez une API Apache Flink et que Realtime Compute for Apache Flink fournit une version optimisée, une exception telle qu'un conflit de package peut se produire.

  • Solution

    Limitez vos appels de méthode à ceux explicitement marqués avec @Public ou @PublicEvolving dans le code source d'Apache Flink. Realtime Compute for Apache Flink garantit la compatibilité avec ces méthodes.

Que faire si le message d'erreur « java.lang.ClassCastException: org.codehaus.janino.CompilerFactory cannot be cast to org.codehaus.commons.compiler.ICompilerFactory » s'affiche ?

  • Description du problème

    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
  • Cause

    • Le package JAR contient une dépendance Janino qui provoque un conflit.

    • Des packages JAR spécifiques commençant par Flink-, tels que flink-table-planner et flink-table-runtime, sont ajoutés au package JAR de l'UDF ou du connecteur.

  • Solutions

    • Vérifiez si le package JAR contient org.codehaus.janino.CompilerFactory. Des conflits de classes peuvent survenir car la séquence de chargement des classes varie selon les machines. Pour résoudre ce problème, procédez comme suit :

      1. Dans le volet de navigation de gauche de la console de développement, choisissez O&M > Deployments. Sur la page Deployments, recherchez le job cible et cliquez sur son nom.

      2. Dans l'onglet Configuration de la page des détails du job, cliquez sur Edit dans le coin supérieur droit de la section Parameters.

      3. Ajoutez le code suivant au champ Other Configuration et cliquez sur Save.

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

        Remplacez la valeur du paramètre classloader.parent-first-patterns.additional par la classe conflictuelle.

    • Spécifiez <scope>provided</scope> pour les dépendances Apache Flink, telles que les dépendances hors connecteurs dont les noms commencent par flink- dans le groupe org.apache.flink.