【问题标题】:How to debug MemoryError in PySpark如何在 PySpark 中调试 MemoryError
【发布时间】:2019-05-21 20:37:52
【问题描述】:

我想下载一些 xml 文件(每个 50MB - 大约 3000 = 150GB),处理它们并使用 pyspark 上传到 BigQuery。出于开发目的,我使用了 jupyter notebook 和少量文件 10。我在 dataproc 上编写了非常复杂的代码设置集群。我的 daproc 集群有 6TB 的 HDFS、10 个节点(每个 4 核)和 120GB 的 RAM。

def context():
    import os
    os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages org.apache.hadoop:hadoop-aws:2.7.3 pyspark-shell'
    import pyspark
    conf = pyspark.SparkConf()

    conf = (conf.setMaster('local[*]')
            .set('spark.executor.memory', '4G')
            .set('spark.driver.memory', '45G')
            .set('spark.driver.maxResultSize', '10G')
            .set("spark.python.profile", "true"))
    sc = pyspark.SparkContext(conf=conf)
    return sc
def job(sc):
    print("Job started")
    RDDread = sc.wholeTextFiles("s3a://custom-bucket/*/*.gz")
    models = RDDread.flatMap(process_xmls).groupByKey()
    tracking_all = (models.filter(lambda x: x[0] == TrackInformation)
                    .flatMap(lambda x: x[1])
                    .map(lambda model: (model.flight_ref, model))
                    .groupByKey())
    tracking_merged = tracking_all.map(lambda x: x[1]).map(merge_ti)
    flight_plans = (models.filter(lambda x: x[0] == FlightPlan).flatMap(lambda x: x[1]).map(lambda fp: (fp.flight_ref, fp)))
    fps_tracking = tracking_merged.union(flight_plans).groupByKey().filter(lambda x: len(x[1]) == 2)
    in_bq_batch = 1000
    n = fps_tracking.count()
    parts = ceil(n / in_bq_batch)
    many_n = fps_tracking.repartition(parts).mapPartitions(upload_fpm2)
    print("Job ended")
    return fps_tracking, tracking_merged, flight_plans, models, many_n

在 200 条消息org.apache.hadoop.io.compress.CodecPool: Got brand-new decompressor [.gz] 之后,我收到 2 个错误:java.lang.OutOfMemoryError 和 MemoryError,主要是 MemoryError。我以为我在 RDDread 之后只有 2 个分区,所以我修改了以下代码: sc.wholeTextFiles("s3a://custom-bucket//.gz", minPartitions=40) -> 并且它破产得更快。我在一些随机的地方添加了持久(磁盘)功能。

File "/usr/lib/spark/python/lib/pyspark.zip/pyspark/serializers.py", line 684, in loads
    return s.decode("utf-8") if self.use_unicode else s
MemoryError
19/05/20 14:09:23 INFO org.apache.hadoop.io.compress.CodecPool: Got brand-new decompressor [.gz]
19/05/20 14:09:30 ERROR org.apache.spark.util.Utils: Uncaught exception in thread stdout writer for /opt/conda/default/bin/python
java.lang.OutOfMemoryError: Java heap space

我做错了什么以及如何调试我的代码?

【问题讨论】:

    标签: apache-spark pyspark


    【解决方案1】:

    您似乎在本地模式下运行 spark (local[*])。这意味着您正在使用具有 45G RAM (spark.driver.memory) 的单个 jvm,并且所有工作线程都在该 jvm 中运行。 spark.executor.memory 选项无效What does setMaster `local[*]` mean in spark?

    您应该将 spark master 设置为 yarn 调度程序,或者如果您没有 yarn 使用独立模式 https://spark.apache.org/docs/latest/spark-standalone.html

    【讨论】:

      猜你喜欢
      • 2010-12-13
      • 1970-01-01
      • 1970-01-01
      • 2017-06-06
      • 2019-03-02
      • 1970-01-01
      • 2021-04-15
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多