Erros comuns de execução de jobs e soluções no Realtime Compute for Apache Flink.
O que fazer se os dados nas tarefas de um job não forem consumidos após a execução?
Por que a saída de dados é suspensa no operador LocalGroupAggregate?
O que fazer se a ociosidade de partições do Kafka atrasar a saída da janela?
Como localizar o erro se o JobManager não estiver em execução?
O que fazer quando a mensagem "INFO: org.apache.flink.fs.osshadoop.shaded.com.aliyun.oss" aparecer?
O que fazer se a mensagem de erro "akka.pattern.AskTimeoutException" aparecer?
O que fazer se a mensagem de erro "Task did not exit gracefully within 180 + seconds." aparecer?
Como corrigir a mensagem de erro "The GRPC call timed out in sqlserver"?
O que fazer se a mensagem de erro "Caused by: java.lang.NoSuchMethodError" aparecer?
O que fazer se um job não puder ser iniciado?
-
Descrição do problema
Ao clicar em Start na coluna Actions, o status do job muda de STARTING para FAILED.
-
Soluções
Verifique a aba Events: Acesse a aba Events na página de detalhes do job. Localize o evento de falha ocorrido durante a inicialização e revise seus detalhes para identificar a causa raiz.
Revise os logs de inicialização: Acesse a aba Logs e selecione a subaba Startup Logs. Verifique nos logs mensagens de erro específicas que expliquem a falha na inicialização do job.
Inspecione os logs do JobManager/TaskManager: Se o JobManager parecer iniciar com sucesso, mas o job ainda falhar, verifique os logs detalhados do JobManager e dos TaskManagers. Esses logs estão disponíveis nas subabas Job Manager ou Running Task Managers, dentro da aba Logs.
-
Erros comuns e soluções
Descrição do problema
Causa
Solução
ERROR:exceeded quota: resourcequota
Recursos insuficientes na Resource Queue atual.
Aumente a capacidade da Resource Queue ou reduza os requisitos de recursos do job.
ERROR:the vswitch ip is not enough
Endereços IP insuficientes no namespace para os TaskManagers necessários.
Reduza o paralelismo do job, ajuste a configuração de slots ou modifique as configurações do vSwitch.
ERROR: pooler: *: authentication failed**
Par de AccessKey inválido ou permissões insuficientes.
Verifique se o par de AccessKey é válido e pertence a uma conta com permissões para executar e gerenciar jobs.
Como corrigir o erro de conexão com o banco de dados?
-
Descrição do problema

-
Causa
O catálogo registrado é inválido ou inacessível.
-
Solução
Acesse a página Catalogs, exclua quaisquer catálogos esmaecidos e registre-os novamente.
O que fazer se os dados nas tarefas de um job não forem consumidos após a execução?
-
Verifique a conectividade de rede
Se os dados não forem gerados ou consumidos no armazenamento upstream e downstream, verifique a aba Startup Logs em busca de mensagens de erro. Caso encontre erros de tempo limite, solucione problemas de conectividade de rede entre os sistemas de armazenamento.
-
Verifique o status de execução da tarefa
Na aba Configuration, verifique se os dados estão sendo lidos da source e gravados no sink para identificar onde o erro ocorre.

-
Verifique a saída do operador
Adicione uma tabela sink de impressão a cada operador para solucionar o problema.
O que fazer se um job reiniciar inesperadamente?
Para solucionar o erro, verifique a aba Logs.
-
Visualize as informações de exceção.
Na subaba JM Exceptions, revise o erro relatado e identifique a causa raiz.
-
Visualize os logs do JobManager e do TaskManager do job.

-
Visualize os logs do TaskManager com falha do job.
Algumas exceções podem causar falhas nos TaskManagers, resultando em logs incompletos. Visualize os últimos logs inválidos do TaskManager para solucionar o problema.

-
Visualize os logs de instâncias históricas do job.
Revise os logs de instâncias históricas do job para identificar a causa da falha.

