You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark/Scala中拆分CSV的JSON列并基于JSON键生成新列

在Spark 2.x(Scala)中拆分Desc列的JSON数据并生成新列

嘿,我来帮你搞定这个需求!针对你给出的DataFrame,我们可以用Spark 2.x里的内置函数轻松把Desc列里的JSON字符串拆成单独的列,下面给你两种实用的方法,按需选择:

方法一:使用from_json(推荐,适合结构化JSON)

这种方法需要先定义JSON对应的Schema,适合JSON结构固定、后续可能扩展的场景,性能也更优。

完整代码示例

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.types.{StructType, StructField, IntegerType}
import org.apache.spark.sql.functions.from_json

object JsonColumnParser {
  def main(args: Array[String]): Unit = {
    // 初始化SparkSession
    val spark = SparkSession.builder()
      .appName("JsonSplitDemo")
      .master("local[*]") // 本地测试用,生产环境请移除这个配置
      .getOrCreate()
    
    import spark.implicits._

    // 创建和你示例一致的DataFrame
    val sourceDF = Seq(
      (201, "MIS20", """{ "Total": 200,"Defective": 21 }"""),
      (202, "MIS30", """{ "Total": 740,"Defective": 58 }""")
    ).toDF("id", "Category", "Desc")

    // 定义JSON数据对应的Schema,和Desc里的键对应
    val jsonSchema = new StructType()
      .add("Total", IntegerType)
      .add("Defective", IntegerType)

    // 解析JSON列并提取字段
    val resultDF = sourceDF
      // 把Desc的JSON字符串转换成Struct类型的列
      .withColumn("parsed_json", from_json($"Desc", jsonSchema))
      // 选择原字段+解析后的JSON字段,并重命名
      .select(
        $"id",
        $"Category",
        $"parsed_json.Total".as("Total"),
        $"parsed_json.Defective".as("Defective")
      )
      // 可选:移除中间生成的parsed_json列
      .drop("parsed_json")

    // 打印结果
    resultDF.show()
  }
}

方法二:使用get_json_object(快速提取,无需Schema)

如果只是简单提取几个字段,不想定义Schema,用这个方法更快捷。它通过JSONPath表达式直接提取指定键的值,不过需要手动转换数据类型。

完整代码示例

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.get_json_object

object QuickJsonSplit {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder()
      .appName("QuickJsonSplit")
      .master("local[*]")
      .getOrCreate()
    
    import spark.implicits._

    val sourceDF = Seq(
      (201, "MIS20", """{ "Total": 200,"Defective": 21 }"""),
      (202, "MIS30", """{ "Total": 740,"Defective": 58 }""")
    ).toDF("id", "Category", "Desc")

    val resultDF = sourceDF
      // 提取Total字段,注意JSONPath是$.Total,然后转成Integer类型
      .withColumn("Total", get_json_object($"Desc", "$.Total").cast(IntegerType))
      // 提取Defective字段
      .withColumn("Defective", get_json_object($"Desc", "$.Defective").cast(IntegerType))
      // 可选:移除原Desc列
      .drop("Desc")

    resultDF.show()
  }
}

最终结果

两种方法都会得到如下的DataFrame:

+---+--------+-----+---------+
| id|Category|Total|Defective|
+---+--------+-----+---------+
|201|   MIS20|  200|       21|
|202|   MIS30|  740|       58|
+---+--------+-----+---------+

注意事项

  • 确保Desc列的JSON格式是合法的,如果有格式错误,解析后会得到null值,你可以用isnull或filter来处理脏数据
  • from_json是Spark 2.1及以上版本才支持的,如果你的Spark版本更早,建议升级或者使用get_json_object
  • 如果JSON里有嵌套结构,from_json的Schema可以嵌套定义(比如StructType里再包含StructType),灵活性很强

内容的提问来源于stack exchange,提问作者Prashant

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.21 04:24:55