原文:towardsdatascience.com/pyspark-explained-dealing-with-invalid-records-when-reading-csv-and-json-files-4c671feb4d8e

如果你经常使用 PySpark,你将执行的最常见操作之一就是从外部文件中读取 CSV 或 JSON 数据到 DataFrame 中。如果你的输入数据与一个用户指定的模式定义相关联,你可能会发现你正在处理的记录中并非所有都符合模式规范。

换句话说,可能存在无效记录。也许某些字段会缺失,会有额外的未记录字段,或者某些字段包含的数据类型与模式中指定的类型不符。对于小文件,追踪这些记录并不是问题,但对于大数据世界中的大型文件,这些可能会成为一个真正的头疼问题。问题是,在这些情况下我们能做什么?

如你很快就会发现的,这个问题的答案之一是在将 CSV 或 JSON 文件读取到 DataFrame 时使用各种 PySpark parse 选项。

访问免费的 PySpark 开发环境

在我们继续之前,如果你想要跟随这篇文章中的代码,你需要访问 PySpark 开发环境。

如果你很幸运,可以通过工作、云服务或本地安装来访问 PySpark,那么请继续使用。如果没有,请查看下面的链接,其中我详细介绍了如何访问一个名为 Databricks Community Edition 的优秀的免费在线 PySpark 开发环境。

Databricks 是一个基于云的数据工程、机器学习和分析平台,围绕 Apache Spark 构建,为处理大数据工作负载提供了一个统一的环境。Databricks 的创始人创建了 Spark,因此他们了解自己的产品。

如何访问免费的在线 Spark 开发环境

样本数据

对于这篇文章,假设我们正在处理一个包含关于世界上人口最多的城市信息的 CSV 或 JSON 数据文件。该文件包含城市名称、所在国家、纬度和经度坐标,以及最后,2015 年的人口。

下面是一个包含有效记录和一些包含无效数据的记录的文件示例。

City, Country, Latitude,Longitude,Population
Tokyo,Japan,35.6895,139.69171,38001000
Delhi,India,28.66667,77.21667,25703000
Shanghai,China,31.22,121.46,23741000,ABCDE
São Paulo,Brazil,-23.55
Mumbai,India, India,72.880838,21043000

以下是上述 CSV 数据的 JSON 等价物。

{"City": "Tokyo", "Country": "Japan", "Latitude": 35.6895, "Longitude": 139.69171, "Population": 38001000}
{"City": "Delhi", "Country": "India", "Latitude": 28.66667, "Longitude": 77.21667, "Population": 25703000}
{"City": "Shanghai", "Country": "China", "Latitude": 31.22, "Longitude": 121.46, "Population": "23741000,ABCDE"}
{"City": "São Paulo", "Country": "Brazil", }
{"City": "Mumbai","Country": "India","Latitude": "India","Longitude": 72.880838,"Population": 21043000}

我在我的系统上创建了这两个文件,然后将它们上传到 Databricks 文件系统,分别命名为 cities.csv 和 cities.json。

记录 1 和 2 是有效的,记录 3 在末尾有一个虚假的额外文本值 “ABCDE”。记录 4 缺少经纬度字段的数据,而记录 5 在纬度值应该出现的地方重复了字符串 “India”。

显然,在这种情况下,格式不正确的记录很容易被发现,这些记录会在我们用 PySpark 处理数据之前很久就被处理掉。但考虑一下这些记录被埋藏在多兆字节 CSV 或 JSON 文件中间的情况。或者,数据处理可能是自动化管道的一部分,你并不一定事先能看到数据。

Spark 提供了什么可以帮助我们识别和处理此类数据的东西吗?当然,它提供了,否则这篇文章就没有什么意义了!

当读取 CSV 或 JSON 数据时,PySpark 可以使用三种不同的解析模式来帮助处理数据问题。这些是,

宽容

这是默认行为,告诉 Spark 在无法正确解析的字段中插入空值。如果任何输入字段被认为不符合定义的模式,所有字段都将设置为空值。当你想尽可能多地读取数据并在稍后阶段处理无效数据时,请使用此模式。

Dropmalformed

在这种模式下,Spark 将丢弃任何无法正确解析一个或多个字段的记录。当你只想读取完全符合你的模式定义的数据时,请使用此模式。

Failfast

在这种模式下,如果 Spark 无法正确解析任何输入数据,它将抛出一个异常,并且不会加载任何数据。当输入数据中的错误无法容忍时,请使用此模式。

让我们看看这些模式在实际操作中的例子。

PySpark 代码

当使用 PySpark 处理像“cities”文件这样的数据时,通常我们首先会设置与输入数据相关的模式。换句话说,我们向 PySpark 描述构成我们输入数据记录的字段名称和数据类型。

from pyspark.sql.types import StructField, StructType, DoubleType,StringType,LongType