Por que a saída de dados é suspensa no operador LocalGroupAggregate?
-
Código
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; -
Descrição do problema
A saída de dados fica suspensa no operador LocalGroupAggregate por um longo período, e o operador MiniBatchAssigner está ausente na topologia do job.

-
Causa
O job inclui os operadores WindowAggregate e GroupAggregate. O operador WindowAggregate usa proctime como coluna de tempo. A memória gerenciada armazena dados em cache no modo de processamento miniBatch se o parâmetro
table.exec.mini-batch.sizenão estiver configurado ou estiver definido com um valor negativo.O operador MiniBatchAssigner falha ao gerar e não consegue enviar mensagens de watermark para os operadores de computação, impedindo o acionamento do cálculo final e da saída de dados. O cálculo final e a saída de dados são acionados apenas quando uma das seguintes condições é atendida: a memória gerenciada está cheia, um comando CHECKPOINT é recebido sem que o checkpointing tenha sido realizado, ou o job é cancelado. Para obter mais informações, consulte table.exec.mini-batch.size. Se o intervalo de checkpoint estiver definido com um valor excessivamente grande, o operador LocalGroupAggregate não acionará a saída de dados por um longo período.
-
Soluções
Diminua o intervalo de checkpoint para que o operador LocalGroupAggregate acione a saída de dados antes do checkpointing. Tuning Checkpointing.
Use a memória heap para armazenar dados em cache. A saída é acionada automaticamente quando os dados em cache atingem o valor de
table.exec.mini-batch.size. Defina este parâmetro com um valor positivo N. Como configuro parâmetros de tempo de execução personalizados para um job?
O que fazer se a ociosidade de partições do Kafka atrasar a saída da janela?
Se o conector Kafka upstream tiver várias partições, mas apenas algumas receberem dados, as partições ociosas impedirão o avanço do watermark. As janelas não conseguem fechar prontamente, atrasando a saída em tempo real.
Configure um tempo limite para marcar partições ociosas. As partições ociosas são excluídas dos cálculos de watermark até que recebam dados novamente. Configuration.
Adicione a seguinte configuração ao campo Other Configuration na seção Parameters da aba Configuration. Como configuro parâmetros de tempo de execução personalizados para um job?
table.exec.source.idle-timeout: 1s
Como localizar o erro se o JobManager não estiver em execução?
A página Flink UI não aparece porque o JobManager não está em execução. Para identificar a causa, execute as etapas a seguir:
No painel de navegação à esquerda do Development Console, escolha . Na página Deployments, localize a implantação do job de destino e clique em seu nome.
Clique na aba Events.
-
Pesquise erros usando o atalho de teclado do seu sistema operacional:
Windows: Ctrl+F
macOS: Command+F

O que fazer quando a mensagem "INFO: org.apache.flink.fs.osshadoop.shaded.com.aliyun.oss" aparecer?
Descrição do problema

