【问题标题】:Java samples for GraphXGraphX 的 Java 示例
【发布时间】:2014-04-28 22:25:06
【问题描述】:

在哪里可以找到与 Java 中的 GraphX example 等效的内容?例如以下如何翻译:

val users: RDD[(VertexId, (String, String))] =
sc.parallelize(Array((3L, ("rxin", "student")), (7L, ("jgonzal", "postdoc")),(5L, ("franklin", "prof")), (2L, ("istoica", "prof"))))

// Create an RDD for edges
val relationships: RDD[Edge[String]] = sc.parallelize(Array(Edge(3L, 7L, "collab"),    Edge(5L, 3L, "advisor"), Edge(2L, 5L, "colleague"), Edge(5L, 7L, "pi")))

// Define a default user in case there are relationship with missing user 
val defaultUser = ("John Doe", "Missing")

// Build the initial Graph
val graph = Graph(users, relationships, defaultUser)

【问题讨论】:

  • 在 dev@spark.apache.org 上得到响应:> 还没有 Java API。
  • 很高兴能用这个回答你自己的问题,并提供一个讨论页面的链接,这样我们也可以保持更新......

标签: apache-spark


【解决方案1】:

在 dev@spark.apache.org 上得到了回复:> 目前还没有 Java API。

users mailing list

【讨论】:

    【解决方案2】:

    仍然没有可用于 Java 的 API: “GraphX 只能从 Scala API 获得。”

    来源:https://forums.databricks.com/questions/3185/i-am-looking-for-some-tutorial-on-graphx-using-jav.html

    【讨论】:

    【解决方案3】:

    GraphX 仅在 Scala 中可用。您可以查看 Graphframe,他们在其中寻找数据帧(Java、Python、Scala)——而不是低级 RDD。它还有一些优势,因为它可以利用查询优化器 Catalyst、Tungsten 项目优化。

    您可以将 GraphX 转换为 Graphframes,反之亦然。 VertexIds 是任何类型的,不像 GraphX 只有 Long。返回类型是 Dataframe 或 Graphframe 但在 GraphX 中只有 Graph[VD,ED], RDD。

    【讨论】:

      【解决方案4】:

      根据[SPARK-3665] Java API for GraphX - ASF JIRA:“目标版本/s:1.3.0”

      【讨论】:

        【解决方案5】:

        虽然没有官方或非官方的文档但是有graphx library for Java。检查thisthis

        【讨论】:

          【解决方案6】:

          你可以这样使用:

              public static void main(String[] args) {
          
          
                  SparkSession spark = SparkSession
                          .builder()
                          .appName("javaGraphx")
                          .getOrCreate();
          
                  List<Tuple2<Object, Data>> vectorList = Arrays.asList(
                          new Tuple2[]{
                                  new Tuple2(1L, new Data("A", 1)),
                                  new Tuple2(2L, new Data("B", 1)),
                                  new Tuple2(3L, new Data("C", 1))}
                  );
          
                  JavaRDD<Tuple2<Object, Data>> users = spark.sparkContext().parallelize(
                          JavaConverters.asScalaIteratorConverter(vectorList.iterator()).asScala()
                                  .toSeq(), 1, scala.reflect.ClassTag$.MODULE$.apply(Tuple2.class)
                  ).toJavaRDD();
          
                  List<Edge<Integer>> edgeList = Arrays.asList(
                          new Edge[]{
                                  new Edge(1L, 2L, 1),
                                  new Edge(2L, 3L, 1)});
          
                  JavaRDD<Edge<Integer>> followers = spark.sparkContext().parallelize(
                          JavaConverters.asScalaIteratorConverter(edgeList.iterator()).asScala()
                                  .toSeq(), 1, scala.reflect.ClassTag$.MODULE$.apply(Edge.class)
                  ).toJavaRDD();
          
                  // construct graph
                  Graph<OddRange, Integer> followerGraph =
                          GraphImpl
                                  .apply(
                                          users.rdd(),
                                          followers.rdd(),
                                          new Data("OVER", 0),
                                          StorageLevel.MEMORY_AND_DISK(),
                                          StorageLevel.MEMORY_AND_DISK(),
                                          scala.reflect.ClassTag$.MODULE$.apply(Data.class),
                                          scala.reflect.ClassTag$.MODULE$.apply(Integer.class)
                                      );
          //If you wanna use pregel, you can use it like this.
                  Graph<OddRange, Integer> dd = Pregel.apply(followerGraph,
                          new Data("", 0),
                          5,
                          EdgeDirection.Out(),
                          new Vprog(),
                          new SendMsg(),
                          new MergemMsg(),
                          scala.reflect.ClassTag$.MODULE$.apply(OddRange.class),
                          scala.reflect.ClassTag$.MODULE$.apply(Integer.class),
                          scala.reflect.ClassTag$.MODULE$.apply(OddRange.class)
                  );
          
                  RDD<OddRange> r = dd
                          .vertices()
                          .toJavaRDD()
                          .map(t->t._2)
                          .rdd();
          
                  r.toJavaRDD().foreach(
                          o->{
                              System.out.println(o.getS());
                          }
                          );
              }
          
          //If you wanna use pregel, you can use it like this.
              static class Vprog extends AbstractFunction3< Object, Data, Data, Data> implements Serializable {
                  @Override
                  public Data apply(Object l, Data self, Data sumOdd) {
                      System.out.println(l + "---" + self.getS()+self.getI()+" ---> "+sumOdd.getS()+sumOdd.getI());
                          self.setS(sumOdd.getS() + self.getS());
                          self.setI(self.getI() + sumOdd.getI());
                      System.out.println(l + "---" + self.getS()+self.getI()+" ---> "+sumOdd.getS()+sumOdd.getI());
                          //Don't just return self here, return a new one;
                          return new Data(self.getS(), self.getI());
                  }
              }
          
              static class SendMsg extends AbstractFunction1<EdgeTriplet<Data,Integer>, scala.collection.Iterator<Tuple2<Object, Data>>> implements Serializable {
                  @Override
                  public scala.collection.Iterator<Tuple2<Object, Data>> apply(EdgeTriplet<Data,Integer> t) {
                      System.out.println(t.srcId()+" ---> "+t.dstId()+" with: "+t.srcAttr().getS()+t.srcAttr().getI()+" ---> "+t.dstAttr().getS()+t.dstAttr().getI());
          
                      if(t.srcAttr().getI() <= 8){
                          List<Tuple2<Object, Data>> data = new ArrayList();
                          data.add(new Tuple2<>( t.dstId(), new Data(t.srcAttr().getS(), t.srcAttr().getI())));
                          return JavaConverters.asScalaIteratorConverter(data.iterator()).asScala();
                      }else{
                          return JavaConverters.asScalaIteratorConverter(new ArrayList<Tuple2<Object, Data>>().iterator()).asScala();
                      }
                  }
              }
          
              static class MergemMsg extends AbstractFunction2< Data, Data, Data> implements Serializable {
                  @Override
                  public Data apply(Data a, Data b) {
                      return new Data( "" + a.getS() + b.getS(), a.getI() + b.getI());
                  }
              }
          

          【讨论】:

            猜你喜欢
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 2011-08-18
            • 2012-01-19
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            相关资源
            最近更新 更多