小编RSG*_*RSG的帖子

遍历 Spark 数据帧的列并更新指定的值

为了遍历从 Hive 表创建的 Spark Dataframe 的列并更新所有出现的所需列值,我尝试了以下代码。

import org.apache.spark.sql.{DataFrame}
import org.apache.spark.sql.functions._
import org.apache.spark.sql.functions.udf

val a: DataFrame = spark.sql(s"select * from default.table_a")

    val column_names: Array[String] = a.columns

    val required_columns: Array[String] = column_names.filter(name => name.endsWith("_date")) 

    val func = udf((value: String) => { if if (value == "XXXX" || value == "WWWW" || value == "TTTT") "NULL" else value } )

    val b = {for (column: String <- required_columns) { a.withColumn(column , func(a(column))) } a}
Run Code Online (Sandbox Code Playgroud)

在 spark shell 中执行代码时,出现以下错误。

scala> val b = {for …
Run Code Online (Sandbox Code Playgroud)

hive scala apache-spark apache-spark-sql

0
推荐指数
1
解决办法
4090
查看次数

标签 统计

apache-spark ×1

apache-spark-sql ×1

hive ×1

scala ×1