GraphFrames中PageRank执行报错问题求助
嘿,我来帮你拆解这些PageRank的报错问题!结合你用的Spark 1.6和GraphFrames环境,这些问题大多和数据格式不匹配以及旧版API的使用方式不对有关,咱们一步步来解决:
先明确核心问题:顶点与边的ID不匹配
从你提供的顶点和边数据截图能看出,顶点的id是整数类型的station_ID,但边的src/dst却是字符串类型的车站名称——这是所有报错的根源!GraphFrame要求边的src/dst必须和顶点的id完全匹配(类型、值都要一致),否则无法关联顶点和边,触发各种异常。
逐个分析你的报错:
1. g.pageRank(resetProbability=0.15, maxIter=10) 报scala.MatchError
原因:边的src/dst是车站名称(字符串),和顶点的整数id类型不匹配,GraphFrame无法匹配顶点,触发Scala的匹配错误。
2. g.pageRank(resetProbability=0.15, maxIter=10, sourceId="id") 报NoSuchVertexException
原因:sourceId需要传入具体的顶点ID值(比如某个车站的整数ID,如10),而你传入的是列名"id",相当于找ID为字符串"id"的顶点,当然不存在。
3. ranks = g.pageRank.resetProbability(0.15).maxIter(10).run() 报AttributeError
原因:这种链式调用是**新版GraphFrames(适配Spark 2.x+)**的API,而Spark 1.6对应的GraphFrames是旧版(比如0.1.0-spark1.6),旧版不支持这种写法,只能通过直接传参数的方式调用。
4. ranks = g.pageRank(resetProbability=0.15, maxIter=10).run() 报Py4JJavaError
原因:一是依然存在src/dst与顶点id不匹配的问题;二是旧版GraphFrames的pageRank方法不需要加.run(),直接调用即可返回结果。
完整修复步骤
步骤1:修复边数据的src/dst,匹配顶点的id
我们需要把边中的车站名称映射成对应的整数station_ID,这里用join的方式(适合大数据量,避免内存溢出):
# 先保留顶点的id和name列,用于关联 station_name_map = verticesRDD.select("id", "name") # 第一步:给边关联起点的station_ID作为src edges_with_src = edgesRDD.join( station_name_map, edgesRDD["Start Station"] == station_name_map["name"], "inner" # 只保留能匹配到顶点的边,避免null值 ).select(edgesRDD["*"], station_name_map["id"].alias("src")) # 第二步:给边关联终点的station_ID作为dst edges_fixed = edges_with_src.join( station_name_map, edges_with_src["End Station"] == station_name_map["name"], "inner" ).select(edges_with_src["*"], station_name_map["id"].alias("dst")) # 去掉不需要的列(原站名、Trip ID),只保留GraphFrame需要的src/dst和其他边属性 edges_fixed = edges_fixed.drop("Start Station", "End Station", "Trip ID")
步骤2:用修复后的边创建GraphFrame
g = GraphFrame(verticesRDD, edges_fixed)
步骤3:用旧版GraphFrames API正确调用PageRank
Spark 1.6对应的GraphFrames版本需要直接传参数调用,不需要链式方法和.run():
# 运行PageRank ranks = g.pageRank(maxIter=10, resetProbability=0.15) # 查看排序后的PageRank结果 ranks.vertices.sort(desc("pagerank")).show()
如果要运行个性化PageRank(指定sourceId),传入具体的顶点ID值即可:
# 假设1是存在的station_ID personalized_ranks = g.pageRank(maxIter=10, resetProbability=0.15, sourceId=1)
额外提醒
- 注册临时表的步骤确实不需要,GraphFrame直接基于DataFrame操作,不需要注册成SQL临时表。
- 确保你的GraphFrames版本是适配Spark 1.6的,比如
graphframes-0.1.0-spark1.6.jar,版本不兼容也会导致各种奇怪的报错。
内容的提问来源于stack exchange,提问作者vikram