citySchema = StructType(
  [

        StructField("city", StringType(), True),
        StructField("country", StringType(), True),
        StructField("latitude", DoubleType(), True),
        StructField("longitude", DoubleType(), True),
        StructField("population", LongType(), True),
  ]
)

在这里,我们说,我们的输入数据每一行应包含 5 个字段,前两个字段是文本,接下来的两个应该是双精度浮点数,最后一个是一个长整数。

使用宽容模式

我们可以使用 spark.read 函数来读取我们的数据文件。请注意,我们指定文件的格式为 CSV,因为解析模式适用于 CSV 或 JSON 格式文件。接下来,我们指出文件包含标题,否则 PySpark 将将文件的第一行视为输入数据。接下来,我们指定要使用的预定义模式,我们的字段分隔符,解析模式,最后是输入文件的存储位置。

df = spark.read.format("csv") 
    .option("header", "true").schema("citySchema") 
    .option("delimiter", ',').option("mode","PERMISSIVE")
    .load("dbfs:/FileStore/shared_uploads/test/cities.csv")

df.show()

+---------+-------+--------+---------+----------+
|     city|country|latitude|longitude|population|
+---------+-------+--------+---------+----------+
|    Tokyo|  Japan| 35.6895|139.69171|  38001000|
|    Delhi|  India|28.66667| 77.21667|  25703000|
| Shanghai|  China|   31.22|   121.46|  23741000|
|Sao Paulo| Brazil|  -23.55|     null|      null|
|     null|   null|    null|     null|      null|
+---------+-------+--------+---------+----------+

如您所见,第 3 行的额外字段被忽略,第 4 行缺失的值被设置为空。数据集的最后一行,因为我们遇到了一个期望双精度浮点数的字符串,被视为损坏的,因此所有字段都被设置为空值。

根据你的情况,记录三中的额外值可能是一个问题。不是 Spark 忽略了它——这一点是可以接受的,但最好能被告知该特定记录有些不太对劲。

幸运的是,有一种方法可以在这种情况下通过向 spark.read 函数添加一个额外的选项来获取信息,该选项称为 columnNameOfCorruptRecord 并将其值设置为 _corrupt_record。此外,我们还需要在我们的模式中添加一个额外的字段 _corrupt_data 来存储任何无效记录的值。

 from pyspark.sql.types import StructField, StructType, DoubleType,StringType,LongType

# Add  the extra field _corrupt_record
# to our original schema definition
#
citySchema = StructType([
    StructField("City", StringType(), True),
    StructField("Country", StringType(), True),
    StructField("Latitude", FloatType(), True),
    StructField("Longitude", FloatType(), True),
    StructField("Population", IntegerType(), True),
    StructField("_corrupt_record", StringType(), True)
])

#Now see what happens when we read the original file again
# after adding the columnNameOfCorruptRecord option
#
df = spark.read.format("csv") 
    .option("header", "true").schema("citySchema") 
    .option("delimiter", ',').option("mode","PERMISSIVE") 
    .option("columnNameOfCorruptRecord", "_corrupt_record") 
    .load("dbfs:/FileStore/shared_uploads/test/cities.csv")

df.show()

+---------+-------+--------+---------+----------+------------------------------------------+
|City     |Country|Latitude|Longitude|Population|_corrupt_record                           |
+---------+-------+--------+---------+----------+------------------------------------------+
|Tokyo    |Japan  |35.6895 |139.69171|38001000  |null                                      |
|Delhi    |India  |28.66667|77.21667 |25703000  |null                                      |
|Shanghai |China  |31.22   |121.46   |23741000  |Shanghai,China,31.22,121.46,23741000,ABCDE|
|São Paulo|Brazil |-23.55  |null     |null      |São Paulo,Brazil,-23.55                   |
|Mumbai   |India  |null    |72.88084 |21043000  |Mumbai,India, India,72.880838,21043000    |
+---------+-------+--------+---------+----------+------------------------------------------+

正如你所见,我们在 DataFrame 中新增了一个完全不同的字段 _corrupt_record。如果输入记录 有效, 则该字段设置为 null,否则它包含整个输入记录。在上海记录的情况下,它捕捉到输入记录中存在未被模式定义所考虑的额外信息,因此将其标记为可疑。这在后续的数据处理流程步骤中可能非常有用。

使用 Dropmalformed

对于这个例子,我们将使用基于城市数据的 JSON 格式文件。我们可以恢复到我们原来的模式,因为我们这次不需要 corruptrecord 选项。记住,dropmalformed 处理有效记录。

from pyspark.sql.types import StructField, StructType, DoubleType,StringType,LongType

# Define the schema
citySchema = StructType([
    StructField("City", StringType(), True),
    StructField("Country", StringType(), True),
    StructField("Latitude", FloatType(), True),
    StructField("Longitude", FloatType(), True),
    StructField("Population", IntegerType(), True)
])
df = spark.read.format("json") 
    .schema(citySchema) 
    .option("mode", "DROPMALFORMED") 
    .load("dbfs:/FileStore/shared_uploads/test/cities.json")

