Apache Spark 2.1.1 Scala:无Schema解析JSON字符串列生成新DataFrame
在Spark 2.1.1中解析无Schema的JSON字符串列为多列
在Spark 2.1.1里处理这种没有预定义Schema的JSON列解析,核心思路是先从现有数据中自动推断出JSON的结构(Schema),再用这个Schema去解析每行的字符串,最后把解析后的结构化数据展开成单独的列。下面是具体的实现步骤,我分别给出Scala和Python版本的代码:
步骤说明
- 准备原始DataFrame:先构造出你示例中的输入DataFrame(实际场景中替换成你的数据源即可)。
- 自动推断JSON Schema:把DataFrame里的JSON字符串提取出来,用Spark的JSON读取器自动推断出Schema——这个读取器会扫描所有JSON数据,合并所有字段类型生成对应的Schema。
- 解析JSON字符串:用
from_json函数,结合推断出的Schema,把JSON字符串列转换成结构化的StructType列。 - 展开结构化列:把
StructType列里的每个字段提取成单独的列,得到最终的结果。
Scala 实现代码
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ // 初始化SparkSession val spark = SparkSession.builder() .appName("ParseUnschemaedJSON") .master("local[*]") // 生产环境请去掉master配置 .getOrCreate() import spark.implicits._ // 1. 构造原始DataFrame(替换成你的实际数据源) val rawDF = Seq( """{"a":2,"b":"hello"}""", """{"a":1,"b":"hi"}""" ).toDF("json_string") // 2. 从现有JSON数据中推断Schema val jsonSchema = spark.read.json(rawDF.select("json_string").as[String].rdd).schema // 3. 解析JSON字符串为结构化列 val parsedDF = rawDF.withColumn("parsed_json", from_json(col("json_string"), jsonSchema)) // 4. 展开结构化列到单独的列 val finalDF = parsedDF.select("parsed_json.*") // 查看结果 finalDF.show()
Python 实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import from_json # 初始化SparkSession spark = SparkSession.builder \ .appName("ParseUnschemaedJSON") \ .master("local[*]") # 生产环境请去掉master配置 .getOrCreate() # 1. 构造原始DataFrame(替换成你的实际数据源) raw_df = spark.createDataFrame([ ("""{"a":2,"b":"hello"}""",), ("""{"a":1,"b":"hi"}""",) ], ["json_string"]) # 2. 从现有JSON数据中推断Schema json_schema = spark.read.json(raw_df.select("json_string").rdd.map(lambda row: row[0])).schema # 3. 解析JSON字符串为结构化列 parsed_df = raw_df.withColumn("parsed_json", from_json(raw_df["json_string"], json_schema)) # 4. 展开结构化列到单独的列 final_df = parsed_df.select("parsed_json.*") # 查看结果 final_df.show()
注意事项
- 如果你的JSON数据存在字段不一致的情况(比如有的行有额外字段,有的行缺失某些字段),自动推断的Schema会包含所有出现过的字段,缺失字段的行对应位置会显示
null。 - Spark 2.1.1的
from_json函数必须依赖明确的Schema,所以推断Schema这一步是必不可少的——这也是它和高版本Spark(比如2.3+支持无Schema直接解析)的核心区别。 - 生产环境中如果数据量极大,推断Schema的过程可能会消耗较多资源,你可以考虑先采样部分数据来推断Schema,能有效提升效率。
内容的提问来源于stack exchange,提问作者Clairton Menezes
相关产品推荐
相关产品推荐

