Spark SQL中用Scala函数处理DataFrame列:生成/替换列方法
嘿,我来帮你搞定这两个Spark SQL的需求!先从你的示例DataFrame入手,咱们一步步用Scala实现。
首先,先创建你的示例DataFrame,方便后续测试:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ // 初始化SparkSession val spark = SparkSession.builder().appName("SparkSQLDemo").master("local[*]").getOrCreate() import spark.implicits._ // 构建示例DataFrame val df = Seq( (1, "Mary", "ABCD"), (2, "Joey", "DOGE"), (3, "Lane", "POOP"), (4, "Jack", "MEGA"), (5, "Lynn", "ARGH") ).toDF("id", "name", "data")
需求1:使用Scala函数生成新列(startsWithA)
我们需要新增一个布尔类型列startsWithA,标识data列是否以"A"开头。这里可以用**UDF(用户自定义函数)**来封装Scala的字符串判断逻辑,再通过withColumn方法添加新列。
// 定义Scala函数:判断字符串是否以"A"开头 def startsWithA(s: String): Boolean = s.startsWith("A") // 将Scala函数注册为Spark UDF val startsWithAUdf = udf(startsWithA _) // 添加新列startsWithA val dfWithNewCol = df.withColumn("startsWithA", startsWithAUdf(col("data"))) // 查看结果 dfWithNewCol.show()
执行后你会得到这样的结果:
+---+----+----+-----------+ | id|name|data|startsWithA| +---+----+----+-----------+ | 1|Mary|ABCD| true| | 2|Joey|DOGE| false| | 3|Lane|POOP| false| | 4|Jack|MEGA| false| | 5|Lynn|ARGH| true| +---+----+----+-----------+
需求2:使用Scala函数替换原有列
withColumn方法有个实用特性:如果指定的列名已经存在,就会直接替换原有列的数据。我们可以利用这个特性,用Scala函数处理后覆盖原列。
举个例子,比如把data列替换为原字符串的反转:
// 定义Scala函数:反转字符串 def reverseString(s: String): String = s.reverse // 注册UDF val reverseStringUdf = udf(reverseString _) // 替换原有data列 val dfWithReplacedCol = df.withColumn("data", reverseStringUdf(col("data"))) // 查看结果 dfWithReplacedCol.show()
执行结果如下:
+---+----+----+ | id|name|data| +---+----+----+ | 1|Mary|DCBA| | 2|Joey|EGOD| | 3|Lane|POOP| | 4|Jack|AGEM| | 5|Lynn|HGRA| +---+----+----+
如果需要同时替换多列,只需要连续调用withColumn即可,比如同时替换name列为名字长度、data列为是否以"A"开头:
// 定义两个Scala函数:获取名字长度、判断是否以"A"开头 def nameLength(name: String): Int = name.length def startsWithA(s: String): Boolean = s.startsWith("A") // 注册UDF val nameLengthUdf = udf(nameLength _) val startsWithAUdf = udf(startsWithA _) // 同时替换name和data列 val dfWithReplacedCols = df .withColumn("name", nameLengthUdf(col("name"))) .withColumn("data", startsWithAUdf(col("data"))) dfWithReplacedCols.show()
结果会变成:
+---+----+-----+ | id|name| data| +---+----+-----+ | 1| 4| true| | 2| 4|false| | 3| 4|false| | 4| 4|false| | 5| 4| true| +---+----+-----+
内容的提问来源于stack exchange,提问作者user3685285
相关产品推荐
相关产品推荐

