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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 04:26:01