Hi Eduard,

You may need to set log level = INFO to see if there are any other error
messages generated in the JM or TM's log. The current exception message
seems to be a result error generated from the JM, but the causing error
message should still be lying somewhere in the TM's log.

Best
Yunfeng

On Fri, May 3, 2024 at 2:01 AM Eduard Skhisov via user <
user@flink.apache.org> wrote:

> 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
> 303.588.2518
>
> Mailing Address: 2500 Dallas Hwy Ste 202, Dept #37049 Marietta, GA 30064
>
>
>
>
>
> <https://bit.ly/3VNHMS5>
>
>
>

Reply via email to