Erreurs d'exécution courantes des jobs et solutions dans Realtime Compute for Apache Flink.
Comment résoudre une erreur de connexion à la base de données ?
Que faire si les données des tâches d'un job ne sont pas consommées après l'exécution du job ?
Pourquoi la sortie des données est-elle suspendue sur l'opérateur LocalGroupAggregate ?
Que faire si l'inactivité d'une partition Kafka retarde la sortie de la fenêtre ?
Comment localiser l'erreur si le JobManager n'est pas en cours d'exécution ?
Que faire si le message d'erreur « akka.pattern.AskTimeoutException » apparaît ?
Que faire si le message d'erreur « Task did not exit gracefully within 180 + seconds. » apparaît ?
Comment résoudre le message d'erreur « The GRPC call timed out in sqlserver » ?
Que faire si le message d'erreur « Caused by: java.lang.NoSuchMethodError » apparaît ?
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.sizen'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 :
Dans le menu de navigation de gauche de la Development Console, choisissez . Sur la page Deployments, recherchez le déploiement du job cible et cliquez sur son nom.
Cliquez sur l'onglet Events.
-
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.timeoutetheartbeat.timeoutsur des valeurs plus élevées.ImportantAjustez
akka.ask.timeoutetheartbeat.timeoutuniquement 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.RemarquePour éviter les erreurs, n'incluez pas l'unité dans la valeur.
-
heartbeat.timeout: Valeur par défaut : 50000. Valeur recommandée : 600000. Unité : millisecondes.RemarquePour é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ètreconnection.pool.sizedécrit dans les paramètres WITH de MySQL. Valeur par défaut : 20.RemarqueDé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 valeurclient.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.timeoutsur 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.ImportantLe paramètre
task.cancellation.timeoutest 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.ttlest 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 queflink-table-planneretflink-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 :Dans le volet de navigation de gauche de la console de développement, choisissez . Sur la page Deployments, recherchez le job cible et cliquez sur son nom.
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.
-
Ajoutez le code suivant au champ Other Configuration et cliquez sur Save.
classloader.parent-first-patterns.additional: org.codehaus.janinoRemplacez la valeur du paramètre
classloader.parent-first-patterns.additionalpar 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 parflink-dans le groupeorg.apache.flink.
-