# Show the DataFrame
df.show(truncate=False)

+-----+-------+--------+---------+----------+
|City |Country|Latitude|Longitude|Population|
+-----+-------+--------+---------+----------+
|Tokyo|Japan  |35.6895 |139.69171|38001000  |
|Delhi|India  |28.66667|77.21667 |25703000  |
+-----+-------+--------+---------+----------+

好的,输出结果是我们预期的。只有我们输入数据文件中的两个有效记录被加载到 DataFrame 中。

使用 Failfast

这种解析模式如果任何输入记录未通过模式验证,则应抛出异常。让我们看看这是否是情况。切换回我们的 CSV 数据文件,我们有 …

from pyspark.sql.types import StructField, StructType, DoubleType,StringType,LongType

# Define the schema
citySchema = StructType([
    StructField("City", StringType(), True),
    StructField("Country", StringType(), True),
    StructField("Latitude", FloatType(), True),
    StructField("Longitude", FloatType(), True),
    StructField("Population", IntegerType(), True)
])

df = spark.read.format("csv") 
    .option("header", "true") 
    .schema("citySchema") 
    .option("delimiter", ',') 
    .option("mode","FAILFAST") 
    .load("dbfs:/FileStore/shared_uploads/test/cities.csv")

df.show()

...
...
Caused by: org.apache.spark.SparkException: Malformed records are detected in record parsing. Parse Mode: FAILFAST. To process malformed records as null result, try setting the option 'mode' as 'PERMISSIVE'.
 at org.apache.spark.sql.errors.QueryExecutionErrors$.malformedRecordsDetectedInRecordParsingError(QueryExecutionErrors.scala:1936)
 at org.apache.spark.sql.catalyst.util.FailureSafeParser.parse(FailureSafeParser.scala:103)
 at org.apache.spark.sql.execution.datasources.json.TextInputJsonDataSource$.$anonfun$readFile$5(JsonDataSource.scala:215)
 at scala.collection.Iterator$$anon$11.nextCur(Iterator.scala:486)
 at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:492)
 at org.apache.spark.util.CompletionIterator.hasNext(CompletionIterator.scala:31)
 at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
 at org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1$$anon$2.getNext(FileScanRDD.scala:608)
 ... 31 more
Caused by: org.apache.spark.sql.catalyst.util.BadRecordException: org.apache.spark.SparkRuntimeException: [CANNOT_PARSE_JSON_FIELD] Cannot parse the field name 'Population' and the value 23741000,ABCDE of the JSON token type VALUE_STRING to target Spark data type "INT".
Caused by: org.apache.spark.SparkRuntimeException: [CANNOT_PARSE_JSON_FIELD] Cannot parse the field name 'Population' and the value 23741000,ABCDE of the JSON token type VALUE_STRING to target Spark data type "INT".

...
... 

为了证明它按预期工作,让我们修复我们输入 CSV 文件中的验证错误,并再次尝试 Failfast。

City, Country, Latitude,Longitude,Population
Tokyo,Japan,35.6895,139.69171,38001000
Delhi,India,28.66667,77.21667,25703000
Shanghai,China,31.22,121.46,23741000
São Paulo,Brazil,-23.55,-46.6396,12330000
Mumbai,India,19.0760,72.880838,21043000

并且我们的输出,运行与之前相同的 FAILFAST 代码,是 …

+---------+-------+--------+---------+----------+
|City     |Country|Latitude|Longitude|Population|
+---------+-------+--------+---------+----------+
|Tokyo    |Japan  |35.6895 |139.69171|38001000  |
|Delhi    |India  |28.66667|77.21667 |25703000  |
|Shanghai |China  |31.22   |121.46   |23741000  |
|São Paulo|Brazil |-23.55  |-46.6396 |12330000  |
|Mumbai   |India  |19.076  |72.88084 |21043000  |
+---------+-------+--------+---------+----------+

摘要

在这篇文章中,我向你展示了如何使用 PySpark 读取 JSON 或 CSV 文件时的解析模式选项。三种不同的模式,宽松的、丢弃错误的Failfast,为你提供了处理无效数据的替代方法。从加载所有内容并在之后担心无效数据,到遇到任何无效数据时抛出异常,PySpark 让你对处理 CSV 或 JSON 数据文件时的意外情况有完全的控制权。

_ 好的,这就是我现在要说的全部内容。希望你觉得这篇文章有用。如果你觉得有用,请通过这个链接查看我的个人资料页面。从那里,你可以看到我其他发布的文章,并订阅以获取我发布新内容的通知。_

如果你喜欢这个内容,我认为你还会对以下这些文章感兴趣。

PySpark 解释:explode 和 collect_list 函数

Luma Labs AI:首次了解他们的文本/图像到视频产品

更多推荐