Hi Flavio! Can you post the JobManager's log here? It should have the message about what is going wrong...
Stephan On Mon, Jun 29, 2015 at 11:43 AM, Flavio Pompermaier <pomperma...@okkam.it> wrote: > Hi to all, > > I'm restarting the discussion about a problem I alredy dicussed on this > mailing list (but that started with a different subject). > I'm running Flink 0.9.0 on CDH 5.1.3 so I compiled the sources as: > > mvn clean install -Dhadoop.version=2.3.0-cdh5.1.3 > -Dhbase.version=0.98.1-cdh5.1.3 -Dhadoop.core.version=2.3.0-mr1-cdh5.1.3 > -DskipTests -Pvendor-repos > > The problem I'm facing is that the cluster start successfully but when I > run my job (from the web-client) I get, after some time, this exception: > > 16:35:41,636 WARN akka.remote.RemoteWatcher > - Detected unreachable: [akka.tcp://flink@192.168.234.83:6123] > 16:35:46,605 INFO org.apache.flink.runtime.taskmanager.TaskManager - > Disconnecting from JobManager: JobManager is no longer reachable > 16:35:46,614 INFO org.apache.flink.runtime.taskmanager.TaskManager - > Cancelling all computations and discarding all cached data. > 16:35:46,644 INFO org.apache.flink.runtime.taskmanager.Task > - Attempting to fail task externally CHAIN GroupReduce (GroupReduce at > compactDataSources(MyClass.java:213)) -> Combine(Distinct at > compactDataSources(MyClass.java:213)) (8/36) > 16:35:46,669 INFO org.apache.flink.runtime.taskmanager.Task > - CHAIN GroupReduce (GroupReduce at compactDataSources(MyClass.java:213)) > -> Combine(Distinct at compactDataSources(MyClass.java:213)) (8/36) > switched to FAILED with exception. > java.lang.Exception: Disconnecting from JobManager: JobManager is no > longer reachable > at org.apache.flink.runtime.taskmanager.TaskManager.org > $apache$flink$runtime$taskmanager$TaskManager$$handleJobManagerDisconnect(TaskManager.scala:741) > at > org.apache.flink.runtime.taskmanager.TaskManager$$anonfun$receiveWithLogMessages$1.applyOrElse(TaskManager.scala:267) > at > scala.runtime.AbstractPartialFunction$mcVL$sp.apply$mcVL$sp(AbstractPartialFunction.scala:33) > at > scala.runtime.AbstractPartialFunction$mcVL$sp.apply(AbstractPartialFunction.scala:33) > at > scala.runtime.AbstractPartialFunction$mcVL$sp.apply(AbstractPartialFunction.scala:25) > at > org.apache.flink.runtime.ActorLogMessages$$anon$1.apply(ActorLogMessages.scala:36) > at > org.apache.flink.runtime.ActorLogMessages$$anon$1.apply(ActorLogMessages.scala:29) > at > scala.PartialFunction$class.applyOrElse(PartialFunction.scala:118) > at > org.apache.flink.runtime.ActorLogMessages$$anon$1.applyOrElse(ActorLogMessages.scala:29) > at akka.actor.Actor$class.aroundReceive(Actor.scala:465) > at > org.apache.flink.runtime.taskmanager.TaskManager.aroundReceive(TaskManager.scala:114) > at akka.actor.ActorCell.receiveMessage(ActorCell.scala:516) > at > akka.actor.dungeon.DeathWatch$class.receivedTerminated(DeathWatch.scala:46) > at akka.actor.ActorCell.receivedTerminated(ActorCell.scala:369) > at akka.actor.ActorCell.autoReceiveMessage(ActorCell.scala:501) > at akka.actor.ActorCell.invoke(ActorCell.scala:486) > at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:254) > at akka.dispatch.Mailbox.run(Mailbox.scala:221) > at akka.dispatch.Mailbox.exec(Mailbox.scala:231) > at > scala.concurrent.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260) > at > scala.concurrent.forkjoin.ForkJoinPool$WorkQueue.pollAndExecAll(ForkJoinPool.java:1253) > at > scala.concurrent.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1346) > at > scala.concurrent.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979) > at > scala.concurrent.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107) > 16:35:46,767 INFO org.apache.flink.runtime.taskmanager.Task > - Triggering cancellation of task code CHAIN GroupReduce (GroupReduce > at compactDataSources(MyClass.java:213)) -> Combine(Distinct at > compactDataSources(MyClass.java:213)) (8/36) > (57a0ad78726d5ba7255aa87038250c51). > > The job instead runs correctly from the IDE (Eclipse). How can I > understand/debug what's wrong? > > Best, > Flavio > >