From 374d30d41b8009ac9d0f038c61fb157c3f730cf1 Mon Sep 17 00:00:00 2001 From: hguturu Date: Wed, 20 Dec 2017 14:57:36 -0800 Subject: [PATCH 1/2] upgraded dga to giraph 1.2.0 --- dga-giraph/build.gradle | 4 ++-- .../com/soteradefense/dga/hbse/HBSEComputation.java | 8 +++----- .../dga/lc/LeafCompressionComputation.java | 10 ++++------ .../dga/louvain/giraph/LouvainComputation.java | 7 +++---- .../com/soteradefense/dga/pr/PageRankComputation.java | 7 +++---- .../dga/wcc/WeaklyConnectedComponentComputation.java | 7 +++---- 6 files changed, 18 insertions(+), 25 deletions(-) diff --git a/dga-giraph/build.gradle b/dga-giraph/build.gradle index 841252f..944f641 100644 --- a/dga-giraph/build.gradle +++ b/dga-giraph/build.gradle @@ -22,12 +22,12 @@ configurations { dependencies { compile project(':dga-core') - compile('org.apache.giraph:giraph-core:1.1.0-hadoop2') { + compile('org.apache.giraph:giraph-core:1.2.0-hadoop2') { exclude module: 'guava' exclude module: 'zookeeper' } hadoopProvided group: 'org.apache.hadoop', name: 'hadoop-client', version: cdh_version - testCompile 'org.apache.giraph:giraph-core:1.1.0-hadoop2' + testCompile 'org.apache.giraph:giraph-core:1.2.0-hadoop2' testCompile "org.mockito:mockito-core:1.9.5" testCompile 'junit:junit:4.11' testCompile 'commons-httpclient:commons-httpclient:3.0.1' diff --git a/dga-giraph/src/main/java/com/soteradefense/dga/hbse/HBSEComputation.java b/dga-giraph/src/main/java/com/soteradefense/dga/hbse/HBSEComputation.java index b7e9e1a..3d4e009 100644 --- a/dga-giraph/src/main/java/com/soteradefense/dga/hbse/HBSEComputation.java +++ b/dga-giraph/src/main/java/com/soteradefense/dga/hbse/HBSEComputation.java @@ -18,14 +18,12 @@ package com.soteradefense.dga.hbse; import com.soteradefense.dga.DGALoggingUtil; +import org.apache.giraph.bsp.CentralizedServiceWorker; import org.apache.giraph.comm.WorkerClientRequestProcessor; import org.apache.giraph.edge.Edge; import org.apache.giraph.graph.AbstractComputation; import org.apache.giraph.graph.GraphState; -import org.apache.giraph.graph.GraphTaskManager; import org.apache.giraph.graph.Vertex; -import org.apache.giraph.worker.WorkerAggregatorUsage; -import org.apache.giraph.worker.WorkerContext; import org.apache.giraph.worker.WorkerGlobalCommUsage; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; @@ -55,8 +53,8 @@ public class HBSEComputation extends AbstractComputation workerClientRequestProcessor, GraphTaskManager graphTaskManager, WorkerGlobalCommUsage workerGlobalCommUsage, WorkerContext workerContext) { - super.initialize(graphState, workerClientRequestProcessor, graphTaskManager, workerGlobalCommUsage, workerContext); + public void initialize(GraphState graphState, WorkerClientRequestProcessor workerClientRequestProcessor, CentralizedServiceWorker centralizedServiceWorker, WorkerGlobalCommUsage workerGlobalCommUsage) { + super.initialize(graphState, workerClientRequestProcessor, centralizedServiceWorker, workerGlobalCommUsage); DGALoggingUtil.setDGALogLevel(this.getConf()); } diff --git a/dga-giraph/src/main/java/com/soteradefense/dga/lc/LeafCompressionComputation.java b/dga-giraph/src/main/java/com/soteradefense/dga/lc/LeafCompressionComputation.java index 961c8bb..af2fac3 100644 --- a/dga-giraph/src/main/java/com/soteradefense/dga/lc/LeafCompressionComputation.java +++ b/dga-giraph/src/main/java/com/soteradefense/dga/lc/LeafCompressionComputation.java @@ -18,14 +18,12 @@ package com.soteradefense.dga.lc; import com.soteradefense.dga.DGALoggingUtil; +import org.apache.giraph.bsp.CentralizedServiceWorker; import org.apache.giraph.comm.WorkerClientRequestProcessor; import org.apache.giraph.edge.Edge; import org.apache.giraph.graph.BasicComputation; import org.apache.giraph.graph.GraphState; -import org.apache.giraph.graph.GraphTaskManager; import org.apache.giraph.graph.Vertex; -import org.apache.giraph.worker.WorkerAggregatorUsage; -import org.apache.giraph.worker.WorkerContext; import org.apache.giraph.worker.WorkerGlobalCommUsage; import org.apache.hadoop.io.Text; import org.slf4j.Logger; @@ -44,8 +42,8 @@ public class LeafCompressionComputation extends BasicComputation workerClientRequestProcessor, GraphTaskManager graphTaskManager, WorkerGlobalCommUsage workerGlobalCommUsage, WorkerContext workerContext) { - super.initialize(graphState, workerClientRequestProcessor, graphTaskManager, workerGlobalCommUsage, workerContext); + public void initialize(GraphState graphState, WorkerClientRequestProcessor workerClientRequestProcessor, CentralizedServiceWorker centralizedServiceWorker, WorkerGlobalCommUsage workerGlobalCommUsage) { + super.initialize(graphState, workerClientRequestProcessor, centralizedServiceWorker, workerGlobalCommUsage); DGALoggingUtil.setDGALogLevel(getConf()); } @@ -107,4 +105,4 @@ private int getVertexValue(String value) { } return vertexValue; } -} \ No newline at end of file +} diff --git a/dga-giraph/src/main/java/com/soteradefense/dga/louvain/giraph/LouvainComputation.java b/dga-giraph/src/main/java/com/soteradefense/dga/louvain/giraph/LouvainComputation.java index 1f3b2ad..703d643 100644 --- a/dga-giraph/src/main/java/com/soteradefense/dga/louvain/giraph/LouvainComputation.java +++ b/dga-giraph/src/main/java/com/soteradefense/dga/louvain/giraph/LouvainComputation.java @@ -18,14 +18,13 @@ package com.soteradefense.dga.louvain.giraph; import com.soteradefense.dga.DGALoggingUtil; +import org.apache.giraph.bsp.CentralizedServiceWorker; import org.apache.giraph.comm.WorkerClientRequestProcessor; import org.apache.giraph.edge.Edge; import org.apache.giraph.edge.EdgeFactory; import org.apache.giraph.graph.AbstractComputation; import org.apache.giraph.graph.GraphState; -import org.apache.giraph.graph.GraphTaskManager; import org.apache.giraph.graph.Vertex; -import org.apache.giraph.worker.WorkerContext; import org.apache.giraph.worker.WorkerGlobalCommUsage; import org.apache.hadoop.io.DoubleWritable; import org.apache.hadoop.io.LongWritable; @@ -76,8 +75,8 @@ private void aggregateQ(Double q) { } @Override - public void initialize(GraphState graphState, WorkerClientRequestProcessor workerClientRequestProcessor, GraphTaskManager graphTaskManager, WorkerGlobalCommUsage workerGlobalCommUsage, WorkerContext workerContext) { - super.initialize(graphState, workerClientRequestProcessor, graphTaskManager, workerGlobalCommUsage, workerContext); + public void initialize(GraphState graphState, WorkerClientRequestProcessor workerClientRequestProcessor, CentralizedServiceWorker centralizedServiceWorker, WorkerGlobalCommUsage workerGlobalCommUsage) { + super.initialize(graphState, workerClientRequestProcessor, centralizedServiceWorker, workerGlobalCommUsage); DGALoggingUtil.setDGALogLevel(this.getConf()); } diff --git a/dga-giraph/src/main/java/com/soteradefense/dga/pr/PageRankComputation.java b/dga-giraph/src/main/java/com/soteradefense/dga/pr/PageRankComputation.java index f490d83..fb068f5 100644 --- a/dga-giraph/src/main/java/com/soteradefense/dga/pr/PageRankComputation.java +++ b/dga-giraph/src/main/java/com/soteradefense/dga/pr/PageRankComputation.java @@ -19,12 +19,11 @@ package com.soteradefense.dga.pr; import com.soteradefense.dga.DGALoggingUtil; +import org.apache.giraph.bsp.CentralizedServiceWorker; import org.apache.giraph.comm.WorkerClientRequestProcessor; import org.apache.giraph.graph.BasicComputation; import org.apache.giraph.graph.GraphState; -import org.apache.giraph.graph.GraphTaskManager; import org.apache.giraph.graph.Vertex; -import org.apache.giraph.worker.WorkerContext; import org.apache.giraph.worker.WorkerGlobalCommUsage; import org.apache.hadoop.io.DoubleWritable; import org.apache.hadoop.io.Text; @@ -42,8 +41,8 @@ public class PageRankComputation extends BasicComputation workerClientRequestProcessor, GraphTaskManager graphTaskManager, WorkerGlobalCommUsage workerGlobalCommUsage, WorkerContext workerContext) { - super.initialize(graphState, workerClientRequestProcessor, graphTaskManager, workerGlobalCommUsage, workerContext); + public void initialize(GraphState graphState, WorkerClientRequestProcessor workerClientRequestProcessor, CentralizedServiceWorker centralizedServiceWorker, WorkerGlobalCommUsage workerGlobalCommUsage) { + super.initialize(graphState, workerClientRequestProcessor, centralizedServiceWorker, workerGlobalCommUsage); DGALoggingUtil.setDGALogLevel(this.getConf()); } diff --git a/dga-giraph/src/main/java/com/soteradefense/dga/wcc/WeaklyConnectedComponentComputation.java b/dga-giraph/src/main/java/com/soteradefense/dga/wcc/WeaklyConnectedComponentComputation.java index e69935f..d206d0f 100644 --- a/dga-giraph/src/main/java/com/soteradefense/dga/wcc/WeaklyConnectedComponentComputation.java +++ b/dga-giraph/src/main/java/com/soteradefense/dga/wcc/WeaklyConnectedComponentComputation.java @@ -18,13 +18,12 @@ package com.soteradefense.dga.wcc; import com.soteradefense.dga.DGALoggingUtil; +import org.apache.giraph.bsp.CentralizedServiceWorker; import org.apache.giraph.comm.WorkerClientRequestProcessor; import org.apache.giraph.edge.Edge; import org.apache.giraph.graph.BasicComputation; import org.apache.giraph.graph.GraphState; -import org.apache.giraph.graph.GraphTaskManager; import org.apache.giraph.graph.Vertex; -import org.apache.giraph.worker.WorkerContext; import org.apache.giraph.worker.WorkerGlobalCommUsage; import org.apache.hadoop.io.Text; import org.slf4j.Logger; @@ -40,8 +39,8 @@ public class WeaklyConnectedComponentComputation extends BasicComputation workerClientRequestProcessor, GraphTaskManager graphTaskManager, WorkerGlobalCommUsage workerGlobalCommUsage, WorkerContext workerContext) { - super.initialize(graphState, workerClientRequestProcessor, graphTaskManager, workerGlobalCommUsage, workerContext); + public void initialize(GraphState graphState, WorkerClientRequestProcessor workerClientRequestProcessor, CentralizedServiceWorker centralizedServiceWorker, WorkerGlobalCommUsage workerGlobalCommUsage) { + super.initialize(graphState, workerClientRequestProcessor, centralizedServiceWorker, workerGlobalCommUsage); DGALoggingUtil.setDGALogLevel(this.getConf()); } From 157fb812b15ace91d72bccb60ba52a25ba8dd5f6 Mon Sep 17 00:00:00 2001 From: hguturu Date: Wed, 20 Dec 2017 14:59:21 -0800 Subject: [PATCH 2/2] patch to handle writing to hdfs on newer hadoop (e.g. 2.7.3) --- .../dga/louvain/giraph/LouvainMasterCompute.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/dga-giraph/src/main/java/com/soteradefense/dga/louvain/giraph/LouvainMasterCompute.java b/dga-giraph/src/main/java/com/soteradefense/dga/louvain/giraph/LouvainMasterCompute.java index fb2b36e..32be99e 100644 --- a/dga-giraph/src/main/java/com/soteradefense/dga/louvain/giraph/LouvainMasterCompute.java +++ b/dga-giraph/src/main/java/com/soteradefense/dga/louvain/giraph/LouvainMasterCompute.java @@ -191,7 +191,8 @@ private void writeFile(String path, String message) { Path pt = new Path(path); logger.debug("Writing file out to {}, message {}", path, message); try { - FileSystem fs = FileSystem.get(new Configuration()); + //FileSystem fs = FileSystem.get(new Configuration()); + FileSystem fs = FileSystem.get(getConf()); BufferedWriter br = new BufferedWriter(new OutputStreamWriter(fs.create(pt, true))); br.write(message); br.close(); @@ -206,7 +207,8 @@ private String readFile(String path) { StringBuilder builder = new StringBuilder(); try { Path pt = new Path(path); - FileSystem fs = FileSystem.get(new Configuration()); + //FileSystem fs = FileSystem.get(new Configuration()); + FileSystem fs = FileSystem.get(getConf()); BufferedReader br = new BufferedReader(new InputStreamReader(fs.open(pt))); String line; line = br.readLine();