基于Spark DataFrame字符串列创建新列:提取指定序列目标数字
解决方案:从Spark DataFrame字符串列提取指定数字
咱们一步步来实现这个需求,先看你的原始数据:
val df= Seq(("0003C32C-FC1D-482F-B543-3CBD7F0A0E36 0,8,1,799,300:3 0,6,1,330,300:1 2,6,1,15861:1 0,7,1,734,300:1 0,6,0,95,300:1 2,7,1,15861:1 0,8,0,134,300:3")).toDF("col_str")
核心需求是:从col_str中以0开头的序列里,提取第二个数字和冒号后的最后一个数字,生成新列。下面分两种场景给出实现方式:
方式一:用Spark内置函数(Spark 3.0+推荐)
这种方式不用写UDF,性能更优,先导入必要的函数:
import org.apache.spark.sql.functions._
然后执行数据处理:
val resultDF = df // 1. 按空格拆分字符串为数组,过滤出以"0,"开头的有效序列(避开开头的UUID) .withColumn("valid_seqs", split(col("col_str"), " ").filter(_.startsWith("0,"))) // 2. 遍历每个有效序列,提取第二个数字和最后一个数字,封装成结构体数组 .withColumn("extracted_pairs", expr(""" transform(valid_seqs, seq -> struct( cast(split(split(seq, ':')[0], ',')[1] as int) as second_num, cast(split(seq, ':')[1] as int) as last_num )) """)) // 可选:如果需要把每个数字对展开成单独行,添加这一步 .withColumn("extracted_pair", explode(col("extracted_pairs"))) .select("col_str", "extracted_pair.second_num", "extracted_pair.last_num")
代码解释:
split(col("col_str"), " "):把字符串按空格拆成数组,第一个元素是UUID,后面是目标序列filter(_.startsWith("0,")):精准筛选我们要的以0开头的序列(避免误匹配UUID)transform:遍历数组中的每个序列,做两步拆分:- 先按
:拆分,冒号后面的直接是最后一个数字 - 冒号前面的部分按
,拆分,取索引1的元素(第二个数字)
- 先按
cast(..., as int):把提取的字符串转成数字类型,方便后续计算explode:把数组展开成多行,每一行对应一个提取的数字对,适合查看或进一步分析
运行resultDF.show(false)会得到这样的结果:
+--------------------------------------------------------------------------------------------------------------------------------------------------------------------+----------+---------+ |col_str |second_num|last_num | +--------------------------------------------------------------------------------------------------------------------------------------------------------------------+----------+---------+ |0003C32C-FC1D-482F-B543-3CBD7F0A0E36 0,8,1,799,300:3 0,6,1,330,300:1 2,6,1,15861:1 0,7,1,734,300:1 0,6,0,95,300:1 2,7,1,15861:1 0,8,0,134,300:3|8 |3 | |0003C32C-FC1D-482F-B543-3CBD7F0A0E36 0,8,1,799,300:3 0,6,1,330,300:1 2,6,1,15861:1 0,7,1,734,300:1 0,6,0,95,300:1 2,7,1,15861:1 0,8,0,134,300:3|6 |1 | |0003C32C-FC1D-482F-B543-3CBD7F0A0E36 0,8,1,799,300:3 0,6,1,330,300:1 2,6,1,15861:1 0,7,1,734,300:1 0,6,0,95,300:1 2,7,1,15861:1 0,8,0,134,300:3|7 |1 | |0003C32C-FC1D-482F-B543-3CBD7F0A0E36 0,8,1,799,300:3 0,6,1,330,300:1 2,6,1,15861:1 0,7,1,734,300:1 0,6,0,95,300:1 2,7,1,15861:1 0,8,0,134,300:3|6 |1 | |0003C32C-FC1D-482F-B543-3CBD7F0A0E36 0,8,1,799,300:3 0,6,1,330,300:1 2,6,1,15861:1 0,7,1,734,300:1 0,6,0,95,300:1 2,7,1,15861:1 0,8,0,134,300:3|8 |3 | +--------------------------------------------------------------------------------------------------------------------------------------------------------------------+----------+---------+
方式二:用UDF兼容低版本Spark(Spark <3.0)
如果你的Spark版本不支持transform函数,就用自定义UDF来处理:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ // 定义UDF:输入字符串数组,返回(第二个数字, 最后一个数字)的数组 val extractPairsUDF = udf((seqs: Array[String]) => { seqs.filter(_.startsWith("0,")) .map(seq => { val colonParts = seq.split(':') val commaParts = colonParts(0).split(',') (commaParts(1).toInt, colonParts(1).toInt) }) }, ArrayType(StructType(Seq( StructField("second_num", IntegerType), StructField("last_num", IntegerType) )))) val resultDF = df .withColumn("valid_seqs", split(col("col_str"), " ")) .withColumn("extracted_pairs", extractPairsUDF(col("valid_seqs")))
这个UDF的逻辑和内置函数一致,只是用Scala代码实现了遍历和提取,兼容性更好。
内容的提问来源于stack exchange,提问作者mikeL
相关产品推荐
相关产品推荐

