【问题标题】:Why pyspark runs out of memory for a map only program?为什么 pyspark 的仅地图程序的内存不足?
【发布时间】: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


    【解决方案1】:

    Spark 有很多活动部件。它将数据从文本读取到分区(通常在内存中),您正在解码 json,如果您的行很长(即大型 json 对象)可能会导致问题,您执行 partitionBy 可能有太多元素。

    我会首先尝试增加分区的数量(即使用重新分区而不是只减少分区数量的合并),我也会尝试在没有 partitionBy 的情况下进行写入,如果尝试查找失败最长的 json 并尝试仅分析它(即将 json 字符串映射到长度并取最长)。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-05-13
      • 2012-08-12
      • 1970-01-01
      • 1970-01-01
      • 2023-03-19
      • 2012-06-19
      • 2010-10-10
      相关资源
      最近更新 更多