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

Spark 2.0(Scala2.11)DataFrame动态子串转换实现

Hey there! Let's work through this DataFrame transformation problem you've got. The core task is parsing the B column by using numeric values from each segment to dictate how much of the next segment to keep, repeating until we've processed all parts. Here's how to do it in Scala for Spark 2.0 (Scala 2.11):

First, clarify the parsing logic

For each row's B value (split by *):

  • Grab the first segment as your initial substring length n.
  • Take the first n characters from the next segment, add that to your result.
  • The leftover part of that segment becomes the new n for the next round.
  • Keep going until all segments are done, then join the results with *.

Build a custom UDF to handle the parsing

Spark's built-in functions aren't great for this kind of iterative, stateful parsing, so we'll write a Scala UDF to wrap the logic:

import org.apache.spark.sql.functions._
import org.apache.spark.sql.DataFrame

// The core parsing function
def parseBValue(bStr: String): String = {
  val segments = bStr.split("\\*")
  if (segments.size < 2) return bStr // Bail out if there's nothing to parse
  
  var currentLength = segments.head.toInt
  val parsedParts = scala.collection.mutable.ListBuffer[String]()
  
  // Iterate through each segment after the first
  segments.tail.foreach { seg =>
    // Extract up to currentLength characters (handle if seg is shorter than expected)
    val extracted = if (seg.length >= currentLength) seg.substring(0, currentLength) else seg
    parsedParts += extracted
    
    // Update currentLength with the remaining part of the segment
    if (seg.length > currentLength) {
      val remaining = seg.substring(currentLength)
      currentLength = remaining.toInt
    } else {
      // If we've used the entire segment, set length to max to take full next segments (edge case handling)
      currentLength = Int.MaxValue
    }
  }
  
  parsedParts.mkString("*")
}

// Register the UDF for use in Spark
val parseBUdf = udf(parseBValue _)

Apply the UDF to your DataFrame

Assuming your original DataFrame is called inputDF, just call withColumn to replace the B column:

val transformedDF = inputDF.withColumn("B", parseBUdf(col("B")))

// Check the result
transformedDF.show(truncate = false)

Expected Output

Running this will produce exactly the format you specified:

AB
aZ1CAThisnThat5182763032fa2
bCFriendsnames*O
c1945helpmeplease
dM1finna3212345*fancy2

Quick Notes

  • We added basic edge case handling for segments shorter than the required length (it just takes the whole segment) and rows with too few segments.
  • This code is fully compatible with Spark 2.0 and Scala 2.11 as you requested.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 07:40:05