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