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

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为非正数时出现异常。

两种方法运行后,都会得到你预期的结果:

ocIdfreqnameFilewordrepeat_word
12fileA.txtanalysesanalyses,analyses
23fileB.txtcarriedcarried,carried,carried
31fileC.txtatlanticatlantic
42fileD.txthellohello,hello

内容的提问来源于stack exchange,提问作者jawad adari

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 10:12:42