【问题标题】:PySpark GC overhead limit exceeded and timeout超过 PySpark GC 开销限制并超时
【发布时间】:2020-07-12 07:26:30
【问题描述】:

我在 pyspark 中编写了代码,任务是查找与前一天 parquet 文件相比的 Delta/更改记录,并将其写入 csv 并带有标记新记录或现有记录的附加列。这是通过连接列并应用 base64(encode(columns here)) 将其命名为 hash_key 来完成的。

  1. 读取 9 个 parquet 文件并针对它们注册临时表。

    • 对其应用查询,将输出命名为 df2。
    • 从 S3 读取前一天的 parquet 文件并将其命名为 df1 并获取它的计数。
    • 针对两个数据帧将临时表注册为 t1 和 t2。
  2. 现在我使用 sqlContext.sql 查找 t2 与 t1 相比的更改记录,使用 hash_key 即 df3 并获取行数。

  3. 将不同的保存在 csv 文件中,并将 parquet 文件替换为 df2 中的数据。

到目前为止,一切都运行良好。但后来的任务是显示增量记录中哪些行是全新的,哪些是旧行,其中一些列已更改。在这里,我尝试使用 PySpark 代码并加入了两个数据帧 df1 和 df3。但我收到以下错误 GC 开销限制超出超时

作为替代方案,我尝试针对数据帧注册临时表并对其执行 sql 查询。 但问题仍然存在。我真的被它困住了。不知道为什么会消耗资源。

配置:

conf = (SparkConf()
        .setAppName("GD_Regex")
        .set("spark.executor.instances", "3")
        .set("spark.executor.cores", "3")
        .set("spark.sql.parquet.enableVectorizedReader", "false")
        .set("spark.executor.memory", "4g")
        .set("fs.s3a.server-side-encryption-algorithm", "AES256")
        .set("spark.sql.autoBroadcastJoinThreshold","-1")
        )

使用 hash_key 发现差异

s3_parquet_file = sqlContext.read.parquet(parquet_path)
historical_file_row_count = s3_parquet_file.count()
s3_parquet_file.registerTempTable("s3_parquet_file_temp")
query_response.registerTempTable("query_response_temp")
filtered_data = sqlContext.sql(
        """select * from query_response_temp where hash_key NOT IN ( SELECT hash_key FROM s3_parquet_file_temp )""")

方法一:

ff = filtered_data.alias('df2').join(s3_parquet_file.alias('df1'), ( filtered_data.CTMS_STUDY_NUMBER == s3_parquet_file.CTMS_STUDY_NUMBER ) & ( filtered_data.CTMS_SITE_NUMBER == s3_parquet_file.CTMS_SITE_NUMBER ) 
        & (filtered_data.FULL_NAME_OF_PI == s3_parquet_file.FULL_NAME_OF_PI ) & ( filtered_data.PRIMARY_ROLE_NAME == s3_parquet_file.PRIMARY_ROLE_NAME) & ( filtered_data.CENTER_NAME == s3_parquet_file.CENTER_NAME ) , "outer" )\
        .select('df2.*', 'df1.CTMS_STUDY_NUMBER')\
        .withColumn('status', when(s3_parquet_file.CTMS_STUDY_NUMBER.isNotNull() & filtered_data.CTMS_STUDY_NUMBER.isNotNull(), 'existing').otherwise('new'))\
        .filter(filtered_data.CTMS_STUDY_NUMBER.isNotNull())\
        .select('df2.*', 'status')

方法二:

q = """SELECT query_response_temp.*, IF( query_response_temp.CTMS_STUDY_NUMBER IS NOT NULL AND s3_parquet_file_temp.CTMS_STUDY_NUMBER IS NOT NULL, "existing", "new" ) as status FROM query_response_temp FULL OUTER JOIN s3_parquet_file_temp on query_response_temp.CTMS_STUDY_NUMBER = s3_parquet_file_temp.CTMS_STUDY_NUMBER AND query_response_temp.CTMS_SITE_NUMBER = s3_parquet_file_temp.CTMS_SITE_NUMBER AND query_response_temp.FULL_NAME_OF_PI = s3_parquet_file_temp.FULL_NAME_OF_PI AND query_response_temp.PRIMARY_ROLE_NAME = s3_parquet_file_temp.PRIMARY_ROLE_NAME AND query_response_temp.CENTER_NAME = s3_parquet_file_temp.CENTER_NAME"""
        delta_records = sqlContext.sql(q)
        delta_records = delta_records.filter(delta_records.CTMS_STUDY_NUMBER.isNotNull())

发生错误时发生错误的代码:

dataset = dataset.drop("hash_key")
print('hash_key dropped.')
dataset.coalesce(1).write.format('csv') \
  .option('header', 'true') \
  .option("compression", "none") \
  .mode('overwrite') \
  .save(_path)

【问题讨论】:

    标签: python pandas amazon-s3 pyspark parquet


    【解决方案1】:

    增加executor memory 并尝试重新运行作业。

    【讨论】:

    • 我已将它增加到 6g,但它不起作用。并且绑定为它的开发服务器,我不能增加太多。
    猜你喜欢
    • 1970-01-01
    • 2020-07-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-05-21
    相关资源
    最近更新 更多