Hi Biao,
I figured that the error happens only when there is a JOIN in the select.  But 
I will put together a simple example.

Thank you,
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:image001.png@01DAA062.D9330690]<https://bit.ly/3VNHMS5>

From: Biao Geng <biaoge...@gmail.com>
Sent: Monday, May 6, 2024 8:20 PM
To: Eduard Skhisov <eduard.skhi...@intradiem.com>
Cc: user@flink.apache.org
Subject: [EXTERNAL] Re: Coordinator of operator ... does not exist or the job 
vertex this operator belongs to is not initialized.

Hi Ed,
Would you mind giving a minimal example to reproduce your case?
I tried a pretty simple case like this in a mini cluster:
```
        tEnv.createTemporaryView("test", env.fromData(1, 2, 3));
        Table table = tEnv.sqlQuery("SELECT * FROM test");
        table.execute().collect().next();
```
But I failed to reproduce the exception you attached :(

Best,
Biao Geng


Eduard Skhisov via user <user@flink.apache.org<mailto:user@flink.apache.org>> 
于2024年5月1日周三 05:09写道:
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:image001.png@01DAA062.D9330690]<https://bit.ly/3VNHMS5>

Reply via email to