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
ncharacters from the next segment, add that to your result. - The leftover part of that segment becomes the new
nfor 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:
| A | B |
|---|---|
| a | Z1CAThisnThat5182763032fa2 |
| b | CFriendsnames*O |
| c | 1945helpmeplease |
| d | M1finna3212345*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
相关产品推荐
相关产品推荐

