Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Job errors FAQ

Última atualização: Jun 27, 2026

Erros comuns de execução de jobs e soluções no Realtime Compute for Apache Flink.

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

    image

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

    image

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

    image

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

    image

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

    image

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.

    image

  • 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.size nã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

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:

  1. No painel de navegação à esquerda do Development Console, escolha O&M > Deployments. Na página Deployments, localize a implantação do job de destino e clique em seu nome.

  2. Clique na aba Events.

  3. Pesquise erros usando o atalho de teclado do seu sistema operacional:

    • Windows: Ctrl+F

    • macOS: Command+F

    example

O que fazer quando a mensagem "INFO: org.apache.flink.fs.osshadoop.shaded.com.aliyun.oss" aparecer?

  • Descrição do problemaerror details

  • 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.timeout e heartbeat.timeout com valores maiores.

      Importante
      • Ajuste akka.ask.timeout e heartbeat.timeout apenas 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.

        Nota

        Para evitar erros, não inclua a unidade no valor.

      • heartbeat.timeout: Valor padrão: 50000. Valor recomendado: 600000. Unidade: milissegundos.

        Nota

        Para 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âmetro connection.pool.size descrito nos parâmetros WITH do MySQL. Valor padrão: 20.

      Nota

      Determine 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 de client.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.timeout de 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.timeout como 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.

    Importante

    O parâmetro task.cancellation.timeout serve 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.ttl está 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-, como flink-table-planner e flink-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:

      1. No painel de navegação à esquerda do Development Console, escolha O&M > Deployments. Na página Deployments, localize o job de destino e clique em seu nome.

      2. Na aba Configuration da página de detalhes do job, clique em Edit no canto superior direito da seção Parameters.

      3. Adicione o seguinte código ao campo Other Configuration e clique em Save.

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

        Substitua o valor do parâmetro classloader.parent-first-patterns.additional pela classe em conflito.

    • Especifique <scope>provided</scope> para dependências do Apache Flink, como dependências não-conector cujos nomes começam com flink- no grupo org.apache.flink.