【问题标题】:Spark df partitioniong after partioning by yy/mm/dd按 yy/mm/dd 分区后的 Spark df 分区
【发布时间】:2020-05-04 20:43:58
【问题描述】:

S3 托管一个非常大的压缩文件(20gb 压缩 -> 200gb 未压缩)。 我想读取这个文件(不幸的是在单核上解压),转换一些sql列,然后以s3_path/year=2020/month=01/day=01/[files 1-200].parquet格式输出到S3。

整个文件将包含同一日期的数据。这使我相信我应该将"year={year}/month={month}/day={day}/" 附加到s3 路径而不是使用partitionBy('year','month','day'),因为目前spark 一次将一个文件写入s3(每个文件大小为1gb)。我的想法对吗?

这是我目前正在做的事情:

df = df\
    .withColumn('year', lit(datetime_object.year))\
    .withColumn('month', lit(datetime_object.month))\
    .withColumn('day', lit(datetime_object.day))

df\
    .write\
    .partitionBy('year','month','day')\
    .parquet(s3_dest_path, mode='overwrite')

我在想什么:

df = spark.read.format('json')\
    .load(s3_file, schema=StructType.fromJson(my_schema))\
    .repartition(200)
# currently takes a long time decompressing the 20gb s3_file.json.gz

# transform
df.write\
    .parquet(s3_dest_path + 'year={}/month={}/day={}/'.format(year,month,day))

【问题讨论】:

  • 我认为这样做没有任何害处。

标签: python dataframe apache-spark


【解决方案1】:

您可能遇到了 spark 先将数据写入某个 _temporary 目录,然后才将其提交到最终位置的问题。在 HDFS 中,这是通过重命名来完成的。但是 S3 不支持重命名,而是完全复制数据(仅使用一个执行程序)。有关此主题的更多信息,请参见例如此帖子:Extremely slow S3 write times from EMR/ Spark

常见的解决方法是写入 hdfs,然后使用 distcp 将分布式从 hdfs 复制到 s3

【讨论】:

  • 我不断听到人们说 S3 怎么可能。现在是 2020 年,这似乎很奇怪。
  • 我完全同意。我的知识是从 2018 年开始的......所以如果从那以后发生了什么,我很高兴得到更新
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-07-24
  • 2023-03-09
  • 2012-01-27
  • 1970-01-01
  • 2022-11-19
  • 2018-10-13
  • 1970-01-01
相关资源
最近更新 更多