Repository: incubator-zeppelin Updated Branches: refs/heads/master 18e599162 -> 80edf58ad
[ZEPPELIN-433] Bump Flink version to 0.10 Bumps the Flink version to the latest released version 0.10. The new version contains many bug fixes which make the runtime more stable. Author: Till Rohrmann <[email protected]> Closes #442 from tillrohrmann/flink-0.10 and squashes the following commits: 4f00305 [Till Rohrmann] [ZEPPELIN-433] Bump Flink version to 0.10 Project: http://git-wip-us.apache.org/repos/asf/incubator-zeppelin/repo Commit: http://git-wip-us.apache.org/repos/asf/incubator-zeppelin/commit/80edf58a Tree: http://git-wip-us.apache.org/repos/asf/incubator-zeppelin/tree/80edf58a Diff: http://git-wip-us.apache.org/repos/asf/incubator-zeppelin/diff/80edf58a Branch: refs/heads/master Commit: 80edf58adcff163356c435f192b0389f165591b1 Parents: 18e5991 Author: Till Rohrmann <[email protected]> Authored: Tue Nov 17 14:48:17 2015 +0100 Committer: Lee moon soo <[email protected]> Committed: Wed Nov 18 20:02:54 2015 +0900 ---------------------------------------------------------------------- flink/pom.xml | 2 +- .../org/apache/zeppelin/flink/FlinkInterpreter.java | 13 +++++++------ 2 files changed, 8 insertions(+), 7 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/incubator-zeppelin/blob/80edf58a/flink/pom.xml ---------------------------------------------------------------------- diff --git a/flink/pom.xml b/flink/pom.xml index 27d8034..3a8c36c 100644 --- a/flink/pom.xml +++ b/flink/pom.xml @@ -35,7 +35,7 @@ <url>http://zeppelin.incubator.apache.org</url> <properties> - <flink.version>0.9.0</flink.version> + <flink.version>0.10.0</flink.version> <flink.akka.version>2.3.7</flink.akka.version> <flink.scala.binary.version>2.10</flink.scala.binary.version> <flink.scala.version>2.10.4</flink.scala.version> http://git-wip-us.apache.org/repos/asf/incubator-zeppelin/blob/80edf58a/flink/src/main/java/org/apache/zeppelin/flink/FlinkInterpreter.java ---------------------------------------------------------------------- diff --git a/flink/src/main/java/org/apache/zeppelin/flink/FlinkInterpreter.java b/flink/src/main/java/org/apache/zeppelin/flink/FlinkInterpreter.java index b106c7d..8a022bd 100644 --- a/flink/src/main/java/org/apache/zeppelin/flink/FlinkInterpreter.java +++ b/flink/src/main/java/org/apache/zeppelin/flink/FlinkInterpreter.java @@ -20,8 +20,6 @@ package org.apache.zeppelin.flink; import java.io.BufferedReader; import java.io.ByteArrayOutputStream; import java.io.File; -import java.io.IOException; -import java.io.OutputStream; import java.io.PrintStream; import java.io.PrintWriter; import java.net.URL; @@ -31,7 +29,6 @@ import java.util.List; import java.util.Map; import java.util.Properties; -import org.apache.flink.api.java.ExecutionEnvironment; import org.apache.flink.api.scala.FlinkILoop; import org.apache.flink.configuration.Configuration; import org.apache.flink.runtime.minicluster.LocalFlinkMiniCluster; @@ -46,7 +43,6 @@ import org.slf4j.LoggerFactory; import scala.Console; import scala.None; -import scala.Option; import scala.Some; import scala.runtime.AbstractFunction0; import scala.tools.nsc.Settings; @@ -137,7 +133,7 @@ public class FlinkInterpreter extends Interpreter { private int getPort() { if (localMode()) { - return localFlinkCluster.getJobManagerRPCPort(); + return localFlinkCluster.getLeaderRPCPort(); } else { return Integer.parseInt(getProperty("port")); } @@ -332,7 +328,12 @@ public class FlinkInterpreter extends Interpreter { private void startFlinkMiniCluster() { localFlinkCluster = new LocalFlinkMiniCluster(flinkConf, false); - localFlinkCluster.waitForTaskManagersToBeRegistered(); + + try { + localFlinkCluster.start(true); + } catch (Exception e){ + throw new RuntimeException("Could not start Flink mini cluster.", e); + } } private void stopFlinkMiniCluster() {
