Update:
I figured out that error happens only of the SQL contains JOIN of any kind. If
there are no JOINs, everything works fine.
Any help?
Hello,
I am trying to use CloseableIterator, but next() operation reliably generates
the following error:
java.util.concurrent.ExecutionException: org.apache.flink.util.FlinkException:
Coordinator of operator 4596fb32cad14208ec80c1cae8623e11 does not exist or the
job vertex this operator belongs to is not initialized.
at
java.base/java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:395)
~[na:na]
at
java.base/java.util.concurrent.CompletableFuture.get(CompletableFuture.java:2005)
~[na:na]
at
org.apache.flink.streaming.api.operators.collect.CollectResultFetcher.sendRequest(CollectResultFetcher.java:171)
~[flink-streaming-java-1.18.0.jar:1.18.0]
at
org.apache.flink.streaming.api.operators.collect.CollectResultFetcher.next(CollectResultFetcher.java:129)
~[flink-streaming-java-1.18.0.jar:1.18.0]
at
org.apache.flink.streaming.api.operators.collect.CollectResultIterator.nextResultFromFetcher(CollectResultIterator.java:106)
~[flink-streaming-java-1.18.0.jar:1.18.0]
at
org.apache.flink.streaming.api.operators.collect.CollectResultIterator.next(CollectResultIterator.java:88)
~[flink-streaming-java-1.18.0.jar:1.18.0]
at
org.apache.flink.table.planner.connectors.CollectDynamicSink$CloseableRowIteratorWrapper.next(CollectDynamicSink.java:229)
~[flink-table-planner_2.12-1.18.0.jar:1.18.0]
at
com.intradiem.service.flink.job.UserSnapshotJob.createSnapshot(UserSnapshotJob.java:108)
~[classes/:na]
at
com.intradiem.service.quartz.TriggerUserSnapshot.execute(TriggerUserSnapshot.java:68)
~[classes/:na]
at org.quartz.core.JobRunShell.run(JobRunShell.java:202)
~[quartz-2.3.2.jar:na]
at
org.quartz.simpl.SimpleThreadPool$WorkerThread.run(SimpleThreadPool.java:573)
~[quartz-2.3.2.jar:na]
Caused by: org.apache.flink.util.FlinkException: Coordinator of operator
4596fb32cad14208ec80c1cae8623e11 does not exist or the job vertex this operator
belongs to is not initialized.
at
org.apache.flink.runtime.scheduler.DefaultOperatorCoordinatorHandler.deliverCoordinationRequestToCoordinator(DefaultOperatorCoordinatorHandler.java:135)
~[flink-runtime-1.18.0.jar:1.18.0]
at
org.apache.flink.runtime.scheduler.SchedulerBase.deliverCoordinationRequestToCoordinator(SchedulerBase.java:1070)
~[flink-runtime-1.18.0.jar:1.18.0]
at
org.apache.flink.runtime.jobmaster.JobMaster.sendRequestToCoordinator(JobMaster.java:616)
~[flink-runtime-1.18.0.jar:1.18.0]
at
org.apache.flink.runtime.jobmaster.JobMaster.deliverCoordinationRequestToCoordinator(JobMaster.java:937)
~[flink-runtime-1.18.0.jar:1.18.0]
at
jdk.internal.reflect.GeneratedMethodAccessor62.invoke(Unknown Source) ~[na:na]
at
java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
~[na:na]
at java.base/java.lang.reflect.Method.invoke(Method.java:566)
~[na:na]
at
org.apache.flink.runtime.rpc.pekko.PekkoRpcActor.lambda$handleRpcInvocation$1(PekkoRpcActor.java:309)
~[na:na]
at
org.apache.flink.runtime.concurrent.ClassLoadingUtils.runWithContextClassLoader(ClassLoadingUtils.java:83)
~[flink-rpc-core-1.18.0.jar:1.18.0]
at
org.apache.flink.runtime.rpc.pekko.PekkoRpcActor.handleRpcInvocation(PekkoRpcActor.java:307)
~[na:na]
at
org.apache.flink.runtime.rpc.pekko.PekkoRpcActor.handleRpcMessage(PekkoRpcActor.java:222)
~[na:na]
at
org.apache.flink.runtime.rpc.pekko.FencedPekkoRpcActor.handleRpcMessage(FencedPekkoRpcActor.java:85)
~[na:na]
at
org.apache.flink.runtime.rpc.pekko.PekkoRpcActor.handleMessage(PekkoRpcActor.java:168)
~[na:na]
at
org.apache.pekko.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:33)
~[na:na]
at
org.apache.pekko.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:29)
~[na:na]
at scala.PartialFunction.applyOrElse(PartialFunction.scala:127)
~[scala-library-2.12.7.jar:na]
at
scala.PartialFunction.applyOrElse$(PartialFunction.scala:126)
~[scala-library-2.12.7.jar:na]
at
org.apache.pekko.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:29)
~[na:na]
at
scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:175)
~[scala-library-2.12.7.jar:na]
at
scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:176)
~[scala-library-2.12.7.jar:na]
at
scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:176)
~[scala-library-2.12.7.jar:na]
at org.apache.pekko.actor.Actor.aroundReceive(Actor.scala:547)
~[na:na]
at org.apache.pekko.actor.Actor.aroundReceive$(Actor.scala:545)
~[na:na]
at
org.apache.pekko.actor.AbstractActor.aroundReceive(AbstractActor.scala:229)
~[na:na]
at
org.apache.pekko.actor.ActorCell.receiveMessage(ActorCell.scala:590) ~[na:na]
at org.apache.pekko.actor.ActorCell.invoke(ActorCell.scala:557)
~[na:na]
at
org.apache.pekko.dispatch.Mailbox.processMailbox(Mailbox.scala:280) ~[na:na]
at org.apache.pekko.dispatch.Mailbox.run(Mailbox.scala:241)
~[na:na]
at org.apache.pekko.dispatch.Mailbox.exec(Mailbox.scala:253)
~[na:na]
at
java.base/java.util.concurrent.ForkJoinTask.doExec(ForkJoinTask.java:290)
~[na:na]
at
java.base/java.util.concurrent.ForkJoinPool$WorkQueue.topLevelExec(ForkJoinPool.java:1020)
~[na:na]
at
java.base/java.util.concurrent.ForkJoinPool.scan(ForkJoinPool.java:1656)
~[na:na]
at
java.base/java.util.concurrent.ForkJoinPool.runWorker(ForkJoinPool.java:1594)
~[na:na]
at
java.base/java.util.concurrent.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:183)
~[na:na]
Code that generates exception:
String snapshotSQl = <Some hardcoded SQL string>
org.apache.flink.table.api.Table snapshotTable = tableEnv.sqlQuery(snapshotSQl);
snapshotTable.execute().collect().next();
Any help will be greatly appreciated.
Ed Skhisov
Architect | www.intradiem.com
<https://www.intradiem.com/>303.588.2518
Mailing Address: 2500 Dallas Hwy Ste 202, Dept #37049 Marietta, GA 30064
[cid:[email protected]]<https://bit.ly/3VNHMS5>