如何将PySpark分组后数组拼接为字符串的逻辑转换为Scala Spark实现
Scala Spark 实现方案
方案1:内置函数实现(优先推荐)
Spark 原生提供concat_ws函数,无需自定义UDF即可完成数组转指定分隔符字符串的需求,性能优于自定义实现,完全匹配你的业务逻辑:
// 先导入SQL函数包 import org.apache.spark.sql.functions._ val resultDF = our_dataframe .groupBy("column") .agg( concat_ws(";", collect_set("other_col")).as("other_col_joined") )
concat_ws会自动将数组内的非字符串类型元素转为字符串后拼接,和你原PySpark逻辑中的lambda x: ';'.join(map(str, x))效果完全一致。
方案2:自定义UDF实现
如果后续需要扩展拼接逻辑,可以自定义UDF实现等价效果:
import org.apache.spark.sql.functions._ // 定义接收任意类型数组、返回分号拼接字符串的UDF val semicolonJoinUDF = udf((inputArr: Seq[Any]) => inputArr.mkString(";")) val resultDF = our_dataframe .groupBy("column") .agg( semicolonJoinUDF(collect_set("other_col")).as("other_col_joined") )
注意:无特殊逻辑时不要使用自定义UDF方案,UDF会引入额外的序列化、反序列化开销,大数据量下性能远低于内置
concat_ws函数。
内容的提问来源于stack exchange,提问作者Tom Hill
相关产品推荐
相关产品推荐

