如何使用Spark DataFrame修剪SCHDULE列的值?
问题
我有如下Spark DataFrame,需要修剪SCHDULE列的值,尝试用UDF没得到预期结果。
原始DataFrame:
| SCHDULE | ID | VALUE |
|---|---|---|
| 100H/10AR1 | KL01 | 30 |
| 100H/10TR2 | KL01 | 40 |
| 100H/22TR1 | KL01 | 20 |
| 100H/22TR2 | KL01 | 20 |
| 105JK/12PK1 | AA05 | 10 |
| 105JK/12PK2 | AA05 | 20 |
| 105JH/33PK3 | AA05 | 50 |
| 105JH/33PK4 | AA05 | 30 |
| 110P/1 | BR03 | 20 |
| 110P/2 | BR03 | 10 |
目标输出(注:当前输入输出表格内容一致,推测你需要对SCHDULE做特定字符串修剪,以下基于常见场景给出方案):
| SCHDULE | ID | VALUE |
|---|---|---|
| 100H/10AR1 | KL01 | 30 |
| 100H/10TR2 | KL01 | 40 |
| 100H/22TR1 | KL01 | 20 |
| 100H/22TR2 | KL01 | 20 |
| 105JK/12PK1 | AA05 | 10 |
| 105JK/12PK2 | AA05 | 20 |
| 105JH/33PK3 | AA05 | 50 |
| 105JH/33PK4 | AA05 | 30 |
| 110P/1 | BR03 | 20 |
| 110P/2 | BR03 | 10 |
解决方案
Spark内置字符串函数比UDF更高效,以下是常见修剪场景的实现:
提取斜杠前的部分
// Scala import org.apache.spark.sql.functions._ val trimmedDF = originalDF.withColumn("SCHDULE", split(col("SCHDULE"), "/").getItem(0))
# Python from pyspark.sql.functions import split, col trimmed_df = original_df.withColumn("SCHDULE", split(col("SCHDULE"), "/").getItem(0))
去除首尾空格
// Scala import org.apache.spark.sql.functions._ val trimmedDF = originalDF.withColumn("SCHDULE", trim(col("SCHDULE")))
# Python from pyspark.sql.functions import trim, col trimmed_df = original_df.withColumn("SCHDULE", trim(col("SCHDULE")))
提取斜杠后的部分
// Scala import org.apache.spark.sql.functions._ val trimmedDF = originalDF.withColumn("SCHDULE", split(col("SCHDULE"), "/").getItem(1))
# Python from pyspark.sql.functions import split, col trimmed_df = original_df.withColumn("SCHDULE", split(col("SCHDULE"), "/").getItem(1))
自定义截取长度
比如保留前5个字符:
// Scala import org.apache.spark.sql.functions._ val trimmedDF = originalDF.withColumn("SCHDULE", substring(col("SCHDULE"), 1, 5))
# Python from pyspark.sql.functions import substring, col trimmed_df = original_df.withColumn("SCHDULE", substring(col("SCHDULE"), 1, 5))
如果你的修剪需求更具体,补充规则后可以再调整实现。
内容的提问来源于stack exchange,提问作者RMK
相关产品推荐
相关产品推荐

