如何在PySpark中为字符串类型嵌套JSON列创建Schema
PySpark 解析字符串类型嵌套JSON列的Schema定义与处理方法
第一步:手动定义匹配结构的Schema
根据给出的JSON样例,直接用PySpark内置的Struct类型构造对应Schema即可,所有字段默认设置为可空适配null值场景,其中ArrayType为变长数组类型,天然适配列内任意数量的array_element节点,无需提前固定数组长度:
from pyspark.sql.types import ( StructType, StructField, StringType, ArrayType ) # 定义telephone字段下嵌套的联系方式结构 telephone_element_schema = StructType([ StructField("locationtype", StringType(), nullable=True), StructField("countrycode", StringType(), nullable=True), StructField("phonenumber", StringType(), nullable=True), StructField("phonetechtype", StringType(), nullable=True), StructField("countryaccesscode", StringType(), nullable=True), StructField("phoneremark", StringType(), nullable=True) ]) # 定义整列JSON对应的完整Schema json_col_schema = StructType([ StructField("addressline", ArrayType( StructType([ StructField("array_element", StringType(), nullable=True) ]) ), nullable=True), StructField("telephone", ArrayType( StructType([ StructField("array_element", telephone_element_schema, nullable=True) ]) ), nullable=True) ])
第二步:用定义好的Schema解析字符串列
Parquet文件读入后这两列是String类型,直接用from_json函数转换即可,转换后就可以直接用.语法访问嵌套字段,也支持对数组做explode等常规Spark操作:
from pyspark.sql.functions import from_json, col # 替换下方的col1、col2为你Parquet中实际存储JSON字符串的列名 parsed_df = df.withColumn( "parsed_col1", from_json(col("col1"), json_col_schema) ).withColumn( "parsed_col2", from_json(col("col2"), json_col_schema) )
常用操作示例
解析完成后可直接对嵌套结构做提取、展开等处理:
- 提取地址数组第一个元素的内容:
parsed_df.select(col("parsed_col1.addressline")[0].array_element) - 展开所有telephone节点提取手机号:
from pyspark.sql.functions import explode parsed_df.select( explode(col("parsed_col1.telephone")).alias("tel_item") ).select(col("tel_item.array_element.phonenumber"))
注意:如果JSON字符串存在格式不合法的脏数据,
from_json解析后对应行的结果会返回null,不会直接抛错中断任务,可后续通过过滤null值定位脏数据。
内容的提问来源于stack exchange,提问作者naga satish
相关产品推荐
相关产品推荐

