【发布时间】:2018-04-15 08:51:57
【问题描述】:
所以我有一个问题,即在写入分区 Parquet 文件时,DataFrame 中的某些行会被丢弃。
这是我的步骤:
- 使用指定架构从 S3 读取 CSV 数据文件
- 按“日期”列 (DateType) 分区
- 写成 Parquet 和
mode=append
阅读的第一步按预期工作,没有解析问题。对于质量检查,我执行以下操作:
对于date='2012-11-22' 的特定分区,对 CSV 文件、加载的 DataFrame 和 parquet 文件执行计数。
下面是一些使用 pyspark 重现的代码:
logs_df = spark.read.csv('s3://../logs_2012/', multiLine=True, schema=get_schema()')
logs_df.filter(logs_df.date=='2012-11-22').count() # results in 5000
logs_df.write.partitionBy('date').parquet('s3://.../logs_2012_parquet/', mode='append')
par_df = spark.read.parquet('s3://.../logs_2012_parquet/')
par_df.filter(par_df.date=='2012-11-22').count() # results in 4999, always the same record that is omitted
我也尝试过写入 HDFS,结果是一样的。这发生在多个分区上。默认/空分区中没有记录。以上logs_df准确无误。
我尝试的第二个实验是编写未分区的拼花文件。上面代码的唯一区别是省略了partitionBy():
logs_df.write.parquet('s3://.../logs_2012_parquet/', mode='append')
加载此镶木地板组并按上述方式运行计数会为date='2012-11-22' 和其他日期产生正确的结果 5000。将模式设置为overwrite 或不设置(使用默认值)会导致相同的数据丢失。
我的环境是:
- EMR 5.9.0
- Spark 2.2.0
- Hadoop 发行版:Amazon 2.7.3
- 尝试使用 EMRFS 一致视图,但未尝试。但是,大多数测试都是写入 HDFS 以避免任何 S3 一致性问题。
我非常感谢修复或解决方法或使用 Spark 转换为 parquet 文件的其他方式。
谢谢,
编辑:我无法重现第二个实验。因此,假设分区和未分区在写入 Parquet 或 JSON 时似乎都会删除记录。
【问题讨论】:
-
总是丢失相同的记录吗?您的数据框中的所有记录的日期是否已明确定义?
-
是的,它是同一个。日期格式正确。 DataFrame 包含记录,因为我可以使用
filter()识别它。所以在使用 partitionBy 的时候总是在写的时候出错。 -
知道为什么某些行在写入 hdfs 或 s3 时会丢失吗?我也尝试将列作为字符串。不知道为什么。
-
也作为一个实验尝试写出 JSON 文件并读回,相同的行被删除:'(
标签: amazon-web-services apache-spark pyspark spark-dataframe parquet