import org.apache.spark.api.java.; import org.apache.spark.SparkConf; import org.apache.spark.api.java.function.Function; import org.apache.spark.graphx.; import org.apache.spark.graphx.util.GraphGenerators; import scala.Tuple2; import scala.Tuple3; import scala.collection.Iterator; import scala.collection.JavaConversions; import java.util.ArrayList; import java.util.Arrays;

public class SsspJava {

public static void main(String[] args) { SparkConf conf = new SparkConf().setAppName("SSSPJava").setMaster("local[*]"); JavaSparkContext sc = new JavaSparkContext(conf);

Graph<Object, Double> graph = GraphGenerators.logNormalGraph(sc.sc(), 5, sc.defaultParallelism(), 4.0, 1.3, 0).mapEdges(new Function<EdgeContext<Object, Double, Double>, Double>() {
  public Double call(EdgeContext<Object, Double, Double> e) throws Exception {
    return e.attr().doubleValue();
  }
});

JavaRDD<Tuple2<Object, Object>> edges = graph.edges().toJavaRDD().map(new Function<Edge<Object>, Tuple2<Object, Object>>() {
  public Tuple2<Object, Object> call(Edge<Object> e) throws Exception {
    return new Tuple2<Object, Object>(e.srcId(), e.dstId());
  }
});

JavaRDD<Tuple2<Object, Tuple2<Double, List<Object>>>> vertices = graph.vertices().toJavaRDD().map(new Function<Tuple2<Object, Object>, Tuple2<Object, Tuple2<Double, List<Object>>>>() {
  public Tuple2<Object, Tuple2<Double, List<Object>>> call(Tuple2<Object, Object> v) throws Exception {
    if (v._1().equals(0)) {
      return new Tuple2<Object, Tuple2<Double, List<Object>>>(v._1(), new Tuple2<Double, List<Object>>(0.0, new ArrayList<Object>(Arrays.asList(v._1()))));
    } else {
      return new Tuple2<Object, Tuple2<Double, List<Object>>>(v._1(), new Tuple2<Double, List<Object>>(Double.POSITIVE_INFINITY, new ArrayList<Object>(Arrays.asList(v._1()))));
    }
  }
});

Graph<Tuple2<Double, List<Object>>, Double> initialGraph = Graph.apply(vertices.rdd(), edges.rdd(), new Tuple2<Double, List<Object>>(Double.POSITIVE_INFINITY, new ArrayList<Object>(Arrays.asList(0.0))));

GraphOps<Tuple2<Double, List<Object>>, Double> sssp = initialGraph.ops().pregel(new Tuple2<Double, List<Object>>(Double.POSITIVE_INFINITY, new ArrayList<Object>(Arrays.asList(0.0))), Integer.MAX_VALUE, EdgeDirection.Out(),
  new Function3<Object, Tuple2<Double, List<Object>>, Tuple2<Double, List<Object>>, Tuple2<Double, List<Object>>>() {
    public Tuple2<Double, List<Object>> call(Object id, Tuple2<Double, List<Object>> dist, Tuple2<Double, List<Object>> newDist) throws Exception {
      if (dist._1() < newDist._1()) {
        return dist;
      } else {
        return newDist;
      }
    }
  },
  new Function1<EdgeContext<Object, Double, Tuple2<Double, List<Object>>>, Iterator<Tuple2<Object, Tuple2<Double, List<Object>>>>>() {
    public Iterator<Tuple2<Object, Tuple2<Double, List<Object>>>> call(EdgeContext<Object, Double, Tuple2<Double, List<Object>>> triplet) throws Exception {
      if (triplet.srcAttr()._1() < triplet.dstAttr()._1() - triplet.attr()) {
        return JavaConversions.asScalaIterator(Arrays.asList(new Tuple2<Object, Tuple2<Double, List<Object>>>(triplet.dstId(), new Tuple2<Double, List<Object>>(triplet.srcAttr()._1() + triplet.attr(), new ArrayList<Object>(triplet.srcAttr()._2().$colon$plus(triplet.dstId())))))).iterator();
      } else {
        return JavaConversions.asScalaIterator(new ArrayList<Tuple2<Object, Tuple2<Double, List<Object>>>>().iterator());
      }
    }
  },
  new Function2<Tuple2<Double, List<Object>>, Tuple2<Double, List<Object>>, Tuple2<Double, List<Object>>>() {
    public Tuple2<Double, List<Object>> call(Tuple2<Double, List<Object>> a, Tuple2<Double, List<Object>> b) throws Exception {
      if (a._1() < b._1()) {
        return a;
      } else {
        return b;
      }
    }
  }
);

System.out.println(sssp.vertices().toJavaRDD().collect().mkString("\n"));

sc.stop();
sc.close();

} }

Convert Scala SSSP Implementation in GraphX to Java Code

原文地址: https://www.cveoy.top/t/topic/pfSi 著作权归作者所有。请勿转载和采集!

免费AI点我,无需注册和登录