【发布时间】:2016-12-11 06:03:52
【问题描述】:
我有一个简单的 pyspark 程序,它一次读取 2 个文本文件,将每一行转换为 json 对象并将其写入 parquet 文件,如下所示:
for f in chunk(files, 2):
file_rdd = sc.textFile(f)
df = (file_rdd
.map(decode_to_json).filter(None)
.toDF(schema)
.coalesce(5)
.write
.partitionBy("created_year", "created_month")
.mode("append")
.parquet(file_output))
我用yarn运行作业,配置是这样的:
conf = (SparkConf()
.setAppName(app_name)
.set("spark.executor.memory", '6g')
.set('spark.executor.instances', '6')
.set('spark.executor.cores', '2')
.set("parquet.enable.summary-metadata", "false")
.set("spark.sql.parquet.compression.codec", 'snappy')
)
这看起来像一个仅限地图的程序,为什么它会因大型输入文件而出现内存不足的情况?
【问题讨论】:
标签: pyspark