【发布时间】: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