【问题标题】:Spark - Permissive mode with JSON file moves all records to corrupt columnSpark - 使用 JSON 文件的许可模式将所有记录移动到损坏的列
【发布时间】:2020-10-22 16:57:13
【问题描述】:

我正在尝试使用 spark 摄取 JSON 文件。我正在手动应用架构来创建数据框。问题甚至是模式的单个记录不匹配它会将整个文件(所有记录)移动到损坏的列?

数据

[{
  "RecordNumber": 2,
  "Zipcode": 704,
  "ZipCodeType": "STANDARD",
  "City": "PASEO COSTA DEL SUR",
  "State": "PR"
},
{
  "RecordNumber": 10,
  "Zipcode": 709,
  "ZipCodeType": "STANDARD",
  "City": "BDA SAN LUIS",
  "State": "PR"
},{
  "Zipcode": "709aa",
  "ZipCodeType": "STANDARD",
  "City": "BDA SAN LUIS",
  "State": "PR"
}]

代码

import org.apache.spark.sql.types._
 import org.apache.spark.sql.types.DataTypes._
 val s = StructType(StructField("City",StringType,true) ::
 StructField("RecordNumber",LongType,true) ::
 StructField("State",StringType,true) ::
 StructField("ZipCodeType",StringType,true) ::
 StructField("Zipcode",LongType,true) ::
 StructField("corrupted_record",StringType,true) ::
 Nil)
 val df2=spark.read.
option("multiline","true").
option("mode", "PERMISSIVE").
option("columnNameOfCorruptRecord", "corrupted_record").
schema(s).
json("/tmp/test.json")
df2.show(false)

输出

scala> df2.filter($"corrupted_record".isNotNull).show(false)
+----+------------+-----+-----------+-------+------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
|City|RecordNumber|State|ZipCodeType|Zipcode|corrupted_record                                                                                                                                                                                                                                                                                                                                |
+----+------------+-----+-----------+-------+------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
|null|null        |null |null       |null   |[{
  "RecordNumber": 2,
  "Zipcode": 704,
  "ZipCodeType": "STANDARD",
  "City": "PASEO COSTA DEL SUR",
  "State": "PR"
},
{
  "RecordNumber": 10,
  "Zipcode": 709,
  "ZipCodeType": "STANDARD",
  "City": "BDA SAN LUIS",
  "State": "PR"
},{
  "Zipcode": "709aa",
  "ZipCodeType": "STANDARD",
  "City": "BDA SAN LUIS",
  "State": "PR"
}]

问题 因为只有第三条记录在String 中有邮政编码,而我希望它是integer ("Zipcode": "709aa",) 不应该只有第三条记录转到corrupted_record 列,其他记录应该被正确解析?

【问题讨论】:

    标签: apache-spark pyspark apache-spark-sql


    【解决方案1】:

    您只有一条记录(因为多行,true)已损坏,所以一切都在那里。

    就像 documentation 中所说的,如果您希望 spark 单独处理记录,则需要使用 Json lines format,这对于更大的文件也将更具可扩展性,因为 spark 将能够在多个执行器中分发解析。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2011-08-23
      • 2021-05-09
      • 2019-12-18
      • 2020-02-26
      • 1970-01-01
      • 2020-03-13
      • 2014-12-05
      相关资源
      最近更新 更多