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
相关产品推荐
相关产品推荐

