Spark Scala不使用SQL实现单元格内分隔文本去重的方法
纯Scala实现方案
基于Spark DataFrame API即可实现该需求,无需编写SQL语句,以下是两种可选实现方式:
方式1:使用Spark内置函数(Spark 2.4+推荐)
无需自定义UDF,执行性能更优:
// 导入依赖的内置函数 import org.apache.spark.sql.functions.{split, array_distinct, concat_ws, trim, col, transform, array_remove} val processedDF = rawDF // 按分号拆分技能字符串为数组 .withColumn("skills_arr", split(col("skills"), ";")) // 对每个技能字符串去除前后空格 .withColumn("skills_trimmed", transform(col("skills_arr"), skill => trim(skill))) // 移除空字符串(处理末尾多余分号的场景) .withColumn("skills_clean", array_remove(col("skills_trimmed"), "")) // 数组元素去重 .withColumn("skills_distinct", array_distinct(col("skills_clean"))) // 拼接回分号分隔的字符串 .withColumn("skills", concat_ws("; ", col("skills_distinct"))) // 删除中间临时列 .drop("skills_arr", "skills_trimmed", "skills_clean", "skills_distinct")
方式2:自定义UDF实现(兼容全Spark版本)
如果使用的是Spark 2.4以下版本,可通过UDF实现通用逻辑:
import org.apache.spark.sql.functions.col import org.apache.spark.sql.expressions.UserDefinedFunction import org.apache.spark.sql.types.StringType // 定义去重逻辑的UDF val distinctSkills: UserDefinedFunction = udf((skillsStr: String) => { Option(skillsStr) match { case None => "" case Some(str) => str.split(";") .map(_.trim) .filter(_.nonEmpty) .distinct .mkString("; ") } }, StringType) // 直接替换原skills列即可 val processedDF = rawDF.withColumn("skills", distinctSkills(col("skills")))
两种实现都可以达到预期效果:原始java; c++ ; java的技能列处理后会变为java; c++,其余列(candidate_id、join_date、location)会保持不变。
内容的提问来源于stack exchange,提问作者rapidminer
相关产品推荐
相关产品推荐

