如何基于Spark DataFrame使用GraphX构建酒店-作者加权关系图
Spark GraphX 构建酒店-作者二部图实现方法
基于已有的包含hotel_url、author、rating字段的评论DataFrame,按以下步骤即可生成符合要求的图结构:
前置说明
GraphX 要求所有顶点必须携带全局唯一的Long类型ID,边必须明确指定源顶点ID、目标顶点ID和边属性,因此不能直接用字符串格式的酒店url、作者名构建图,需要先做ID映射。
步骤1:生成全局唯一节点ID映射
先把两类节点从原始数据中提取、去重,统一分配唯一ID,同时保留节点类型、原始属性用于后续识别:
import org.apache.spark.sql.functions._ import org.apache.spark.graphx._ import org.apache.spark.sql.types.LongType // 提取两类去重节点,打类型标签 val hotelNodes = review_df .select(col("hotel_url").as("node_attr")) .withColumn("node_type", lit("hotel")) .distinct() val authorNodes = review_df .select(col("author").as("node_attr")) .withColumn("node_type", lit("author")) .distinct() // 合并全量节点,分配全局唯一Long ID val allNodes = hotelNodes.union(authorNodes) .withColumn("node_id", monotonically_increasing_id().cast(LongType)) .cache()
步骤2:构建顶点RDD
GraphX 顶点RDD的格式要求为RDD[(VertexId, 节点属性)],直接从节点映射表转换即可:
val vertexRDD = allNodes.rdd.map(row => { val nodeId = row.getAs[Long]("node_id") val nodeAttr = row.getAs[String]("node_attr") val nodeType = row.getAs[String]("node_type") (nodeId, (nodeType, nodeAttr)) })
步骤3:构建边RDD
将原始评论数据和节点ID映射表关联,拿到每条评论对应的酒店ID、作者ID,再转换为GraphX要求的Edge格式。
由于需要的是酒店和作者的无向关联关系,这里生成双向边保证图遍历的时候不会因为方向丢失关联;如果只需要单向有向语义,删掉反向生成的Edge即可。
// 拆分出酒店、作者各自的ID映射表 val hotelIdMap = allNodes .filter(col("node_type") === "hotel") .select(col("node_attr").as("hotel_url"), col("node_id").as("hotel_id")) val authorIdMap = allNodes .filter(col("node_type") === "author") .select(col("node_attr").as("author"), col("node_id").as("author_id")) // 关联拿到每条边两端的节点ID val edgeDF = review_df .join(hotelIdMap, Seq("hotel_url")) .join(authorIdMap, Seq("author")) .select("hotel_id", "author_id", "rating") // 转换为EdgeRDD val edgeRDD = edgeDF.rdd.flatMap(row => { val hotelId = row.getAs[Long]("hotel_id") val authorId = row.getAs[Long]("author_id") val rating = row.getAs[Int]("rating") Seq( Edge(hotelId, authorId, rating), Edge(authorId, hotelId, rating) ) })
步骤4:生成图对象并验证
传入顶点RDD和边RDD即可得到最终的Graph对象,可以通过简单统计验证构建结果:
// 构建图 val hotelAuthorGraph = Graph(vertexRDD, edgeRDD) // 基础统计验证 println(s"节点总数: ${hotelAuthorGraph.numVertices}") println(s"边总数: ${hotelAuthorGraph.numEdges}") // 查看边样例 hotelAuthorGraph.triplets.take(10).foreach(trip => { println(s"${trip.srcAttr._2} 与 ${trip.dstAttr._2} 的关联评分: ${trip.attr}") })
注意事项
- 超大数据集场景下,如果担心
monotonically_increasing_id生成ID冲突,可以替换为RDD的zipWithUniqueId方法,或者用StringIndexer对节点属性做统一编码,ID唯一性更有保障。 - 如果后续需要运行PageRank、连通分量等图算法,直接调用GraphX内置API即可,边权重
rating可以直接作为算法输入参数。
内容的提问来源于stack exchange,提问作者Nihad Huseynov
相关产品推荐
相关产品推荐

