使用PySpark读取JSON返回全空DataFrame的问题求助
PySpark读取JSON全为Null问题排查与解决
有如下结构的JSON文件,尝试用PySpark读取并指定自定义Schema后,DataFrame所有字段均为Null:
[{'id': '34556', 'InsuranceProvider': 'sdcsdf', 'Type': {'Client': {'PaidIn': {'Insuranceid': '442211', 'Insurancedesc': 'sdfsdf vdsfs', 'purchaseditems': [{'InsuranceNumber': '1', 'InsuranceLabel': 'SDF', 'Insurancequantity': 1, 'Insuranceprice': 234, 'discountsreceived': [{'amount': 120, 'description': 'Item 1, Discount 1'}], 'childItems': [{'InsuranceNumber': '1', 'InsuranceLabel': 'CSFGG', 'Insurancequantity': 1, 'Insuranceprice': 0, 'discountsreceived': [{'amount': 452, 'description': 'Insurance item 1, Discount 1'}] }]}]}}} 'eventTime': '2022-05-19T01:59:10.379Z' }]使用的Schema定义:
discount_type = StructType([StructField("amount", IntegerType(), True), StructField("description", StringType(), True)]) child_item_type = StructType([StructField("InsuranceNumber", StringType(), True), StructField("InsuranceLabel", StringType(), True), StructField("Insurancequantity", IntegerType(), True), StructField("Insuranceprice", IntegerType(), True), StructField("discountsreceived", ArrayType(discount_type) , True), ]) item_type = StructType([StructField("InsuranceNumber", StringType(), True), StructField("InsuranceLabel", StringType(), True), StructField("Insurancequantity", IntegerType(), True), StructField("Insuranceprice", IntegerType(), True), StructField("discountsreceived", ArrayType(discount_type), True), StructField("childItems",ArrayType(child_item_type) , True), ]) order_paid_type = StructType([StructField("Insuranceid", StringType(), True), StructField("Insurancedesc", StringType(), True), StructField("purchaseditems", ArrayType(item_type), True), ]) message_type = StructType([StructField("PaidIn", order_paid_type, True)]) data_type = StructType([StructField("Client", message_type, True)]) body_type = StructType([StructField("id", StringType(), True), StructField("InsuranceProvider", StringType(), True), StructField("Type", data_type, True), StructField("eventTime", StringType(), True), ])读取代码:
data = spark.read.schema(schema).json(file_path) data.show()输出结果:
+----+-----------------+----+---------+ | id|InsuranceProvider|Type|eventTime| +----+-----------------+----+---------+ |null| null|null| null| |null| null|null| null| |null| null|null| null| +----+-----------------+----+---------+
问题1:JSON格式不符合标准规范
Spark的JSON读取器严格遵循RFC 7159标准JSON格式,提供的JSON存在两处语法错误:
- 使用单引号
',标准JSON要求所有键和字符串值必须用双引号" Type对象结束后,eventTime键前缺少逗号,导致JSON结构不完整
修正后的合法JSON:
[{"id": "34556", "InsuranceProvider": "sdcsdf", "Type": {"Client": {"PaidIn": {"Insuranceid": "442211", "Insurancedesc": "sdfsdf vdsfs", "purchaseditems": [{"InsuranceNumber": "1", "InsuranceLabel": "SDF", "Insurancequantity": 1, "Insuranceprice": 234, "discountsreceived": [{"amount": 120, "description": "Item 1, Discount 1"}], "childItems": [{"InsuranceNumber": "1", "InsuranceLabel": "CSFGG", "Insurancequantity": 1, "Insuranceprice": 0, "discountsreceived": [{"amount": 452, "description": "Insurance item 1, Discount 1"}] }]}]}}, "eventTime": "2022-05-19T01:59:10.379Z" }]
问题2:Schema变量名不匹配
定义的Schema变量是body_type,但读取代码中使用的是未定义的schema,变量名不一致导致Spark无法识别正确的Schema,需修正读取代码:
data = spark.read.schema(body_type).json(file_path) data.show()
验证修正结果
完成以上两处修正后,重新读取JSON,DataFrame将正确解析所有字段:
+------+-----------------+--------------------+--------------------+ | id|InsuranceProvider| Type| eventTime| +------+-----------------+--------------------+--------------------+ |34556 | sdcsdf|{{{442211, sdfsdf...|2022-05-19T01:59:...| +------+-----------------+--------------------+--------------------+
内容的提问来源于stack exchange,提问作者Aman Mishra
相关产品推荐
相关产品推荐

