Spark用自定义Schema读JSON返回全Null及TB级文件处理咨询
问题描述
我是EMR/HDFS/Hive/Spark领域新手,目前要处理一批单文件超50GB的大型JSON文件,目标是加载后查询特定键。这些JSON遵循统一标准,但不是所有文件都包含全部字段。我自定义了覆盖完整标准的Schema并生成DataFrame,但查询时所有字段都返回Null。现在正在EMR集群上用Hive/Spark做POC,同时想找处理TB级JSON文件的最优方案。
简化示例:
- 标准字段:
reporting_entity_name(字符串)、reporting_entity_type(字符串)、version(字符串) - 文件A包含
reporting_entity_name和version,无reporting_entity_type - 文件B包含
reporting_entity_name和reporting_entity_type,无version - 执行查询:
spark.sqlContext.sql("SELECT reporting_entity_name,reporting_entity_type,version FROM vw ").show()
期望结果:文件A返回对应字段,缺失的reporting_entity_type为Null;文件B同理。但实际查询结果全是Null。
实际场景:
JSON遵循复杂的医保价格透明标准,包含多层嵌套的数组和结构体。自定义Schema如下(原代码存在类型定义错误):
val schema = StructType(List( StructField("reporting_entity_name", StringType), StructField("reporting_entity_type", StringType), StructField("plan_name", StringType), StructField("plan_id_type", StringType), StructField("plan_id", StringType), StructField("plan_market_type", StringType), StructField("in_network", ArrayType(StructType(List( StructField("negotiation_arrangement", StringType), StructField("plan_name", StringType), StructField("billing_code_type", StringType), StructField("billing_code_type_version", StringType), StructField("billing_code", StringType), StructField("description", StringType), StructField("negotiated_rates", ArrayType(StructType(Array( StructField("negotiated_prices", ArrayType(StructType(Array( StructField("negotiated_type", StringType), StructField("negotiated_rate", StringType), StructField("expiration_date", StringType), StructField("service_code", ArrayType(StringType)), StructField("billing_class", StringType), StructField("billing_code_modifier", ArrayType(StringType)), StructField("additional_information", StringType)))), true), StructField("provider_groups", ArrayType(StructType(Array( StructField("npi", ArrayType(StringType)), StructField("tin", StringType))))), StructField("provider_references", ArrayType(StructType(Array( StructField("provider_group_id", StringType), StructField("provider_groups", ArrayType(StructType(Array( StructField("npi", ArrayType(StringType)), StructField("tin", StringType))))), StructField("location", StringType))))))))), StructField("bundled_codes", ArrayType(StructType(Array( StructField("billing_code_type", StringType), StructField("billing_code_type_version", StringType), StructField("billing_code", StringType), StructField("description", StringType)))), true), StructField("covered_services", ArrayType(StructType(Array( StructField("billing_code_type", StringType), StructField("billing_code_type_version", StringType), StructField("billing_code", StringType), StructField("description", StringType)))), true), ))), true ), StructField("provider_references",ArrayType(StructType(Array( StructField("provider_group_id",StringType), StructField("provider_groups",ArrayType(StructType(Array( StructField("npi",ArrayType(StringType),true), StructField("tin",StringType)))),true), StructField("location",StringType)))),true), StructField("last_updated_on",StringType), StructField("version",StringType) ))
样本为Gzip压缩的大型JSON文件。
问题排查与解决方案
1. 解决全Null查询问题
核心错误:Schema类型定义错误
原Schema中大量使用ArrayType(StructType(Array(StructField(...)))),但Spark的StructType构造参数必须是List[StructField],而非Array[StructField]。类型不匹配会导致Spark无法正确解析JSON,直接返回全Null。
修正方法:把所有StructType(Array(...))改成StructType(List(...)),同时给每个StructField显式设置nullable=true(缺失字段自动填充Null)。
其他可能原因排查
- JSON格式问题:如果文件是单个大JSON数组(多行),需要在读取时添加配置:
Spark默认解析单行JSON(每个对象占一行),如果是多行格式不配置会解析失败。spark.read.option("multiLine", "true").schema(schema).json("hdfs://path/to/files") - 权限问题:确认EMR集群有HDFS/S3文件的读取权限,无权限也可能返回空数据。
- 字段名匹配问题:检查JSON中的字段名和Schema中的字段名是否完全一致(大小写敏感)。
2. TB级JSON文件处理优化方案
存储格式转换
将JSON转换为Parquet/ORC列式存储:
- 列式存储只读取查询需要的字段,大幅减少IO
- 压缩比远高于JSON,节省存储成本
- 查询速度提升5-10倍
转换代码示例:
spark.read.schema(schema).json("hdfs://path/to/json").write.mode("overwrite").parquet("hdfs://path/to/parquet")
文件拆分与预处理
- Spark自动拆分:Spark会根据文件大小和集群配置自动拆分文件,避免单文件过大导致OOM
- 手动拆分:如果单文件超过100GB,可在HDFS上用
split命令拆分,或者用Hadoop的TextInputFormat按行拆分
EMR集群优化
- 使用EBS优化实例,提升磁盘IO性能
- 调整Spark参数:
spark.sql.shuffle.partitions:设置为Core节点vCPU数的2-3倍(默认200,TB级数据建议设为1000-2000)spark.driver.memory/spark.executor.memory:根据实例规格调整,避免OOM- 开启动态资源分配:
spark.dynamicAllocation.enabled=true
分区与分桶
- 分区:按高频过滤字段(如
last_updated_on、plan_market_type)分区,减少查询时扫描的数据量 - 分桶:按高频查询字段(如
plan_id)分桶,提升聚合和join性能
增量处理
如果是定期更新的文件,只处理新增文件,避免重复处理全量数据:
val newFiles = spark.read.json("hdfs://path/to/new/json") val existingData = spark.read.parquet("hdfs://path/to/parquet") newFiles.union(existingData).write.mode("append").parquet("hdfs://path/to/parquet")
修正后的Schema片段示例
val correctedSchema = StructType(List( StructField("reporting_entity_name", StringType, true), StructField("reporting_entity_type", StringType, true), StructField("plan_name", StringType, true), StructField("in_network", ArrayType(StructType(List( StructField("negotiation_arrangement", StringType, true), StructField("billing_code", StringType, true), // 修正嵌套结构:用List替代Array StructField("negotiated_rates", ArrayType(StructType(List( StructField("negotiated_prices", ArrayType(StructType(List( StructField("negotiated_type", StringType, true), StructField("negotiated_rate", StringType, true) ))), true), StructField("provider_groups", ArrayType(StructType(List( StructField("npi", ArrayType(StringType), true), StructField("tin", StringType, true) ))), true) ))), true) ))), true), StructField("last_updated_on", StringType, true), StructField("version", StringType, true) ))
内容的提问来源于stack exchange,提问作者The Crusher
相关产品推荐
相关产品推荐

