【问题标题】:How to create a graph from Array[(Any, Any)] using Graph.fromEdgeTuples如何使用 Graph.fromEdgeTuples 从 Array[(Any, Any)] 创建图
【发布时间】:2015-08-10 20:02:18
【问题描述】:

我对 spark 很陌生,但我想根据从 Hive 表中获得的关系创建一个图表。我找到了一个函数,它应该在不定义顶点的情况下允许这样做,但我无法让它工作。

我知道这不是一个可重复的示例,但这是我的代码:

import org.apache.spark.SparkContext
import org.apache.spark.graphx._
import org.apache.spark.rdd.RDD
val sqlContext= new org.apache.spark.sql.hive.HiveContext(sc)
val data = sqlContext.sql("select year, trade_flow, reporter_iso, partner_iso, sum(trade_value_us) from comtrade.annual_hs where length(commodity_code)='2' and not partner_iso='WLD' group by year, trade_flow, reporter_iso, partner_iso").collect()
val data_2010 = data.filter(line => line(0)==2010)
val couples = data_2010.map(line=>(line(2),line(3)) //country to country 

val graph = Graph.fromEdgeTuples(couples, 1)

最后一行产生以下错误:

val graph = Graph.fromEdgeTuples(sc.parallelize(couples), 1)
<console>:31: error: type mismatch;
found   : Array[(Any, Any)]
required: Seq[(org.apache.spark.graphx.VertexId,org.apache.spark.graphx.VertexId)]
Error occurred in an application involving default arguments.
val graph = Graph.fromEdgeTuples(sc.parallelize(couples), 1)

情侣长这样:

couples: Array[(Any, Any)] = Array((MWI,MOZ), (WSM,AUS), (MDA,CRI), (KNA,HTI), (PER,ERI), (SWE,CUB), (DEU,PRK), (THA,DJI), (BIH,SVK), (RUS,THA), (SGP,BLR), (MEX,TGO), (TUR,ZAF), (ZWE,SYC), (UGA,GHA), (OMN,SVN), (NZL,SYR), (CHE,SLV), (CZE,LUX), (TGO,COM), (TTO,WLF), (NGA,PAN), (FJI,UKR), (BRA,ECU), (EGY,SWE), (ITA,ARG), (MUS,MLT), (MDG,DZA), (ARE,SUR), (CAN,GUY), (OMN,COG), (NAM,FIN), (ITA,HMD), (SWE,CHE), (SDN,NER), (TUN,USA), (THA,GMB), (HUN,TTO), (FRA,BEN), (NER,TCD), (CHN,JPN), (DNK,ZAF), (MLT,UKR), (ARM,OMN), (PRT,IDN), (BEN,PER), (TTO,BRA), (KAZ,SMR), (CPV,""), (ARG,ZAF), (BLR,TJK), (AZE,SVK), (ITA,STP), (MDA,IRL), (POL,SVN), (PRY,ETH), (HKG,MOZ), (QAT,GAB), (THA,MUS), (PHL,MOZ), (ITA,SGS), (ARM,KHM), (ARG,KOR), (AUT,GMB), (SYR,COM), (CZE,GBR), (DOM,USA), (CYP,LAO), (USA,LBR)

如何转换为合适的格式?

【问题讨论】:

    标签: scala apache-spark apache-spark-sql spark-graphx


    【解决方案1】:

    首先,您不能将String 用作VertexId,因此您必须将标签映射到Long。然后,我们需要准备一个从 label 到 id 的映射。只要唯一值的数量比较少,最简单的方法就是创建一个广播变量:

    val idMap = sc.broadcast(couples // -> Array[(Any, Any)]
      // Make sure we use String not Any returned from Row.apply
      // And convert to Seq so we can flatten results
      .flatMap{case (x: String, y: String) => Seq(x, y)} // -> Array[String]
      // Get different keys
      .distinct // -> Array[String]
      // Create (key, value) pairs
      .zipWithIndex  // -> Array[(String, Int)]
      // Convert values to Long so we can use it as a VertexId
      .map{case (k, v) => (k, v.toLong)}  // -> Array[(String, Long)]
      // Create map
      .toMap) // -> Map[String,Long]
    

    接下来我们可以使用上面的来进行映射:

    val edges: RDD[(VertexId, VertexId)] = sc.parallelize(couples
      .map{case (x: String, y: String) => (idMap.value(x), idMap.value(y))}
    )
    

    最后我们得到一个图表:

    val graph = Graph.fromEdgeTuples(edges, 1)
    

    【讨论】:

    • 哇,谢谢!明天我会先试试这个。您介意详细说明您使用的不同方法吗?我明白了大致的想法,但能够理解每一步而不只是复制它对我来说非常有用
    • 当然,我已经添加了一些 cmets 和类型信息。
    • 完全正常,非常感谢您的解释。你知道这是否可以用 Spark 可视化图表,我正在使用控制台,所以我猜没有图形界面?
    • 我对此表示怀疑。您可以随时收集感兴趣的子图并使用 Gephi 等通用图形处理库。
    • 如果你觉得无聊,我有一个新问题 :-) stackoverflow.com/questions/31944525/…
    猜你喜欢
    • 2019-06-13
    • 1970-01-01
    • 2019-10-10
    • 1970-01-01
    • 2020-11-15
    • 1970-01-01
    • 2015-11-30
    • 2020-04-04
    • 2016-11-06
    相关资源
    最近更新 更多