【问题标题】:pyspark writing lot of smaller files in outputpyspark 在输出中写入大量较小的文件
【发布时间】:2020-04-08 16:13:12
【问题描述】:

我正在使用 pyspark 处理一些数据并将输出写入 S3。我在 athena 中创建了一个表,用于查询这些数据。

数据是json字符串的形式(每行一个),spark代码读取文件,根据某些字段对其进行分区并写入S3。

对于 1.1 GB 的文件,我看到 spark 正在写入 36 个文件,每个文件大小约为 5 MB。在阅读 athena 文档时,我发现最佳文件大小为 ~128 MB。 https://aws.amazon.com/blogs/big-data/top-10-performance-tuning-tips-for-amazon-athena/

sparkSess = SparkSession.builder\
    .appName("testApp")\
    .config("spark.debug.maxToStringFields", "1000")\
    .config("spark.sql.sources.partitionOverwriteMode", "dynamic")\
    .getOrCreate()

sparkCtx = sparkSess.sparkContext
deltaRdd = sparkCtx.textFile(filePath)
df = sparkSess.createDataFrame(deltaRdd, schema)

try:
    df.write.partitionBy('field1','field2','field3')\
        .json(path, mode='overwrite', compression=compression)
except Exception as e:
    print (e)

为什么 spark 会写这么小的文件。有什么方法可以控制文件大小。

【问题讨论】:

  • 您能否发布您的 pyspark 代码以及您的 pyspark 作业的源代码?
  • 这能回答你的问题吗? How do you control the size of the output file?
  • @cool 这更像是一种解决方法。我重新分区了数据并限制了文件大小。我知道平均记录大小。基于此,我设置 maxRecordsPerFile 以保持文件大小无限增长。
  • @cool 是的。这足以限制文件的数量。现在,只有在没有超过指定的记录时才会创建一个新文件。我用该代码创建了一个要点。看看有没有帮助gist.github.com/kapilgarg/ed605408aee21166ba6394a483e1f260
  • @Kapil 再次感谢,这会有所帮助。也许您最好自己在答案中详细说明此讨论并将其标记为已接受,只是为了解决问题。

标签: amazon-s3 pyspark amazon-athena


【解决方案1】:

有没有办法控制文件大小?

有一些控制机制。但是,它们并不明确。

s3 驱动程序不是 spark 本身的一部分。它们是 spark emr 附带的 hadoop 安装的一部分。 s3 块大小可以设置在 /etc/hadoop/core-site.xmlconfig file

但默认情况下它应该是 128 mb 左右。

为什么 spark 会写这么小的文件

Spark 将遵守 hadoop 块大小。但是,您可以在写作之前使用partionBy

假设您使用partionBy("date").write.csv("s3://products/")。 Spark 将为每个分区创建一个带有date 的子文件夹。之内 每个分区文件夹 spark 将再次尝试创建块并尝试遵守fs.s3a.block.size

例如

s3:/products/date=20191127/00000.csv
s3:/products/date=20191127/00001.csv
s3:/products/date=20200101/00000.csv

在上面的示例中 - 特定分区可以小于 128mb 的块大小。

所以只需仔细检查/etc/hadoop/core-site.xml 中的块大小,以及在写入之前是否需要使用partitionBy 对数据帧进行分区。

编辑:

Similar post 还建议重新分区数据帧以匹配partitionBy 方案

df.repartition('field1','field2','field3')
.write.partitionBy('field1','field2','field3')

writer.partitionBy 对现有数据帧分区进行操作。它不会 repartition 原始数据框。因此,如果整个数据帧的分区不同,就会发生嵌套分区。

【讨论】:

  • 即使,完整的数据对应于单个分区,文件大小仍然不是 128。我在 core-site.xml 中没有任何条目,因此假设它采用默认大小,但在更新后分段上传的属性,没有变化
  • 你检查过hadoop s3 conf吗?什么是 s3 block.size?
  • 当我检查 spark-emr hadoop conf 时,s3 block.size 根本没有设置。我们编写的文件大约 150 mb。执行cd /etc/hadoop/conf/ 然后运行grep -rni "s3" * 以检查是否有block.size 设置
  • 刚刚还看到,我们重新分区数据帧以匹配数据写入器模式。为答案添加了另一个提示
猜你喜欢
  • 2013-03-13
  • 1970-01-01
  • 2021-09-12
  • 2019-01-24
  • 2018-05-30
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多