Spark DataFrame:按freq列值重复word列内容生成新列
解决Spark DataFrame按freq重复word生成逗号分隔新列的问题
嘿,这个需求我之前处理过,给你两种实用的解决思路,都是Spark里常用的方案:
方案一:使用Spark内置函数(推荐,性能更优)
Spark自带的array_repeat和concat_ws组合就能完美搞定,不用写自定义UDF,执行效率更高,而且更符合Spark的分布式计算特性。
代码示例
import org.apache.spark.sql.functions.{array_repeat, concat_ws} // 先模拟你的DataStructure创建示例数据 val df = spark.createDataFrame(Seq( (1, 2, "fileA.txt", "analyses"), (2, 3, "fileB.txt", "carried"), (3, 1, "fileC.txt", "atlantic"), (4, 2, "fileD.txt", "hello") )).toDF("ocId", "freq", "nameFile", "word") // 生成目标新列repeat_word val resultDF = df.withColumn("repeat_word", concat_ws(",", array_repeat($"word", $"freq"))) // 查看最终结果 resultDF.show(false)
逻辑拆解
array_repeat($"word", $"freq"):根据freq的数值,把word字符串重复对应次数,生成一个数组。比如freq=2时,会生成["analyses", "analyses"]这样的数组。concat_ws(",", ...):把数组里的所有元素用逗号连接成一个完整字符串,正好匹配你要的格式。
方案二:自定义UDF(适合后续复杂扩展)
如果之后你有更特殊的重复逻辑(比如需要加额外分隔符、过滤空值等),可以写个自定义UDF来实现:
代码示例
import org.apache.spark.sql.functions.udf // 定义UDF:接收word和freq参数,返回重复后的逗号分隔字符串 val repeatWordUdf = udf((word: String, freq: Int) => { // 处理freq为0或负数的边界情况 if (freq <= 0) "" else List.fill(freq)(word).mkString(",") }) // 应用UDF生成新列 val resultDFWithUdf = df.withColumn("repeat_word", repeatWordUdf($"word", $"freq")) // 查看结果 resultDFWithUdf.show(false)
逻辑拆解
自定义UDF里用List.fill(freq)(word)生成包含指定次数word的列表,再用mkString(",")转成逗号分隔的字符串,同时加了边界处理,避免freq为非正数时出现异常。
两种方法运行后,都会得到你预期的结果:
| ocId | freq | nameFile | word | repeat_word |
|---|---|---|---|---|
| 1 | 2 | fileA.txt | analyses | analyses,analyses |
| 2 | 3 | fileB.txt | carried | carried,carried,carried |
| 3 | 1 | fileC.txt | atlantic | atlantic |
| 4 | 2 | fileD.txt | hello | hello,hello |
内容的提问来源于stack exchange,提问作者jawad adari
相关产品推荐
相关产品推荐

