diff --git a/README.md b/README.md
index 81cbcd7..53a70fe 100644
--- a/README.md
+++ b/README.md
@@ -130,7 +130,7 @@ Or on a YARN cluster:
# Building
-To retrieve the latest development version (0.4.0-SNAPSHOT) of the Cascading Connector for Apache Flink™, run the following command
+To retrieve the latest development version (0.4.1-SNAPSHOT) of the Cascading Connector for Apache Flink™, run the following command
git clone https://github.com/dataArtisans/cascading-flink.git
diff --git a/pom.xml b/pom.xml
index 45d9c62..8213b7e 100644
--- a/pom.xml
+++ b/pom.xml
@@ -23,7 +23,7 @@ limitations under the License.
com.data-artisans
cascading-flink
- 0.4.0-SNAPSHOT
+ 0.4.1-SNAPSHOT
Cascading on Flink
jar
@@ -52,9 +52,11 @@ limitations under the License.
- 3.1.0
- 1.0.3
- 1.7.7
+ 3.1.2
+ 1.18.1
+ 1.7.32
+ 2.17.1
+ 2.10.0
UTF-8
UTF-8
@@ -68,27 +70,21 @@ limitations under the License.
-
-
- conjars.org
- https://conjars.org/repo
-
-
-
-
-
+
@@ -107,7 +103,7 @@ limitations under the License.
org.apache.flink
- flink-clients_2.10
+ flink-clients
${flink.version}
@@ -117,6 +113,12 @@ limitations under the License.
${flink.version}
+
+ org.apache.flink
+ flink-hadoop-compatibility_2.12
+ ${flink.version}
+
+
@@ -129,6 +131,12 @@ limitations under the License.
cascading
cascading-hadoop2-io
${cascading.version}
+
+
+ org.apache.hadoop
+ hadoop-core
+
+
@@ -137,6 +145,20 @@ limitations under the License.
${cascading.version}
+
+
+
+ org.apache.hadoop
+ hadoop-common
+ ${hadoop.version}
+
+
+
+ org.apache.hadoop
+ hadoop-mapreduce-client-core
+ ${hadoop.version}
+
+
@@ -164,11 +186,13 @@ limitations under the License.
+
+
org.apache.hadoop
hadoop-hdfs
- 2.2.0
+ ${hadoop.version}
test-jar
test
@@ -186,15 +210,27 @@ limitations under the License.
- org.clapper
- grizzled-slf4j_2.10
- 1.0.2
+ org.apache.logging.log4j
+ log4j-core
+ ${log4j.version}
+
+
+
+ org.apache.logging.log4j
+ log4j-api
+ ${log4j.version}
+
+
+
+ org.apache.logging.log4j
+ log4j-1.2-api
+ ${log4j.version}
- log4j
- log4j
- 1.2.17
+ org.clapper
+ grizzled-slf4j_2.10
+ 1.3.4
@@ -289,6 +325,31 @@ limitations under the License.
+
+ org.apache.maven.plugins
+ maven-shade-plugin
+ 3.4.1
+
+ false
+
+
+
+ package
+
+ shade
+
+
+
+
+
+ com.dataartisans.flink.cascading.example.WordCount
+
+
+
+
+
+
org.apache.maven.plugins
maven-source-plugin
@@ -363,7 +424,8 @@ limitations under the License.
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
- limitations under the License.
+ limitations under the License.mvn exec:exec -Dinput=kinglear.txt -Doutput=wordcounts.txt
+
-->
@@ -406,25 +468,6 @@ limitations under the License.
-
- org.apache.maven.plugins
- maven-checkstyle-plugin
- 2.12.1
-
-
- validate
- validate
-
- check
-
-
-
-
- /tools/maven/checkstyle.xml
- true
-
-
-
org.apache.maven.plugins
@@ -432,8 +475,8 @@ limitations under the License.
3.1
- 1.7
- 1.7
+ 1.8
+ 1.8
diff --git a/src/main/java/com/dataartisans/flink/cascading/example/WordCount.java b/src/main/java/com/dataartisans/flink/cascading/example/WordCount.java
index b13e25d..448e5cb 100644
--- a/src/main/java/com/dataartisans/flink/cascading/example/WordCount.java
+++ b/src/main/java/com/dataartisans/flink/cascading/example/WordCount.java
@@ -44,7 +44,6 @@ public static void main(String[] args) {
RegexSplitGenerator splitter = new RegexSplitGenerator( token, "\\s+" );
// only returns "token"
Pipe docPipe = new Each( "token", text, splitter, Fields.RESULTS );
-
Pipe wcPipe = new Pipe( "wc", docPipe );
wcPipe = new AggregateBy( wcPipe, token, new CountBy(new Fields("count")));
diff --git a/src/main/java/com/dataartisans/flink/cascading/planner/FlinkFlow.java b/src/main/java/com/dataartisans/flink/cascading/planner/FlinkFlow.java
index 025c01c..36933a8 100644
--- a/src/main/java/com/dataartisans/flink/cascading/planner/FlinkFlow.java
+++ b/src/main/java/com/dataartisans/flink/cascading/planner/FlinkFlow.java
@@ -24,6 +24,7 @@
import cascading.flow.planner.PlatformInfo;
import com.dataartisans.flink.cascading.runtime.util.FlinkFlowProcess;
import org.apache.flink.client.program.OptimizerPlanEnvironment;
+import org.apache.flink.client.program.ProgramAbortException;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.mapred.JobConf;
import riffle.process.ProcessComplete;
@@ -81,8 +82,8 @@ public void complete() {
catch(FlowException fe) {
// check if we need to unwrap a ProgramAbortException
Throwable t = fe.getCause();
- if (t instanceof OptimizerPlanEnvironment.ProgramAbortException) {
- throw (OptimizerPlanEnvironment.ProgramAbortException)t;
+ if (t instanceof ProgramAbortException) {
+ throw (ProgramAbortException)t;
}
else {
throw fe;
diff --git a/src/main/java/com/dataartisans/flink/cascading/planner/FlinkFlowStep.java b/src/main/java/com/dataartisans/flink/cascading/planner/FlinkFlowStep.java
index 263e439..1322570 100644
--- a/src/main/java/com/dataartisans/flink/cascading/planner/FlinkFlowStep.java
+++ b/src/main/java/com/dataartisans/flink/cascading/planner/FlinkFlowStep.java
@@ -49,6 +49,8 @@
import com.dataartisans.flink.cascading.runtime.coGroup.regularJoin.CoGroupReducer;
import com.dataartisans.flink.cascading.runtime.coGroup.regularJoin.TupleAppendOuterJoiner;
import com.dataartisans.flink.cascading.runtime.coGroup.regularJoin.TupleOuterJoiner;
+import com.dataartisans.flink.cascading.runtime.each.EachMapper;
+import com.dataartisans.flink.cascading.runtime.groupBy.GroupByReducer;
import com.dataartisans.flink.cascading.runtime.groupBy.GroupByReducer;
import com.dataartisans.flink.cascading.runtime.hashJoin.NaryHashJoinJoiner;
import com.dataartisans.flink.cascading.runtime.util.FlinkFlowProcess;
@@ -57,9 +59,9 @@
import com.dataartisans.flink.cascading.runtime.hashJoin.TupleAppendCrosser;
import com.dataartisans.flink.cascading.runtime.hashJoin.TupleAppendJoiner;
import com.dataartisans.flink.cascading.runtime.hashJoin.HashJoinMapper;
-import com.dataartisans.flink.cascading.runtime.each.EachMapper;
import com.dataartisans.flink.cascading.runtime.sink.TapOutputFormat;
import com.dataartisans.flink.cascading.runtime.source.TapInputFormat;
+import com.dataartisans.flink.cascading.runtime.util.FlinkFlowProcess;
import com.dataartisans.flink.cascading.runtime.util.IdMapper;
import com.dataartisans.flink.cascading.types.tuple.TupleTypeInfo;
import com.dataartisans.flink.cascading.types.tuplearray.TupleArrayTypeInfo;
@@ -71,9 +73,9 @@
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.DataSet;
import org.apache.flink.api.java.ExecutionEnvironment;
+import org.apache.flink.api.java.operators.Operator;
import org.apache.flink.api.java.operators.GroupReduceOperator;
import org.apache.flink.api.java.operators.JoinOperator;
-import org.apache.flink.api.java.operators.Operator;
import org.apache.flink.api.java.operators.PartitionOperator;
import org.apache.flink.api.java.operators.SortPartitionOperator;
import org.apache.flink.api.java.operators.SortedGrouping;
@@ -102,8 +104,8 @@ public class FlinkFlowStep extends BaseFlowStep {
private static final Logger LOG = LoggerFactory.getLogger(FlinkFlowStep.class);
- private ExecutionEnvironment env;
- private List classPath;
+ private final ExecutionEnvironment env;
+ private final List classPath;
public FlinkFlowStep(ExecutionEnvironment env, ElementGraph elementGraph, FlowNodeGraph flowNodeGraph, List classPath) {
super(elementGraph, flowNodeGraph);
diff --git a/src/main/java/com/dataartisans/flink/cascading/planner/FlinkFlowStepJob.java b/src/main/java/com/dataartisans/flink/cascading/planner/FlinkFlowStepJob.java
index 2e02d7f..3ab8a92 100644
--- a/src/main/java/com/dataartisans/flink/cascading/planner/FlinkFlowStepJob.java
+++ b/src/main/java/com/dataartisans/flink/cascading/planner/FlinkFlowStepJob.java
@@ -25,330 +25,285 @@
import org.apache.flink.api.common.ExecutionMode;
import org.apache.flink.api.common.JobID;
import org.apache.flink.api.common.JobSubmissionResult;
-import org.apache.flink.api.common.Plan;
import org.apache.flink.api.java.ExecutionEnvironment;
import org.apache.flink.api.java.LocalEnvironment;
-import org.apache.flink.client.program.Client;
import org.apache.flink.client.program.ContextEnvironment;
-import org.apache.flink.client.program.JobWithJars;
import org.apache.flink.client.program.OptimizerPlanEnvironment;
-import org.apache.flink.configuration.ConfigConstants;
-import org.apache.flink.core.fs.Path;
-import org.apache.flink.optimizer.DataStatistics;
-import org.apache.flink.optimizer.Optimizer;
-import org.apache.flink.optimizer.plan.OptimizedPlan;
-import org.apache.flink.optimizer.plantranslate.JobGraphGenerator;
-import org.apache.flink.runtime.jobgraph.JobGraph;
-import org.apache.flink.runtime.minicluster.FlinkMiniCluster;
-import org.apache.flink.runtime.minicluster.LocalFlinkMiniCluster;
+import org.apache.flink.client.program.ProgramAbortException;
+import org.apache.flink.core.execution.JobClient;
+import org.apache.flink.runtime.minicluster.MiniCluster;
+import org.apache.flink.runtime.minicluster.MiniClusterConfiguration;
+import org.apache.flink.runtime.minicluster.RpcServiceSharing;
import org.apache.hadoop.conf.Configuration;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import scala.concurrent.duration.FiniteDuration;
import java.io.IOException;
-import java.net.MalformedURLException;
-import java.net.URISyntaxException;
-import java.net.URL;
-import java.util.ArrayList;
-import java.util.Collections;
import java.util.List;
-import java.util.concurrent.Callable;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
-import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
-
-
-public class FlinkFlowStepJob extends FlowStepJob
-{
- private static final Logger LOG = LoggerFactory.getLogger( FlinkFlowStepJob.class );
-
- private final Configuration currentConf;
-
- private Client client;
-
- private JobID jobID;
- private Throwable jobException;
-
- private List classPath;
-
- private final ExecutionEnvironment env;
-
- private AccumulatorCache accumulatorCache;
-
- private Future jobSubmission;
-
- private ExecutorService executorService = Executors.newFixedThreadPool(1);
-
- private static final int accumulatorUpdateIntervalSecs = 10;
-
- private volatile static FlinkMiniCluster localCluster;
- private volatile static int localClusterUsers;
- private static final Object lock = new Object();
-
- private static final FiniteDuration DEFAULT_TIMEOUT = new FiniteDuration(60, TimeUnit.SECONDS);
-
-
- public FlinkFlowStepJob( ClientState clientState, FlinkFlowStep flowStep, Configuration currentConf, List classPath ) {
-
- super(clientState, currentConf, flowStep, 1000, 60000, 60000);
-
- this.currentConf = currentConf;
- this.env = ((FlinkFlowStep)this.flowStep).getExecutionEnvironment();
- this.classPath = classPath;
-
- if( flowStep.isDebugEnabled() ) {
- flowStep.logDebug("using polling interval: " + pollingInterval);
- }
- }
-
- @Override
- public Configuration getConfig() {
- return currentConf;
- }
-
- @Override
- protected FlowStepStats createStepStats(ClientState clientState) {
- this.accumulatorCache = new AccumulatorCache(accumulatorUpdateIntervalSecs);
- return new FlinkFlowStepStats(this.flowStep, clientState, accumulatorCache);
- }
-
- protected void internalBlockOnStop() throws IOException {
-
- if (jobSubmission != null && !jobSubmission.isDone()) {
- try {
- client.cancel(jobID);
- } catch (Exception e) {
- throw new IOException("An exception occurred while stopping the Flink job with ID: " + jobID + ": " + e.getMessage());
- }
- }
-
- }
-
- protected void internalNonBlockingStart() throws IOException {
-
- Plan plan = env.createProgramPlan();
-
- // set exchange mode, BATCH is default
- String execMode = getConfig().get(FlinkConfigConstants.EXECUTION_MODE);
- if (execMode == null || FlinkConfigConstants.EXECUTION_MODE_BATCH.equals(execMode)) {
- env.getConfig().setExecutionMode(ExecutionMode.BATCH);
- }
- else if (FlinkConfigConstants.EXECUTION_MODE_PIPELINED.equals(execMode)) {
- env.getConfig().setExecutionMode(ExecutionMode.PIPELINED);
- }
- else {
- LOG.warn("Unknow value for '" + FlinkConfigConstants.EXECUTION_MODE + "' parameter. " +
- "Only '" + FlinkConfigConstants.EXECUTION_MODE_BATCH + "' " +
- "or '" + FlinkConfigConstants.EXECUTION_MODE_PIPELINED + "' supported. " +
- "Using " + FlinkConfigConstants.EXECUTION_MODE_BATCH + " exchange by default.");
- env.getConfig().setExecutionMode(ExecutionMode.BATCH);
- }
-
- Optimizer optimizer = new Optimizer(new DataStatistics(), new org.apache.flink.configuration.Configuration());
- OptimizedPlan optimizedPlan = optimizer.compile(plan);
-
- final JobGraph jobGraph = new JobGraphGenerator().compileJobGraph(optimizedPlan);
- for (String jarPath : classPath) {
- jobGraph.addJar(new Path(jarPath));
- }
-
- jobID = jobGraph.getJobID();
- accumulatorCache.setJobID(jobID);
-
-
- if (isLocalExecution()) {
-
- flowStep.logInfo("Executing in local mode.");
-
- startLocalCluster();
-
- org.apache.flink.configuration.Configuration config = new org.apache.flink.configuration.Configuration();
- config.setString(ConfigConstants.JOB_MANAGER_IPC_ADDRESS_KEY, localCluster.hostname());
-
- client = new Client(config);
- client.setPrintStatusDuringExecution(env.getConfig().isSysoutLoggingEnabled());
-
- } else if (isRemoteExecution()) {
-
- flowStep.logInfo("Executing in cluster mode.");
-
- try {
- String path = this.getClass().getProtectionDomain().getCodeSource().getLocation().toURI().getPath();
- jobGraph.addJar(new Path(path));
- classPath.add(path);
- } catch (URISyntaxException e) {
- throw new IOException("Could not add the submission JAR as a dependency.");
- }
-
- client = ((ContextEnvironment) env).getClient();
- }
-
- List fileList = new ArrayList(classPath.size());
- for (String path : classPath) {
- URL url;
- try {
- url = new URL(path);
- } catch (MalformedURLException e) {
- url = new URL("file://" + path);
- }
- fileList.add(url);
- }
-
- final ClassLoader loader =
- JobWithJars.buildUserCodeClassLoader(fileList, Collections.emptyList(), getClass().getClassLoader());
-
- accumulatorCache.setClient(client);
-
- final Callable callable = new Callable() {
- @Override
- public JobSubmissionResult call() throws Exception {
- return client.runBlocking(jobGraph, loader);
- }
- };
-
- jobSubmission = executorService.submit(callable);
-
- flowStep.logInfo("submitted Flink job: " + jobID);
- }
-
- @Override
- protected void updateNodeStatus( FlowNodeStats flowNodeStats ) {
- try {
- if (internalNonBlockingIsComplete() && internalNonBlockingIsSuccessful()) {
- flowNodeStats.markSuccessful();
- } else if(internalIsStartedRunning()) {
- flowNodeStats.isRunning();
- } else {
- flowNodeStats.markFailed(jobException);
- }
- } catch (IOException e) {
- flowStep.logError("Failed to update node status.");
- }
- }
-
- protected boolean internalNonBlockingIsSuccessful() throws IOException {
- try {
- jobSubmission.get(0, TimeUnit.MILLISECONDS);
- } catch (InterruptedException e) {
- return false;
- } catch (ExecutionException e) {
- jobException = e.getCause();
- return false;
- } catch (TimeoutException e) {
- return false;
- }
-
- boolean isDone = jobSubmission.isDone();
- if (isDone) {
- accumulatorCache.update(true);
- accumulatorCache.setJobID(null);
- accumulatorCache.setClient(null);
- stopCluster();
- }
-
- return isDone;
- }
-
- @Override
- public Throwable call()
- {
- if (env instanceof OptimizerPlanEnvironment) {
- // We have an OptimizerPlanEnvironment.
- // This environment is only used to to fetch the Flink execution plan.
- try {
- // OptimizerPlanEnvironment does not execute but only build the execution plan.
- env.execute("plan generation");
- }
- // execute() throws a ProgramAbortException if everything goes well
- catch(OptimizerPlanEnvironment.ProgramAbortException pae) {
- // Forward call() to get Cascading's internal job stats right.
- // The job will be skipped due to the overridden isSkipFlowStep method.
- super.call();
- // forward expected ProgramAbortException
- return pae;
- }
- //
- catch(Exception e) {
- // forward unexpected exception
- return e;
- }
- }
- // forward to call() if we have a regular ExecutionEnvironment
- return super.call();
-
- }
-
- protected boolean isSkipFlowStep() throws IOException
- {
- if (env instanceof OptimizerPlanEnvironment) {
- // We have an OptimizerPlanEnvironment.
- // This environment is only used to to fetch the Flink execution plan.
- // We do not want to execute the job in this case.
- return true;
- } else {
- return super.isSkipFlowStep();
- }
- }
-
- @Override
- protected boolean isRemoteExecution() {
- return env instanceof ContextEnvironment;
- }
-
- @Override
- protected Throwable getThrowable() {
- return jobException;
- }
-
- protected String internalJobId() {
- return jobID.toString();
- }
-
- protected boolean internalNonBlockingIsComplete() throws IOException {
- return jobSubmission.isDone();
- }
-
- protected void dumpDebugInfo() {
- }
-
- protected boolean internalIsStartedRunning() {
- return jobSubmission != null;
- }
-
- private boolean isLocalExecution() {
- return env instanceof LocalEnvironment;
- }
-
- private void startLocalCluster() {
- synchronized (lock) {
- if (localCluster == null) {
- org.apache.flink.configuration.Configuration configuration = new org.apache.flink.configuration.Configuration();
- configuration.setInteger(ConfigConstants.TASK_MANAGER_NUM_TASK_SLOTS, env.getParallelism() * 2);
- localCluster = new LocalFlinkMiniCluster(configuration, false);
- localCluster.start();
- }
- localClusterUsers++;
- }
- }
-
- private void stopCluster() {
- synchronized (lock) {
- if (localCluster != null) {
- if (--localClusterUsers <= 0) {
- localCluster.shutdown();
- localCluster.awaitTermination();
- localCluster = null;
- localClusterUsers = 0;
- }
- }
- if (executorService != null) {
- executorService.shutdown();
- }
- }
- }
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+
+
+public class FlinkFlowStepJob extends FlowStepJob {
+ private static final Logger LOG = LoggerFactory.getLogger(FlinkFlowStepJob.class);
+ private static final int accumulatorUpdateIntervalSecs = 10;
+ private static final Object lock = new Object();
+ private static final FiniteDuration DEFAULT_TIMEOUT = new FiniteDuration(60, TimeUnit.SECONDS);
+ private volatile static MiniCluster localCluster;
+ private volatile static int localClusterUsers;
+ private final Configuration currentConf;
+ private final List classPath;
+ private final ExecutionEnvironment env;
+ private final ExecutorService executorService = Executors.newFixedThreadPool(1);
+ private JobClient client;
+ private JobID jobID;
+ private Throwable jobException;
+ private AccumulatorCache accumulatorCache;
+ private Future jobSubmission;
+
+
+ public FlinkFlowStepJob(ClientState clientState, FlinkFlowStep flowStep, Configuration currentConf, List classPath) {
+
+ super(clientState, currentConf, flowStep, 1000, 60000, 60000);
+
+ this.currentConf = currentConf;
+ this.env = ((FlinkFlowStep) this.flowStep).getExecutionEnvironment();
+ this.classPath = classPath.stream().map(jarpath -> "file://" + jarpath).collect(Collectors.toList());
+
+ if (flowStep.isDebugEnabled()) {
+ flowStep.logDebug("using polling interval: " + pollingInterval);
+ }
+ }
+
+ @Override
+ public Configuration getConfig() {
+ return currentConf;
+ }
+
+ @Override
+ protected FlowStepStats createStepStats(ClientState clientState) {
+ this.accumulatorCache = new AccumulatorCache(accumulatorUpdateIntervalSecs);
+ return new FlinkFlowStepStats(this.flowStep, clientState, accumulatorCache);
+ }
+
+ protected void internalBlockOnStop() throws IOException {
+ if (client != null) {
+ try {
+ client.cancel().get();
+ } catch (Exception e) {
+ throw new IOException("An exception occurred while stopping the Flink job with ID: " + jobID + ": " + e.getMessage());
+ }
+ }
+
+ }
+
+ protected void internalNonBlockingStart() throws IOException {
+
+ // set exchange mode, BATCH is default
+ String execMode = getConfig().get(FlinkConfigConstants.EXECUTION_MODE);
+ if (execMode == null || FlinkConfigConstants.EXECUTION_MODE_BATCH.equals(execMode)) {
+ env.getConfig().setExecutionMode(ExecutionMode.BATCH);
+ } else if (FlinkConfigConstants.EXECUTION_MODE_PIPELINED.equals(execMode)) {
+ env.getConfig().setExecutionMode(ExecutionMode.PIPELINED);
+ } else {
+ LOG.warn("Unknow value for '" + FlinkConfigConstants.EXECUTION_MODE + "' parameter. " +
+ "Only '" + FlinkConfigConstants.EXECUTION_MODE_BATCH + "' " +
+ "or '" + FlinkConfigConstants.EXECUTION_MODE_PIPELINED + "' supported. " +
+ "Using " + FlinkConfigConstants.EXECUTION_MODE_BATCH + " exchange by default.");
+ env.getConfig().setExecutionMode(ExecutionMode.BATCH);
+ }
+
+ env.getConfiguration().setString(FlinkConfigConstants.PIPELINE_CLASSPATHS, String.join(",", classPath));
+
+ if (!env.getConfiguration().containsKey(FlinkConfigConstants.LEAK_CLASSLOADER_CHECK)) {
+ LOG.warn("disable classloader leak check which failed PlatformTests");
+ env.getConfiguration().setBoolean(FlinkConfigConstants.LEAK_CLASSLOADER_CHECK, false);
+ }
+
+ if (!env.getConfiguration().containsKey(FlinkConfigConstants.NETWORK_MEMORY_MIN)) {
+ env.getConfiguration().setString(FlinkConfigConstants.NETWORK_MEMORY_MIN, "128mb");
+ }
+
+ if (!env.getConfiguration().containsKey(FlinkConfigConstants.TASKMANAGER_MEMORY_MANAGED_SIZE)) {
+ env.getConfiguration().setString(FlinkConfigConstants.TASKMANAGER_MEMORY_MANAGED_SIZE, "512mb");
+ }
+
+ if (isLocalExecution()) {
+
+ flowStep.logInfo("Executing in local mode.");
+
+ try {
+ startLocalCluster();
+ } catch (Exception e) {
+ flowStep.logError("Fail to start local cluster.");
+ throw new RuntimeException(e);
+ }
+ } else if (isRemoteExecution()) {
+ flowStep.logInfo("Executing in cluster mode.");
+ }
+
+ try {
+ client = env.executeAsync();
+ jobID = client.getJobID();
+ accumulatorCache.setClient(client);
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+
+ flowStep.logInfo("submitted Flink job: " + jobID);
+ }
+
+ @Override
+ protected void updateNodeStatus(FlowNodeStats flowNodeStats) {
+ try {
+ if (internalNonBlockingIsComplete() && internalNonBlockingIsSuccessful()) {
+ flowNodeStats.markSuccessful();
+ } else if (internalIsStartedRunning()) {
+ flowNodeStats.isRunning();
+ } else {
+ flowNodeStats.markFailed(jobException);
+ }
+ } catch (IOException e) {
+ flowStep.logError("Failed to update node status.");
+ }
+ }
+
+ protected boolean internalNonBlockingIsSuccessful() throws IOException {
+ try {
+ client.getJobExecutionResult().get(100, TimeUnit.MILLISECONDS);
+ } catch (InterruptedException e) {
+ return false;
+ } catch (ExecutionException e) {
+ jobException = e.getCause();
+ return false;
+ } catch (TimeoutException e) {
+ return false;
+ }
+
+ boolean isDone = client.getJobExecutionResult().isDone();
+ if (isDone) {
+ accumulatorCache.update(true);
+ accumulatorCache.setClient(null);
+ try {
+ stopCluster();
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ return isDone;
+ }
+
+ @Override
+ public Throwable call() {
+ if (env instanceof OptimizerPlanEnvironment) {
+ // We have an OptimizerPlanEnvironment.
+ // This environment is only used to to fetch the Flink execution plan.
+ try {
+ // OptimizerPlanEnvironment does not execute but only build the execution plan.
+ env.execute("plan generation");
+ }
+ // execute() throws a ProgramAbortException if everything goes well
+ catch (ProgramAbortException pae) {
+ // Forward call() to get Cascading's internal job stats right.
+ // The job will be skipped due to the overridden isSkipFlowStep method.
+ super.call();
+ // forward expected ProgramAbortException
+ return pae;
+ }
+ //
+ catch (Exception e) {
+ // forward unexpected exception
+ return e;
+ }
+ }
+ // forward to call() if we have a regular ExecutionEnvironment
+ return super.call();
+
+ }
+
+ protected boolean isSkipFlowStep() throws IOException {
+ if (env instanceof OptimizerPlanEnvironment) {
+ // We have an OptimizerPlanEnvironment.
+ // This environment is only used to to fetch the Flink execution plan.
+ // We do not want to execute the job in this case.
+ return true;
+ } else {
+ return super.isSkipFlowStep();
+ }
+ }
+
+ @Override
+ protected boolean isRemoteExecution() {
+ return env instanceof ContextEnvironment;
+ }
+
+ @Override
+ protected Throwable getThrowable() {
+ return jobException;
+ }
+
+ protected String internalJobId() {
+ return jobID.toString();
+ }
+
+ protected boolean internalNonBlockingIsComplete() throws IOException {
+ try {
+ return client.getJobExecutionResult().isDone() || client.getJobExecutionResult().isCompletedExceptionally();
+ } catch (Exception e) {
+ return false;
+ }
+ }
+
+ protected void dumpDebugInfo() {
+ }
+
+ protected boolean internalIsStartedRunning() {
+ try {
+ return client.getJobExecutionResult() != null;
+ } catch (Exception e) {
+ return false;
+ }
+ }
+
+ private boolean isLocalExecution() {
+ return env instanceof LocalEnvironment;
+ }
+
+ private void startLocalCluster() throws Exception {
+ synchronized (lock) {
+ if (localCluster == null) {
+ final MiniClusterConfiguration miniClusterConfiguration = new MiniClusterConfiguration.Builder()
+ .setNumSlotsPerTaskManager(env.getParallelism() * 2)
+ .setRpcServiceSharing(RpcServiceSharing.SHARED)
+ .withRandomPorts()
+ .build();
+ localCluster = new MiniCluster(miniClusterConfiguration);
+ localCluster.start();
+ }
+ localClusterUsers++;
+ }
+ }
+
+ private void stopCluster() throws Exception {
+ synchronized (lock) {
+ if (localCluster != null) {
+ if (--localClusterUsers <= 0) {
+ localCluster.close();
+ localCluster = null;
+ localClusterUsers = 0;
+ }
+ }
+ if (executorService != null) {
+ executorService.shutdown();
+ }
+ }
+ }
}
diff --git a/src/main/java/com/dataartisans/flink/cascading/planner/FlinkFlowStepStats.java b/src/main/java/com/dataartisans/flink/cascading/planner/FlinkFlowStepStats.java
index eaec507..b0c8566 100644
--- a/src/main/java/com/dataartisans/flink/cascading/planner/FlinkFlowStepStats.java
+++ b/src/main/java/com/dataartisans/flink/cascading/planner/FlinkFlowStepStats.java
@@ -19,8 +19,8 @@
import cascading.flow.FlowStep;
import cascading.management.state.ClientState;
import cascading.stats.FlowStepStats;
-import com.dataartisans.flink.cascading.runtime.stats.EnumStringConverter;
import com.dataartisans.flink.cascading.runtime.stats.AccumulatorCache;
+import com.dataartisans.flink.cascading.runtime.stats.EnumStringConverter;
import java.util.Collection;
import java.util.HashSet;
@@ -29,82 +29,82 @@
public class FlinkFlowStepStats extends FlowStepStats {
- private AccumulatorCache accumulatorCache;
-
- protected FlinkFlowStepStats(FlowStep flowStep, ClientState clientState, AccumulatorCache accumulatorCache) {
- super(flowStep, clientState);
- this.accumulatorCache = accumulatorCache;
- }
-
- @Override
- public void recordChildStats() {
- // TODO
- }
-
- @Override
- public String getProcessStepID() {
- return null;
- }
-
- @Override
- public Collection getCounterGroupsMatching(String regex) {
- return null;
- }
-
- @Override
- public void captureDetail(Type depth) {
- // TODO
- }
-
- @Override
- public long getLastSuccessfulCounterFetchTime() {
- return accumulatorCache.getLastUpdateTime();
- }
-
- @Override
- public Collection getCounterGroups() {
- accumulatorCache.update();
- Map currentAccumulators = accumulatorCache.getCurrentAccumulators();
- Set result = new HashSet();
-
- for (String key : currentAccumulators.keySet()) {
- result.add(EnumStringConverter.groupCounterToGroup(key));
- }
- return result;
- }
-
- @Override
- public Collection getCountersFor(String group) {
- accumulatorCache.update();
- Map currentAccumulators = accumulatorCache.getCurrentAccumulators();
- Set result = new HashSet();
-
- for (String key : currentAccumulators.keySet()) {
- if (EnumStringConverter.accInGroup(group, key)) {
- result.add(EnumStringConverter.groupCounterToCounter(key));
- }
- }
- return result;
- }
-
- @Override
- public long getCounterValue(Enum counter) {
- return getCounterValue(EnumStringConverter.enumToGroup(counter), EnumStringConverter.enumToCounter(counter));
- }
-
- @Override
- public long getCounterValue(String group, String counter) {
- accumulatorCache.update();
- Map currentAccumulators = accumulatorCache.getCurrentAccumulators();
-
- for (String key : currentAccumulators.keySet()) {
- if (EnumStringConverter.accMatchesGroupCounter(key, group, counter)) {
- Object o = currentAccumulators.get(key);
- return (Long) o;
- }
- }
- // Cascading returns 0 in case of empty accumulators
- return 0;
- }
+ private final AccumulatorCache accumulatorCache;
+
+ protected FlinkFlowStepStats(FlowStep flowStep, ClientState clientState, AccumulatorCache accumulatorCache) {
+ super(flowStep, clientState);
+ this.accumulatorCache = accumulatorCache;
+ }
+
+ @Override
+ public void recordChildStats() {
+ // TODO
+ }
+
+ @Override
+ public String getProcessStepID() {
+ return null;
+ }
+
+ @Override
+ public Collection getCounterGroupsMatching(String regex) {
+ return null;
+ }
+
+ @Override
+ public void captureDetail(Type depth) {
+ // TODO
+ }
+
+ @Override
+ public long getLastSuccessfulCounterFetchTime() {
+ return accumulatorCache.getLastUpdateTime();
+ }
+
+ @Override
+ public Collection getCounterGroups() {
+ accumulatorCache.update();
+ Map currentAccumulators = accumulatorCache.getCurrentAccumulators();
+ Set result = new HashSet();
+
+ for (String key : currentAccumulators.keySet()) {
+ result.add(EnumStringConverter.groupCounterToGroup(key));
+ }
+ return result;
+ }
+
+ @Override
+ public Collection getCountersFor(String group) {
+ accumulatorCache.update();
+ Map currentAccumulators = accumulatorCache.getCurrentAccumulators();
+ Set result = new HashSet();
+
+ for (String key : currentAccumulators.keySet()) {
+ if (EnumStringConverter.accInGroup(group, key)) {
+ result.add(EnumStringConverter.groupCounterToCounter(key));
+ }
+ }
+ return result;
+ }
+
+ @Override
+ public long getCounterValue(Enum counter) {
+ return getCounterValue(EnumStringConverter.enumToGroup(counter), EnumStringConverter.enumToCounter(counter));
+ }
+
+ @Override
+ public long getCounterValue(String group, String counter) {
+ accumulatorCache.update();
+ Map currentAccumulators = accumulatorCache.getCurrentAccumulators();
+
+ for (String key : currentAccumulators.keySet()) {
+ if (EnumStringConverter.accMatchesGroupCounter(key, group, counter)) {
+ Object o = currentAccumulators.get(key);
+ return (Long) o;
+ }
+ }
+ // Cascading returns 0 in case of empty accumulators
+ return 0;
+ }
}
diff --git a/src/main/java/com/dataartisans/flink/cascading/planner/FlinkPlanner.java b/src/main/java/com/dataartisans/flink/cascading/planner/FlinkPlanner.java
index 885dbe5..f6f4b3c 100644
--- a/src/main/java/com/dataartisans/flink/cascading/planner/FlinkPlanner.java
+++ b/src/main/java/com/dataartisans/flink/cascading/planner/FlinkPlanner.java
@@ -30,7 +30,7 @@
import cascading.tap.Tap;
import com.dataartisans.flink.cascading.util.Version;
import org.apache.flink.api.java.ExecutionEnvironment;
-import org.apache.flink.client.CliFrontend;
+import org.apache.flink.client.cli.CliFrontend;
import org.apache.flink.configuration.ConfigConstants;
import org.apache.flink.configuration.GlobalConfiguration;
import org.apache.hadoop.conf.Configuration;
@@ -45,19 +45,18 @@ public class FlinkPlanner extends FlowPlanner {
private Configuration defaultConfig;
- private List classPath;
+ private final List classPath;
- private ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
+ private final ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
public FlinkPlanner(List classPath) {
super();
this.classPath = classPath;
- env.getConfig().disableSysoutLogging();
if (env.getParallelism() <= 0) {
// load the default parallelism from config
GlobalConfiguration.loadConfiguration(new File(CliFrontend.getConfigurationDirectoryFromEnv()).getAbsolutePath());
- org.apache.flink.configuration.Configuration configuration = GlobalConfiguration.getConfiguration();
+ org.apache.flink.configuration.Configuration configuration = GlobalConfiguration.loadConfiguration();
int parallelism = configuration.getInteger(ConfigConstants.DEFAULT_PARALLELISM_KEY, -1);
if (parallelism <= 0) {
throw new RuntimeException("Please set the default parallelism via the -p command-line flag");
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/boundaryStages/BoundaryInStage.java b/src/main/java/com/dataartisans/flink/cascading/runtime/boundaryStages/BoundaryInStage.java
index 17eb818..d79569a 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/boundaryStages/BoundaryInStage.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/boundaryStages/BoundaryInStage.java
@@ -64,8 +64,8 @@ public void run(Object input) throws Throwable {
{
tuple = tupleIterator.next();
tupleEntry.setTuple(tuple);
- flowProcess.increment( StepCounters.Tuples_Read, 1 );
- flowProcess.increment( SliceCounters.Tuples_Read, 1 );
+ //flowProcess.increment( StepCounters.Tuples_Read, 1 );
+ //flowProcess.increment( SliceCounters.Tuples_Read, 1 );
}
catch( OutOfMemoryError error ) {
handleReThrowableException( "out of memory, try increasing task memory allocation", error );
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/bufferJoin/BufferJoinKeyExtractor.java b/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/bufferJoin/BufferJoinKeyExtractor.java
index 44daab9..5c21ac4 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/bufferJoin/BufferJoinKeyExtractor.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/bufferJoin/BufferJoinKeyExtractor.java
@@ -22,9 +22,9 @@
public class BufferJoinKeyExtractor implements MapFunction> {
- private int[] keyPos;
- private Tuple3 outT;
- private Tuple defaultKey = new Tuple(1);
+ private final int[] keyPos;
+ private final Tuple3 outT;
+ private final Tuple defaultKey = new Tuple(1);
public BufferJoinKeyExtractor(int inputId, int[] keyPos) {
this.keyPos = keyPos;
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/bufferJoin/CoGroupBufferClosure.java b/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/bufferJoin/CoGroupBufferClosure.java
index 14942ef..f02eccf 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/bufferJoin/CoGroupBufferClosure.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/bufferJoin/CoGroupBufferClosure.java
@@ -143,7 +143,7 @@ private TupleBuilder makeJoinedBuilder( final Fields[] joinFields )
return new TupleBuilder()
{
- Tuple result = TupleViews.createComposite(fields);
+ final Tuple result = TupleViews.createComposite(fields);
@Override
public Tuple makeResult( Tuple[] tuples )
@@ -158,7 +158,7 @@ private Iterator makeIterator( final int pos, final Iterator values )
return new Iterator()
{
final int cleanPos = valueFields.length == 1 ? 0 : pos; // support repeated pipes
- cascading.tuple.util.TupleBuilder[] valueBuilder = new cascading.tuple.util.TupleBuilder[ valueFields.length ];
+ final cascading.tuple.util.TupleBuilder[] valueBuilder = new cascading.tuple.util.TupleBuilder[ valueFields.length ];
{
for( int i = 0; i < valueFields.length; i++ ) {
@@ -184,7 +184,7 @@ public Tuple makeResult( Tuple valueTuple, Tuple groupTuple )
else {
return new cascading.tuple.util.TupleBuilder() {
- Tuple result = TupleViews.createOverride(valueField, joinField);
+ final Tuple result = TupleViews.createOverride(valueField, joinField);
@Override
public Tuple makeResult(Tuple valueTuple, Tuple groupTuple) {
@@ -291,13 +291,13 @@ public void remove() {
};
}
- static interface TupleBuilder {
+ interface TupleBuilder {
Tuple makeResult(Tuple[] tuples);
}
private static class FlinkUnwrappingIterator implements Iterator {
- private Iterator> flinkIterator;
+ private final Iterator> flinkIterator;
public FlinkUnwrappingIterator(Iterable> vals) {
@@ -359,7 +359,7 @@ public Iterator iterator() {
try {
if (iterator == null) {
// use emptyList() iterator for java 6 compatibility
- return Collections.emptyList().iterator();
+ return Collections.emptyIterator();
}
return iterator;
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/regularJoin/CoGroupInGate.java b/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/regularJoin/CoGroupInGate.java
index a7900f8..8faf99e 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/regularJoin/CoGroupInGate.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/regularJoin/CoGroupInGate.java
@@ -133,8 +133,8 @@ private static class JoinResultIterator implements Iterator {
private Iterator> input;
- private JoinClosure closure;
- private Joiner joiner;
+ private final JoinClosure closure;
+ private final Joiner joiner;
private Iterator joinedTuples;
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/regularJoin/TupleAppendOuterJoiner.java b/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/regularJoin/TupleAppendOuterJoiner.java
index e15028a..2d5871a 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/regularJoin/TupleAppendOuterJoiner.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/regularJoin/TupleAppendOuterJoiner.java
@@ -24,10 +24,10 @@
public class TupleAppendOuterJoiner extends RichJoinFunction, Tuple, Tuple2> {
- private int tupleListPos;
- private int tupleListSize;
- private Fields inputFields;
- private Fields keyFields;
+ private final int tupleListPos;
+ private final int tupleListSize;
+ private final Fields inputFields;
+ private final Fields keyFields;
private transient Tuple2 outT;
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/regularJoin/TupleOuterJoiner.java b/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/regularJoin/TupleOuterJoiner.java
index e9291fd..fa9ed72 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/regularJoin/TupleOuterJoiner.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/coGroup/regularJoin/TupleOuterJoiner.java
@@ -24,11 +24,11 @@
public class TupleOuterJoiner extends RichJoinFunction> {
- private int tupleListSize;
- private Fields inputFieldsLeft;
- private Fields keyFieldsLeft;
- private Fields inputFieldsRight;
- private Fields keyFieldsRight;
+ private final int tupleListSize;
+ private final Fields inputFieldsLeft;
+ private final Fields keyFieldsLeft;
+ private final Fields inputFieldsRight;
+ private final Fields keyFieldsRight;
private transient Tuple2 outT;
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/JoinBoundaryInStage.java b/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/JoinBoundaryInStage.java
index 1a71a99..9f62427 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/JoinBoundaryInStage.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/JoinBoundaryInStage.java
@@ -62,8 +62,8 @@ public void run(Object input) throws Throwable {
throw new RuntimeException("JoinBoundaryInStage expects Tuple2", cce);
}
- flowProcess.increment( StepCounters.Tuples_Read, 1 );
- flowProcess.increment( SliceCounters.Tuples_Read, 1 );
+ //flowProcess.increment( StepCounters.Tuples_Read, 1 );
+ //flowProcess.increment( SliceCounters.Tuples_Read, 1 );
next.receive(this, joinInputTuples);
}
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/JoinBoundaryMapperInStage.java b/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/JoinBoundaryMapperInStage.java
index 0e8324d..8dd9a18 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/JoinBoundaryMapperInStage.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/JoinBoundaryMapperInStage.java
@@ -61,8 +61,8 @@ public void run(Object input) throws Throwable {
try {
joinListTuple = joinInputIterator.next();
- flowProcess.increment( StepCounters.Tuples_Read, 1 );
- flowProcess.increment( SliceCounters.Tuples_Read, 1 );
+ //flowProcess.increment( StepCounters.Tuples_Read, 1 );
+ //flowProcess.increment( SliceCounters.Tuples_Read, 1 );
}
catch( CascadingException exception ) {
handleException( exception, null );
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/JoinClosure.java b/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/JoinClosure.java
index 27a3a47..a4bfd89 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/JoinClosure.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/JoinClosure.java
@@ -99,7 +99,7 @@ private TupleBuilder makeJoinedBuilder( final Fields[] joinFields )
return new TupleBuilder()
{
- Tuple result = TupleViews.createComposite(fields);
+ final Tuple result = TupleViews.createComposite(fields);
@Override
public Tuple makeResult( Tuple[] tuples )
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/TupleAppendCrosser.java b/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/TupleAppendCrosser.java
index d83f9a9..d32e81e 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/TupleAppendCrosser.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/TupleAppendCrosser.java
@@ -22,7 +22,7 @@
public class TupleAppendCrosser implements CrossFunction, Tuple, Tuple2> {
- private int tupleListPos;
+ private final int tupleListPos;
public TupleAppendCrosser(int tupleListPos) {
this.tupleListPos = tupleListPos;
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/TupleAppendJoiner.java b/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/TupleAppendJoiner.java
index 1219cb5..9884e90 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/TupleAppendJoiner.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/hashJoin/TupleAppendJoiner.java
@@ -22,7 +22,7 @@
public class TupleAppendJoiner implements JoinFunction, Tuple, Tuple2> {
- private int tupleListPos;
+ private final int tupleListPos;
public TupleAppendJoiner(int tupleListPos) {
this.tupleListPos = tupleListPos;
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/sink/SinkBoundaryInStage.java b/src/main/java/com/dataartisans/flink/cascading/runtime/sink/SinkBoundaryInStage.java
index 830c90d..5fc92d7 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/sink/SinkBoundaryInStage.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/sink/SinkBoundaryInStage.java
@@ -34,7 +34,7 @@
public class SinkBoundaryInStage extends ElementStage implements InputSource {
private boolean nextStarted;
- private TupleEntry tupleEntry;
+ private final TupleEntry tupleEntry;
public SinkBoundaryInStage(FlowProcess flowProcess, FlowElement flowElement, FlowNode node) {
super(flowProcess, flowElement);
@@ -69,8 +69,8 @@ public void run(Object input) throws Throwable {
try {
Tuple tuple = (Tuple)input;
tupleEntry.setTuple(tuple);
- flowProcess.increment( StepCounters.Tuples_Read, 1 );
- flowProcess.increment(SliceCounters.Tuples_Read, 1);
+ //flowProcess.increment( StepCounters.Tuples_Read, 1 );
+ //flowProcess.increment(SliceCounters.Tuples_Read, 1);
}
catch( OutOfMemoryError error ) {
handleReThrowableException("out of memory, try increasing task memory allocation", error);
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/sink/TapOutputFormat.java b/src/main/java/com/dataartisans/flink/cascading/runtime/sink/TapOutputFormat.java
index 7fcb154..82783bc 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/sink/TapOutputFormat.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/sink/TapOutputFormat.java
@@ -47,7 +47,7 @@ public class TapOutputFormat extends RichOutputFormat implements Finalize
private static final Logger LOG = LoggerFactory.getLogger(TapOutputFormat.class);
- private FlowNode flowNode;
+ private final FlowNode flowNode;
private transient org.apache.hadoop.conf.Configuration config;
private transient FlinkFlowProcess flowProcess;
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/source/SourceStreamGraph.java b/src/main/java/com/dataartisans/flink/cascading/runtime/source/SourceStreamGraph.java
index 9859ef6..a3cb8af 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/source/SourceStreamGraph.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/source/SourceStreamGraph.java
@@ -29,7 +29,7 @@
public class SourceStreamGraph extends NodeStreamGraph {
- private TapSourceStage sourceStage;
+ private final TapSourceStage sourceStage;
private SingleOutBoundaryStage sinkStage;
public SourceStreamGraph(FlowProcess flowProcess, FlowNode node, Tap tap) {
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/source/TapInputFormat.java b/src/main/java/com/dataartisans/flink/cascading/runtime/source/TapInputFormat.java
index 15e0738..71d28e2 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/source/TapInputFormat.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/source/TapInputFormat.java
@@ -64,7 +64,7 @@ public class TapInputFormat extends RichInputFormat {
private static final Logger LOG = LoggerFactory.getLogger(TapInputFormat.class);
- private FlowNode flowNode;
+ private final FlowNode flowNode;
private transient SourceStreamGraph streamGraph;
private transient TapSourceStage sourceStage;
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/source/TapSourceStage.java b/src/main/java/com/dataartisans/flink/cascading/runtime/source/TapSourceStage.java
index 558a0a2..88cc4ed 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/source/TapSourceStage.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/source/TapSourceStage.java
@@ -33,7 +33,7 @@ public class TapSourceStage extends SourceStage {
private static final Logger LOG = LoggerFactory.getLogger(TapSourceStage.class);
- private Tap source;
+ private final Tap source;
private TupleEntryIterator iterator;
public TapSourceStage(FlowProcess flowProcess, Tap tap) {
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/stats/AccumulatorCache.java b/src/main/java/com/dataartisans/flink/cascading/runtime/stats/AccumulatorCache.java
index 22a24f6..942ae9a 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/stats/AccumulatorCache.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/stats/AccumulatorCache.java
@@ -18,7 +18,8 @@
import org.apache.flink.api.common.JobID;
-import org.apache.flink.client.program.Client;
+import org.apache.flink.client.program.MiniClusterClient;
+import org.apache.flink.core.execution.JobClient;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -31,7 +32,7 @@ public class AccumulatorCache {
private JobID jobID;
- private Client client;
+ private JobClient client;
private volatile Map currentAccumulators = Collections.emptyMap();
@@ -39,7 +40,7 @@ public class AccumulatorCache {
private long lastUpdateTime;
public AccumulatorCache(int updateIntervalSecs) {
- this.updateIntervalMillis = updateIntervalSecs * 1000;
+ this.updateIntervalMillis = updateIntervalSecs * 1000L;
}
public void update() {
@@ -53,35 +54,25 @@ public void update(boolean force) {
return;
}
- if (jobID == null) {
- return;
- }
-
if (client != null) {
try {
- currentAccumulators = client.getAccumulators(jobID);
+ currentAccumulators = client.getAccumulators().get();
lastUpdateTime = currentTime;
LOG.debug("Updated accumulators: {}", currentAccumulators);
} catch (Exception e) {
LOG.error("Failed to fetch accumulators for job {}.", jobID);
}
-
}
-
}
public Map getCurrentAccumulators() {
return currentAccumulators;
}
-
- public void setJobID(JobID jobID) {
- this.jobID = jobID;
- }
-
- public void setClient(Client client) {
+ public void setClient(JobClient client) {
this.client = client;
+ if(client != null) this.jobID = client.getJobID();
}
public long getLastUpdateTime() {
diff --git a/src/main/java/com/dataartisans/flink/cascading/runtime/util/FlinkFlowProcess.java b/src/main/java/com/dataartisans/flink/cascading/runtime/util/FlinkFlowProcess.java
index 68868f0..2535fe8 100644
--- a/src/main/java/com/dataartisans/flink/cascading/runtime/util/FlinkFlowProcess.java
+++ b/src/main/java/com/dataartisans/flink/cascading/runtime/util/FlinkFlowProcess.java
@@ -42,7 +42,7 @@
public class FlinkFlowProcess extends FlowProcess {
private transient RuntimeContext runtimeContext;
- private Configuration conf;
+ private final Configuration conf;
private String taskId;
public FlinkFlowProcess() {
diff --git a/src/main/java/com/dataartisans/flink/cascading/types/field/CustomFieldComparator.java b/src/main/java/com/dataartisans/flink/cascading/types/field/CustomFieldComparator.java
index 7819d3e..1b9af35 100644
--- a/src/main/java/com/dataartisans/flink/cascading/types/field/CustomFieldComparator.java
+++ b/src/main/java/com/dataartisans/flink/cascading/types/field/CustomFieldComparator.java
@@ -37,7 +37,7 @@ public class CustomFieldComparator extends TypeComparator {
private final Hasher hasher;
- private TypeSerializer serializer;
+ private final TypeSerializer serializer;
private transient Comparable ref;
@@ -75,12 +75,7 @@ public int hash(Comparable record) {
}
public void setReference(Comparable toCompare) {
- if(toCompare == null) {
- this.ref = null;
- }
- else {
- this.ref = toCompare;
- }
+ this.ref = toCompare;
}
public boolean equalToReference(Comparable candidate) {
diff --git a/src/main/java/com/dataartisans/flink/cascading/types/field/FieldComparator.java b/src/main/java/com/dataartisans/flink/cascading/types/field/FieldComparator.java
index eca2d68..c8fff94 100644
--- a/src/main/java/com/dataartisans/flink/cascading/types/field/FieldComparator.java
+++ b/src/main/java/com/dataartisans/flink/cascading/types/field/FieldComparator.java
@@ -30,7 +30,7 @@ public class FieldComparator> extends TypeComparator
private final boolean ascending;
private final Class type;
- private TypeSerializer serializer;
+ private final TypeSerializer serializer;
private transient T ref;
@@ -59,12 +59,7 @@ public boolean equalToReference(T t) {
if(t != null && ref != null) {
return t.equals(this.ref);
}
- else if(t == null && ref == null) {
- return true;
- }
- else {
- return false;
- }
+ else return t == null && ref == null;
}
@Override
diff --git a/src/main/java/com/dataartisans/flink/cascading/types/field/WrappingFieldComparator.java b/src/main/java/com/dataartisans/flink/cascading/types/field/WrappingFieldComparator.java
index 459dabc..1c97e6d 100644
--- a/src/main/java/com/dataartisans/flink/cascading/types/field/WrappingFieldComparator.java
+++ b/src/main/java/com/dataartisans/flink/cascading/types/field/WrappingFieldComparator.java
@@ -67,12 +67,7 @@ public boolean equalToReference(T t) {
if(t != null && !this.refNull) {
return this.wrappedComparator.equalToReference(t);
}
- else if(t == null && this.refNull) {
- return true;
- }
- else {
- return false;
- }
+ else return t == null && this.refNull;
}
@Override
diff --git a/src/main/java/com/dataartisans/flink/cascading/types/tuple/DefinedTupleSerializer.java b/src/main/java/com/dataartisans/flink/cascading/types/tuple/DefinedTupleSerializer.java
index 7c32b64..a3e5219 100644
--- a/src/main/java/com/dataartisans/flink/cascading/types/tuple/DefinedTupleSerializer.java
+++ b/src/main/java/com/dataartisans/flink/cascading/types/tuple/DefinedTupleSerializer.java
@@ -20,6 +20,7 @@
import cascading.tuple.Fields;
import cascading.tuple.Tuple;
import org.apache.flink.api.common.typeutils.TypeSerializer;
+import org.apache.flink.api.common.typeutils.TypeSerializerSnapshot;
import org.apache.flink.core.memory.DataInputView;
import org.apache.flink.core.memory.DataOutputView;
@@ -227,6 +228,10 @@ public int hashCode() {
}
@Override
+ public TypeSerializerSnapshot snapshotConfiguration() {
+ return null;
+ }
+
public boolean canEqual(Object obj) {
return obj instanceof DefinedTupleSerializer;
}
diff --git a/src/main/java/com/dataartisans/flink/cascading/types/tuple/TupleTypeInfo.java b/src/main/java/com/dataartisans/flink/cascading/types/tuple/TupleTypeInfo.java
index 4531446..e065429 100644
--- a/src/main/java/com/dataartisans/flink/cascading/types/tuple/TupleTypeInfo.java
+++ b/src/main/java/com/dataartisans/flink/cascading/types/tuple/TupleTypeInfo.java
@@ -35,11 +35,11 @@ public class TupleTypeInfo extends CompositeType {
private final static int NEG_FIELD_POS_OFFSET = Integer.MAX_VALUE / 2;
- private Fields schema;
+ private final Fields schema;
private final int length;
- private LinkedHashMap fieldTypes;
- private HashMap fieldIndexes;
+ private final LinkedHashMap fieldTypes;
+ private final HashMap fieldIndexes;
public TupleTypeInfo(Fields schema) {
super(Tuple.class);
diff --git a/src/main/java/com/dataartisans/flink/cascading/types/tuple/UnknownTupleSerializer.java b/src/main/java/com/dataartisans/flink/cascading/types/tuple/UnknownTupleSerializer.java
index 7bdcdea..b64a4f4 100644
--- a/src/main/java/com/dataartisans/flink/cascading/types/tuple/UnknownTupleSerializer.java
+++ b/src/main/java/com/dataartisans/flink/cascading/types/tuple/UnknownTupleSerializer.java
@@ -18,6 +18,7 @@
import cascading.tuple.Tuple;
import org.apache.flink.api.common.typeutils.TypeSerializer;
+import org.apache.flink.api.common.typeutils.TypeSerializerSnapshot;
import org.apache.flink.core.memory.DataInputView;
import org.apache.flink.core.memory.DataOutputView;
@@ -201,7 +202,12 @@ public int hashCode() {
return this.fieldSer.hashCode();
}
+ // TODO: address checkpoint serailization
@Override
+ public TypeSerializerSnapshot snapshotConfiguration() {
+ return null;
+ }
+
public boolean canEqual(Object obj) {
return obj instanceof UnknownTupleSerializer;
}
diff --git a/src/main/java/com/dataartisans/flink/cascading/types/tuplearray/TupleArraySerializer.java b/src/main/java/com/dataartisans/flink/cascading/types/tuplearray/TupleArraySerializer.java
index 657e34f..e34b007 100644
--- a/src/main/java/com/dataartisans/flink/cascading/types/tuplearray/TupleArraySerializer.java
+++ b/src/main/java/com/dataartisans/flink/cascading/types/tuplearray/TupleArraySerializer.java
@@ -19,6 +19,7 @@
import cascading.tuple.Tuple;
import com.dataartisans.flink.cascading.types.tuple.NullMaskSerDeUtils;
import org.apache.flink.api.common.typeutils.TypeSerializer;
+import org.apache.flink.api.common.typeutils.TypeSerializerSnapshot;
import org.apache.flink.core.memory.DataInputView;
import org.apache.flink.core.memory.DataOutputView;
@@ -177,6 +178,10 @@ public int hashCode() {
}
@Override
+ public TypeSerializerSnapshot snapshotConfiguration() {
+ return null;
+ }
+
public boolean canEqual(Object obj) {
return obj instanceof TupleArraySerializer;
}
diff --git a/src/main/java/com/dataartisans/flink/cascading/util/FlinkConfigConstants.java b/src/main/java/com/dataartisans/flink/cascading/util/FlinkConfigConstants.java
index be7fcdf..fbbaeb7 100644
--- a/src/main/java/com/dataartisans/flink/cascading/util/FlinkConfigConstants.java
+++ b/src/main/java/com/dataartisans/flink/cascading/util/FlinkConfigConstants.java
@@ -18,8 +18,16 @@
public class FlinkConfigConstants {
- public static final String EXECUTION_MODE = "flink.executionMode";
- public static final String EXECUTION_MODE_BATCH = "BATCH";
- public static final String EXECUTION_MODE_PIPELINED = "PIPELINED";
+ public static final String EXECUTION_MODE = "flink.executionMode";
+ public static final String EXECUTION_MODE_BATCH = "BATCH";
+ public static final String EXECUTION_MODE_PIPELINED = "PIPELINED";
+
+ public static final String PIPELINE_CLASSPATHS = "pipeline.classpaths";
+
+ public static final String LEAK_CLASSLOADER_CHECK = "classloader.check-leaked-classloader";
+
+ public static final String NETWORK_MEMORY_MIN = "taskmanager.memory.network.min";
+
+ public static final String TASKMANAGER_MEMORY_MANAGED_SIZE = "taskmanager.memory.managed.size";
}