-
Causa
Os dados são armazenados em um bucket do OSS. Quando o OSS cria um diretório, ele verifica se o diretório existe. Caso contrário, esta mensagem INFO é impressa. Isso não afeta seus jobs.
-
Solução
Adicione a seguinte configuração de logger ao seu modelo de log para suprimir esta mensagem:
<Logger level="ERROR" name="org.apache.flink.fs.osshadoop.shaded.com.aliyun.oss"/>. o tópico de configuração de log.
O que fazer se a mensagem de erro "akka.pattern.AskTimeoutException" aparecer?
-
Causas
Causa 1: Coleta de lixo (GC) frequente. Memória insuficiente no JobManager ou TaskManager causa GC frequente, levando a tempos limite de heartbeat e RPC entre o JobManager e os TaskManagers.
Causa 2: Alto volume de solicitações RPC. Muitas solicitações RPC sobrecarregam o JobManager, causando um backlog de RPC e tempos limite de heartbeat e RPC.
Causa 3: Valores de tempo limite muito pequenos. As configurações de tempo limite estão muito baixas. Quando o Realtime Compute for Apache Flink tenta reconectar a um serviço de terceiros, o tempo limite expira antes que a falha seja relatada.
-
Soluções
Solução 1: Verifique a frequência e a duração do GC nos logs de uso de memória e GC do job. Se o GC for frequente ou longo, aumente a memória do JobManager e do TaskManager.
-
Solução 2: Para lidar com um alto volume de solicitações RPC, aumente o número de núcleos de CPU e o tamanho da memória do JobManager, e defina os parâmetros
akka.ask.timeouteheartbeat.timeoutcom valores maiores.ImportanteAjuste
akka.ask.timeouteheartbeat.timeoutapenas quando existir um grande número de solicitações RPC. Para jobs com poucas solicitações RPC, valores menores geralmente não causam esse problema.Defina valores com base nos seus requisitos de negócio. Valores excessivamente grandes aumentam o tempo de recuperação quando um TaskManager sai inesperadamente.
-
Solução 3: Para lidar com falhas de conexão de serviços de terceiros, aumente os seguintes parâmetros para que as falhas de conexão sejam relatadas prontamente:
client.timeout: Valor padrão: 60. Valor recomendado: 600. Unidade: segundos.akka.ask.timeout: Valor padrão: 10. Valor recomendado: 600. Unidade: segundos.-
client.heartbeat.timeout: Valor padrão: 180000. Valor recomendado: 600000. Unidade: segundos.NotaPara evitar erros, não inclua a unidade no valor.
-
heartbeat.timeout: Valor padrão: 50000. Valor recomendado: 600000. Unidade: milissegundos.NotaPara evitar erros, não inclua a unidade no valor.
Por exemplo, se a mensagem de erro
"Caused by: java.sql.SQLTransientConnectionException: connection-pool-xxx.mysql.rds.aliyuncs.com:3306 - Connection is not available, request timed out after 30000ms"aparecer, o pool de conexões do MySQL está cheio. Nesse caso, aumente o valor do parâmetroconnection.pool.sizedescrito nos parâmetros WITH do MySQL. Valor padrão: 20.NotaDetermine os valores mínimos a partir da mensagem de erro de tempo limite. O valor mostrado no erro indica a configuração atual. Por exemplo, "60000 ms" em
"pattern.AskTimeoutException: Ask timed out on [Actor[akka://flink/user/rpc/dispatcher_1#1064915964]] after [60000 ms]."é o valor declient.timeout.
O que fazer se a mensagem de erro "Task did not exit gracefully within 180 + seconds." aparecer?
-
Descrição do problema
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 -
Causa
Este erro não indica a causa raiz. Ele significa que a saída da tarefa ficou travada durante o failover ou cancelamento por mais tempo do que o padrão de
task.cancellation.timeoutde 180 segundos. O Realtime Compute for Apache Flink trata a tarefa como irrecuperável, interrompe o TaskManager afetado e permite que o failover ou cancelamento continue.Isso geralmente é causado por funções definidas pelo usuário (UDFs). Por exemplo, se o método close em uma UDF bloquear ou não retornar, a tarefa não poderá sair.
-
Solução
Para depuração, defina
task.cancellation.timeoutcomo 0. Como configuro parâmetros de tempo de execução personalizados para um job? Quando definido como 0, uma tarefa bloqueada aguarda indefinidamente para sair sem acionar um tempo limite. Se o failover for acionado novamente ou uma tarefa permanecer travada após a reinicialização, localize a tarefa no estado CANCELLING, inspecione seu stack trace e corrija a causa raiz.ImportanteO parâmetro
task.cancellation.timeoutserve apenas para depuração. Não o defina como 0 em produção. Use um tempo limite apropriado e corrija o problema subjacente da UDF ou da lógica de negócio.
O que fazer quando a mensagem de erro "Can not retract a non-existent record. This should never happen." aparecer?
-
Descrição do problema
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) -
Causas e soluções
Cenário
Causa
Solução
Cenário 1
O problema é causado pela função
now()no código.O algoritmo TopN não permite que um campo não determinístico seja usado na cláusula ORDER BY ou PARTITION BY. Se um campo não determinístico for usado, os valores retornados pela função
now()serão diferentes para cada registro, e o valor anterior não poderá ser encontrado no estado.Use um campo determinístico na cláusula ORDER BY ou PARTITION BY.
Cenário 2
O parâmetro
table.exec.state.ttlestá definido com um valor excessivamente pequeno. Como resultado, as entradas de estado expiram e são excluídas, e o estado de chave necessário não pode ser encontrado no estado.Aumente o valor de
table.exec.state.ttl. Como configuro parâmetros de tempo de execução personalizados para um job?
Como corrigir a mensagem de erro "The GRPC call timed out in sqlserver"?
-
Descrição do problema
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 -
Causa
SQL complexo no rascunho causa tempo limite de execução de RPC.
-
Solução
Adicione o seguinte código ao campo Other Configuration na seção Parameters da aba Configuration para aumentar o tempo limite de RPC. O padrão é 120 segundos. Para obter mais informações, consulte Configurar parâmetros de execução personalizados
flink.sqlserver.rpc.execution.timeout: 600s
Como resolver a mensagem de erro "RESOURCE_EXHAUSTED: gRPC message exceeds maximum size 41943040: 58384051"?
-
Descrição do problema
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) -
Causa
O JobGraph é muito grande devido à lógica complexa do rascunho. Isso causa erros de verificação ou impede que o job de rascunho seja iniciado ou cancelado.
-
Solução
Adicione o seguinte código ao campo Other Configuration na seção Parameters da aba Configuration. Como configuro parâmetros de tempo de execução personalizados para um job?
table.exec.operator-name.max-length: 1000
O que fazer se a mensagem de erro "Caused by: java.lang.NoSuchMethodError" aparecer?
-
Descrição do problema
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; -
Causa
Se você chamar uma API do Apache Flink e o Realtime Compute for Apache Flink fornecer uma versão otimizada, uma exceção como conflito de pacotes poderá ocorrer.
-
Solução
Restrinja suas chamadas de método àquelas explicitamente marcadas com @Public ou @PublicEvolving no código-fonte do Apache Flink. O Realtime Compute for Apache Flink garante compatibilidade com esses métodos.
O que fazer se a mensagem de erro "java.lang.ClassCastException: org.codehaus.janino.CompilerFactory cannot be cast to org.codehaus.commons.compiler.ICompilerFactory" aparecer?
-
Descrição do problema
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 -
Causa
O pacote JAR contém uma dependência Janino que causa conflito.
Pacotes JAR específicos que começam com
Flink-, comoflink-table-plannereflink-table-runtime, foram adicionados ao pacote JAR da UDF ou do conector.
-
Soluções
-
Verifique se o pacote JAR contém
org.codehaus.janino.CompilerFactory. Conflitos de classe podem ocorrer porque a sequência de carregamento de classes varia entre máquinas. Para resolver esse problema, execute as etapas a seguir:No painel de navegação à esquerda do Development Console, escolha . Na página Deployments, localize o job de destino e clique em seu nome.
Na aba Configuration da página de detalhes do job, clique em Edit no canto superior direito da seção Parameters.
-
Adicione o seguinte código ao campo Other Configuration e clique em Save.
classloader.parent-first-patterns.additional: org.codehaus.janinoSubstitua o valor do parâmetro
classloader.parent-first-patterns.additionalpela classe em conflito.
Especifique
<scope>provided</scope>para dependências do Apache Flink, como dependências não-conector cujos nomes começam comflink-no grupoorg.apache.flink.
-