【发布时间】: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 作业的源代码?
-
@cool 这更像是一种解决方法。我重新分区了数据并限制了文件大小。我知道平均记录大小。基于此,我设置 maxRecordsPerFile 以保持文件大小无限增长。
-
@cool 是的。这足以限制文件的数量。现在,只有在没有超过指定的记录时才会创建一个新文件。我用该代码创建了一个要点。看看有没有帮助gist.github.com/kapilgarg/ed605408aee21166ba6394a483e1f260
-
@Kapil 再次感谢,这会有所帮助。也许您最好自己在答案中详细说明此讨论并将其标记为已接受,只是为了解决问题。
标签: amazon-s3 pyspark amazon-athena