如何用PySpark在Hadoop中输出PageRank算法的Top10结果并写入CSV
PySpark PageRank输出Top10并写入CSV的解决方案
为什么你的方法失败
- 方法1问题:你把
ranks.collect()的结果存成了Python普通列表,而show()是Spark DataFrame的专属方法,普通列表没有这个方法,必然报错。 - 方法2问题:一是
rank_rdd.zip(link_rdd)搞反了键值对顺序(原ranks是(link, rank)结构);二是排序用了asc('rank'),要取Top10高排名得用降序;另外代码里top_10 = sorted_df.limit(10)后面缺了分号,属于语法错误。
两种可行的解决办法
办法1:直接用RDD处理
不需要转DataFrame,直接对RDD做排序取Top10,再写入CSV:
# 按rank降序排序,取前10条数据 top10_rdd = ranks.sortBy(lambda x: -x[1]).take(10) # 转成CSV格式的RDD,写入HDFS路径 sc.parallelize(top10_rdd).map(lambda x: f"{x[0]},{x[1]}").saveAsTextFile("hdfs:///your/output/path/pagerank_top10")
注意:
saveAsTextFile会生成多个分区文件,如果必须要单个文件,可以加.coalesce(1),但大数据量下不推荐这么做,会影响性能。
办法2:用Spark DataFrame处理(更简洁)
利用DataFrame的API完成排序、取Top10和写入CSV:
from pyspark.sql import SparkSession from pyspark.sql.functions import desc # 初始化SparkSession(如果代码里还没初始化) spark = SparkSession.builder.appName("PageRankTop10").getOrCreate() # 把ranks RDD转换成DataFrame,指定列名 ranks_df = spark.createDataFrame(ranks, ["link", "rank"]) # 按rank降序排序,取前10 top10_df = ranks_df.orderBy(desc("rank")).limit(10) # 可选:查看Top10结果 top10_df.show() # 写入CSV,带表头,覆盖已有文件 top10_df.write \ .mode("overwrite") \ .option("header", "true") \ .csv("hdfs:///your/output/path/pagerank_top10")
说明:同样,若需要单个CSV文件,可在
write前加.coalesce(1),比如top10_df.coalesce(1).write...。
注意事项
- 确保你指定的HDFS路径有写入权限,路径格式正确(比如完整路径
hdfs://namenode:9000/user/xxx/pagerank_top10,或者用相对路径)。 - 如果在YARN集群运行,要保证Spark的资源配置满足需求。
内容的提问来源于stack exchange,提问作者Vitaliy Plokhovskyy
相关产品推荐
相关产品推荐

