You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何对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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.29 21:36:06