【问题标题】:Spark 2.2.0 on AWS EMR writing to Parquet drops rowsAWS EMR 上的 Spark 2.2.0 写入 Parquet 删除行
【发布时间】:2018-04-15 08:51:57
【问题描述】:

所以我有一个问题,即在写入分区 Parquet 文件时,DataFrame 中的某些行会被丢弃。

这是我的步骤:

  1. 使用指定架构从 S3 读取 CSV 数据文件
  2. 按“日期”列 (DateType) 分区
  3. 写成 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


【解决方案1】:

所以谜团肯定在模式定义中。然而,出乎意料的是它不是日期或时间戳。它实际上是布尔值。

我已经从 Redshift 导出了 CSV,它将布尔值写为 tf。当我检查推断的架构时,这些字段被标记为字符串类型。在 CSV 文件中使用 truefalse 进行的简单测试将它们识别为布尔值。

所以我预计日期和时间戳解析会像往常一样出错,但它是布尔值。经验教训。

感谢指点。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-04-28
    • 2018-10-27
    • 2015-10-06
    • 2019-06-07
    • 2018-09-19
    • 2016-10-03
    • 1970-01-01
    • 2018-07-07
    相关资源
    最近更新 更多