如何基于CSV邻接矩阵在Spark GraphX中创建边?
基于CSV邻接矩阵生成Spark GraphX Edges RDD的实现方案
场景说明
我有如下格式的CSV邻接矩阵(空格分隔):
A B C D A 0 3 2 5 B -1 0 2 9 C -1 -1 0 8 D -1 -1 -1 0
已通过其他文件获取Vertices RDD:
val Vertices: RDD[(VertexId, String)] = data.map(_.split(",")).map { arr => val id = arr(0) val place = arr(1) (id.toLong, place) }
需要基于上述邻接矩阵创建Edges RDD,当前代码卡在目标顶点的获取逻辑,具体实现方案如下:
实现步骤
1. 构建顶点标签到ID的映射
首先从已有的Vertices RDD中提取标签(如A/B/C/D)与VertexId的对应关系,方便后续将矩阵中的标签转换为GraphX识别的顶点ID:
// 收集标签到ID的映射,若顶点规模较大,需注意collectAsMap的内存占用,可按需优化 val labelToIdMap: Map[String, VertexId] = Vertices.map { case (id, label) => (label, id) }.collectAsMap()
2. 解析邻接矩阵生成Edges RDD
邻接矩阵的第一行是目标顶点标签,每行的第一个元素是源顶点标签,后续元素为对应目标顶点的边权重(0表示自环,-1表示无连接)。通过以下代码解析并生成有效边:
// 拆分表头与数据行 val header = edgesData.first().split("\\s+").filter(_.nonEmpty) // 适配空格分隔的格式,过滤空字符串 val dataRows = edgesData.filter(_ != header.mkString(" ")) // 跳过表头行 val edges: RDD[Edge[Double]] = dataRows.flatMap { line => val parts = line.split("\\s+").filter(_.nonEmpty) val sourceLabel = parts(0) val sourceId = labelToIdMap(sourceLabel) // 遍历权重与目标标签的对应关系 parts.drop(1).zip(header).map { case (weightStr, targetLabel) => val weight = weightStr.toDouble (sourceId, labelToIdMap(targetLabel), weight) }.filter { case (src, dst, weight) => // 过滤自环和无连接的无效边 src != dst && weight != -1.0 }.map { case (src, dst, weight) => Edge(src, dst, weight) } }
关键说明
- 邻接矩阵为空格分隔,因此使用
\\s+拆分字符串,并过滤空字符串以处理行首的空格 - 使用
flatMap是因为每行源顶点对应多条目标顶点的边,需要将单行转换为多个边对象 - 通过过滤条件排除自环(源顶点等于目标顶点)和无连接(权重为-1)的无效边
内容的提问来源于stack exchange,提问作者dfouheqoijefoih
相关产品推荐
相关产品推荐

