您可以使用 from_json 方法将您的 json 列转换为 structtype 列。然后,您可以根据您的情况将此列分成不同的列。但是,您必须记住,json 应该具有统一的格式,否则结果可能不理想。
可以参考以下代码:
val df = spark.createDataFrame(Seq(
("A", "B", "{\"Name\":\"xyz\",\"Address\":\"NYC\",\"title\":\"engg\"}"),
("C", "D", "{\"Name\":\"mnp\",\"Address\":\"MIC\",\"title\":\"data\"}"),
("E", "F", "{\"Name\":\"pqr\",\"Address\":\"MNN\",\"title\":\"bi\"}")
)).toDF("col_1", "col_2", "col_json")
输入数据框如下:
scala> df.show(false)
+-----+-----+---------------------------------------------+
|col_1|col_2|col_json |
+-----+-----+---------------------------------------------+
|A |B |{"Name":"xyz","Address":"NYC","title":"engg"}|
|C |D |{"Name":"mnp","Address":"MIC","title":"data"}|
|E |F |{"Name":"pqr","Address":"MNN","title":"bi"} |
+-----+-----+---------------------------------------------+
现在,我们将找出 json 列的架构 col_schema,以便将其应用于 col_json 列
val col_schema = spark.read.json(df.select(col("col_json")).as[String]).schema
val outputDF = df.withColumn("new_col", from_json(col("col_json"), col_schema)).select("col_1", "col_2", "new_col.*")
结果如下:
scala> outputDF.show(false)
+-----+-----+-------+----+-----+
|col_1|col_2|Address|Name|title|
+-----+-----+-------+----+-----+
|A |B |NYC |xyz |engg |
|C |D |MIC |mnp |data |
|E |F |MNN |pqr |bi |
+-----+-----+-------+----+-----+
对我有用的 scala 代码是:
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.{col, from_json}
import scala.collection.Seq
object Sample{
def main(args: Array[String]): Unit = {
val spark = SparkSession.builder().master("local[*]").getOrCreate()
val df = spark.createDataFrame(Seq(
("A", "B", "{\"Name\":\"xyz\",\"Address\":\"NYC\",\"title\":\"engg\"}"),
("C", "D", "{\"Name\":\"mnp\",\"Address\":\"MIC\",\"title\":\"data\"}"),
("E", "F", "{\"Name\":\"pqr\",\"Address\":\"MNN\",\"title\":\"bi\"}")
)).toDF("col_1", "col_2", "col_json")
import spark.implicits._
val col_schema = spark.read.json(df.select("col_json").as[String]).schema
val outputDF = df.withColumn("new_col", from_json(col("col_json"), col_schema)).select("col_1", "col_2", "new_col.*")
outputDF.show(false)
}
}