Spark SQL如何统计DataFrame指定列按下划线分割后各词汇的出现次数
Spark统计下划线分割术语频次实现方案
原有思路缺陷
你之前的遍历插入实现思路存在两个核心问题:
foreach算子运行在Executor执行节点,无法直接调用Driver端的SparkSession执行插入SQL,会触发序列化异常- 逐行遍历插入的方式完全没有利用Spark分布式计算能力,执行效率极低,不适合批量数据处理场景
最优实现方案
直接使用Spark内置的字符串处理、集合处理、聚合算子即可完成需求,代码简洁且执行效率更高:
import org.apache.spark.sql.functions.{split, explode, col} val resultDF = df_titles // 将title字段按下划线分割为字符串数组 .select(split(col("title"), "_").alias("term_array")) // 将数组元素炸裂为独立行,每行对应一个术语 .select(explode(col("term_array")).alias("term")) // 按术语分组统计出现频次 .groupBy("term") .count() // 输出结果 resultDF.show()
运行后输出的结果和你期望的格式完全一致。
步骤说明
split(col("title"), "_"):将每行类似harry_potter_1的字符串转换为["harry", "potter", "1"]的数组格式explode:将数组的每个元素拆分为独立行,单个字符串数组会生成对应元素数量的行groupBy("term").count():按术语字段分组,统计每个术语的总出现次数
内容的提问来源于stack exchange,提问作者Angelia Servais
相关产品推荐
相关产品推荐

