报错信息:

Exception in thread “main” org.apache.flink.table.api.TableException: Failed to execute sql
at org.apache.flink.table.api.internal.TableEnvironmentImpl.executeQueryOperation(TableEnvironmentImpl.java:810)
at org.apache.flink.table.api.internal.TableEnvironmentImpl.executeInternal(TableEnvironmentImpl.java:1223)
at org.apache.flink.table.api.internal.TableImpl.execute(TableImpl.java:577)
at org.caigou._1_cdc_mysql._2.main(_2.java:1159)
Caused by: org.apache.flink.util.FlinkException: Failed to execute job ‘collect’.
at org.apache.flink.streaming.api.environment.StreamExecutionEnvironment.executeAsync(StreamExecutionEnvironment.java:1970)
at org.apache.flink.table.planner.delegation.ExecutorBase.executeAsync(ExecutorBase.java:55)
at org.apache.flink.table.api.internal.TableEnvironmentImpl.executeQueryOperation(TableEnvironmentImpl.java:793)
… 3 more
Caused by: java.lang.RuntimeException: Error while waiting for job to be initialized
at org.apache.flink.client.ClientUtils.waitUntilJobInitializationFinished(ClientUtils.java:160)
at org.apache.flink.client.program.PerJobMiniClusterFactory.lambda$submitJob2(PerJobMiniClusterFactory.java:83)atorg.apache.flink.util.function.FunctionUtils.lambda2(PerJobMiniClusterFactory.java:83) at org.apache.flink.util.function.FunctionUtils.lambda2(PerJobMiniClusterFactory.java:83)atorg.apache.flink.util.function.FunctionUtils.lambdauncheckedFunction2(FunctionUtils.java:73)atjava.base/java.util.concurrent.CompletableFuture2(FunctionUtils.java:73) at java.base/java.util.concurrent.CompletableFuture2(FunctionUtils.java:73)atjava.base/java.util.concurrent.CompletableFutureUniApply.tryFire(CompletableFuture.java:646)
at java.base/java.util.concurrent.CompletableFutureCompletion.exec(CompletableFuture.java:483)atjava.base/java.util.concurrent.ForkJoinTask.doExec(ForkJoinTask.java:373)atjava.base/java.util.concurrent.ForkJoinPoolCompletion.exec(CompletableFuture.java:483) at java.base/java.util.concurrent.ForkJoinTask.doExec(ForkJoinTask.java:373) at java.base/java.util.concurrent.ForkJoinPoolCompletion.exec(CompletableFuture.java:483)atjava.base/java.util.concurrent.ForkJoinTask.doExec(ForkJoinTask.java:373)atjava.base/java.util.concurrent.ForkJoinPoolWorkQueue.topLevelExec(ForkJoinPool.java:1182)
at java.base/java.util.concurrent.ForkJoinPool.scan(ForkJoinPool.java:1655)
at java.base/java.util.concurrent.ForkJoinPool.runWorker(ForkJoinPool.java:1622)
at java.base/java.util.concurrent.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:165)
Caused by: java.util.concurrent.ExecutionException: java.util.concurrent.TimeoutException: Invocation of public default java.util.concurrent.CompletableFuture org.apache.flink.runtime.webmonitor.RestfulGateway.requestJobStatus(org.apache.flink.api.common.JobID,org.apache.flink.api.common.time.Time) timed out.
at java.base/java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:396)
at java.base/java.util.concurrent.CompletableFuture.get(CompletableFuture.java:2073)
at org.apache.flink.client.program.PerJobMiniClusterFactory.lambda$null0(PerJobMiniClusterFactory.java:89)atorg.apache.flink.client.ClientUtils.waitUntilJobInitializationFinished(ClientUtils.java:144)...9moreCausedby:java.util.concurrent.TimeoutException:Invocationofpublicdefaultjava.util.concurrent.CompletableFutureorg.apache.flink.runtime.webmonitor.RestfulGateway.requestJobStatus(org.apache.flink.api.common.JobID,org.apache.flink.api.common.time.Time)timedout.atjdk.proxy2/jdk.proxy2.0(PerJobMiniClusterFactory.java:89) at org.apache.flink.client.ClientUtils.waitUntilJobInitializationFinished(ClientUtils.java:144) ... 9 more Caused by: java.util.concurrent.TimeoutException: Invocation of public default java.util.concurrent.CompletableFuture org.apache.flink.runtime.webmonitor.RestfulGateway.requestJobStatus(org.apache.flink.api.common.JobID,org.apache.flink.api.common.time.Time) timed out. at jdk.proxy2/jdk.proxy2.0(PerJobMiniClusterFactory.java:89)atorg.apache.flink.client.ClientUtils.waitUntilJobInitializationFinished(ClientUtils.java:144)...9moreCausedby:java.util.concurrent.TimeoutException:Invocationofpublicdefaultjava.util.concurrent.CompletableFutureorg.apache.flink.runtime.webmonitor.RestfulGateway.requestJobStatus(org.apache.flink.api.common.JobID,org.apache.flink.api.common.time.Time)timedout.atjdk.proxy2/jdk.proxy2.Proxy175.requestJobStatus(Unknown Source)
at org.apache.flink.runtime.minicluster.MiniCluster.lambda$getJobStatus6(MiniCluster.java:704)atjava.base/java.util.concurrent.CompletableFuture.uniApplyNow(CompletableFuture.java:684)atjava.base/java.util.concurrent.CompletableFuture.uniApplyStage(CompletableFuture.java:662)atjava.base/java.util.concurrent.CompletableFuture.thenApply(CompletableFuture.java:2168)atorg.apache.flink.runtime.minicluster.MiniCluster.runDispatcherCommand(MiniCluster.java:751)atorg.apache.flink.runtime.minicluster.MiniCluster.getJobStatus(MiniCluster.java:703)atorg.apache.flink.client.program.PerJobMiniClusterFactory.lambda6(MiniCluster.java:704) at java.base/java.util.concurrent.CompletableFuture.uniApplyNow(CompletableFuture.java:684) at java.base/java.util.concurrent.CompletableFuture.uniApplyStage(CompletableFuture.java:662) at java.base/java.util.concurrent.CompletableFuture.thenApply(CompletableFuture.java:2168) at org.apache.flink.runtime.minicluster.MiniCluster.runDispatcherCommand(MiniCluster.java:751) at org.apache.flink.runtime.minicluster.MiniCluster.getJobStatus(MiniCluster.java:703) at org.apache.flink.client.program.PerJobMiniClusterFactory.lambda6(MiniCluster.java:704)atjava.base/java.util.concurrent.CompletableFuture.uniApplyNow(CompletableFuture.java:684)atjava.base/java.util.concurrent.CompletableFuture.uniApplyStage(CompletableFuture.java:662)atjava.base/java.util.concurrent.CompletableFuture.thenApply(CompletableFuture.java:2168)atorg.apache.flink.runtime.minicluster.MiniCluster.runDispatcherCommand(MiniCluster.java:751)atorg.apache.flink.runtime.minicluster.MiniCluster.getJobStatus(MiniCluster.java:703)atorg.apache.flink.client.program.PerJobMiniClusterFactory.lambdanullKaTeX parse error: Expected 'EOF', got '#' at position 153: …pc/dispatcher_2#̲673109285]] aft…$anonfun2.apply(AskSupport.scala:635)atakka.pattern.PromiseActorRef2.apply(AskSupport.scala:635) at akka.pattern.PromiseActorRef2.apply(AskSupport.scala:635)atakka.pattern.PromiseActorRef$anonfun2.apply(AskSupport.scala:635)atakka.pattern.PromiseActorRef2.apply(AskSupport.scala:635) at akka.pattern.PromiseActorRef2.apply(AskSupport.scala:635)atakka.pattern.PromiseActorRef$anonfun1.apply1.apply1.applymcVsp(AskSupport.scala:648)atakka.actor.Schedulersp(AskSupport.scala:648) at akka.actor.Schedulersp(AskSupport.scala:648)atakka.actor.Scheduler$anon4.run(Scheduler.scala:205)atscala.concurrent.Future4.run(Scheduler.scala:205) at scala.concurrent.Future4.run(Scheduler.scala:205)atscala.concurrent.FutureInternalCallbackExecutor.unbatchedExecute(Future.scala:601)atscala.concurrent.BatchingExecutor.unbatchedExecute(Future.scala:601) at scala.concurrent.BatchingExecutor.unbatchedExecute(Future.scala:601)atscala.concurrent.BatchingExecutorclass.execute(BatchingExecutor.scala:109)
at scala.concurrent.FutureInternalCallbackExecutorInternalCallbackExecutorInternalCallbackExecutor.execute(Future.scala:599)
at akka.actor.LightArrayRevolverSchedulerTaskHolder.executeTask(LightArrayRevolverScheduler.scala:328)atakka.actor.LightArrayRevolverSchedulerTaskHolder.executeTask(LightArrayRevolverScheduler.scala:328) at akka.actor.LightArrayRevolverSchedulerTaskHolder.executeTask(LightArrayRevolverScheduler.scala:328)atakka.actor.LightArrayRevolverScheduler$anon$4.executeBucket1(LightArrayRevolverScheduler.scala:279)atakka.actor.LightArrayRevolverScheduler1(LightArrayRevolverScheduler.scala:279) at akka.actor.LightArrayRevolverScheduler1(LightArrayRevolverScheduler.scala:279)atakka.actor.LightArrayRevolverScheduler$anon4.nextTick(LightArrayRevolverScheduler.scala:283)atakka.actor.LightArrayRevolverScheduler4.nextTick(LightArrayRevolverScheduler.scala:283) at akka.actor.LightArrayRevolverScheduler4.nextTick(LightArrayRevolverScheduler.scala:283)atakka.actor.LightArrayRevolverScheduler$anon$4.run(LightArrayRevolverScheduler.scala:235)
at java.base/java.lang.Thread.run(Thread.java:840)

解决方案:添加2个配置

Configuration conf = new Configuration();
//设置WebUI绑定的本地端口
 conf.setString(RestOptions.BIND_PORT, "8081");
// 设置 akka.ask.timeout 参数为 100s
conf.setString("akka.ask.timeout", "100s");
//创建执行环境
StreamExecutionEnvironment env =
                StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(conf);
Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