如何对PySpark DataFrame前100行的Row列表执行单词统计
问题解决方法
你的报错核心原因有两个:
df6.head(100)返回的是pyspark.sql.Row对象组成的列表,你直接将该列表转RDD后,传入flatMap的参数是Row对象而非字符串,无法直接调用split方法- 原代码没有对
rdd100rows做定义就直接调用,属于变量未定义错误
修改后的完整代码
# 配置环境变量 注意要放到SparkContext初始化之前 import os os.environ["JAVA_HOME"] = "/usr/lib/jvm/java-8-openjdk-amd64" os.environ["SPARK_HOME"] = "/content/drive/MyDrive/IDS561_BigData/spark-3.1.2-bin-hadoop3.2" # 取前100行text_字段的内容,提取为纯字符串列表 df6 = df5.select("text_") df6_100rows = df6.head(100) # 从Row对象中提取text_字段的字符串内容 text_list = [row["text_"] for row in df6_100rows] # 初始化SparkContext并生成RDD from pyspark import SparkContext sc = SparkContext.getOrCreate() rdd100rows = sc.parallelize(text_list) # 单词统计 counts = rdd100rows.flatMap(lambda line: line.split(" "))\ .map(lambda word: (word, 1))\ .reduceByKey(lambda a, b: a + b)\ .collect() print(counts)
可选优化方案(无需转RDD,直接用DataFrame API实现,性能更高)
from pyspark.sql.functions import split, explode, count # 取前100行,拆分单词,统计计数 word_count = df5.select("text_").limit(100)\ .select(explode(split("text_", " ")).alias("word"))\ .groupBy("word")\ .agg(count("*").alias("count"))\ .orderBy("count", ascending=False) # 打印结果 word_count.show()
内容的提问来源于stack exchange,提问作者Lina
相关产品推荐
相关产品推荐

