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@01DA9B10.1A7C32A0]<https://bit.ly/3VNHMS5>