AWS Glue作业向S3写入Parquet时任务中止解决方案咨询
问题描述
我的代码如下,包含数据转换逻辑:
dictionaryDf = spark.read.option("header", "true").csv( "s3://...../.csv") web_notif_data = fullLoad.cache() web_notif_data.persist(StorageLevel.MEMORY_AND_DISK) print("::::::data has been loaded::::::::::::") distinct_campaign_name = web_notif_data.select( trim(web_notif_data.campaign_name).alias("campaign_name")).distinct() web_notif_data.createOrReplaceTempView("temp") variablesList = Config.get('web', 'variablesListWeb') web_notif_data = spark.sql(variablesList) web_notif_data.persist(StorageLevel.MEMORY_AND_DISK) web_notif_data = web_notif_data.withColumn("camp", regexp_replace("campaign_name", "_", "")) web_notif_data = web_notif_data.drop("campaign_name") web_notif_data = web_notif_data.withColumnRenamed("camp", "campaign_name") web_notif_data = web_notif_data.withColumn("channel", lit("web_notification")) web_notif_data.createOrReplaceTempView("data") campaignTeamWeb = Config.get('web', 'campaignTeamWeb') web_notif_data = spark.sql(campaignTeamWeb) web_notif_data.persist(StorageLevel.MEMORY_AND_DISK) distinct_campaign_name = distinct_campaign_name.withColumn("camp", F.regexp_replace( F.lower(F.trim(col("campaign_name"))), "[^a-zA-Z0-9]", "")) output_df3 = ( distinct_campaign_name.withColumn("cname_split", F.explode(F.split(F.lower(F.trim(col("campaign_name"))), "_"))) .join( dictionaryDf, ( ( (F.col("function") == "contains") & F.col("camp").contains(F.col("terms")) ) | ( (F.col("function") == "match") & F.col("campaign_name").contains("_") & (F.col("cname_split") == F.col("terms")) ) ), "left" ) .withColumn( "empty_is_other", F.when( ( F.col("product").isNull() & F.col("product_category").isNull() ), "other" ) ) .withColumn( "rn", F.row_number().over( Window.partitionBy("campaign_name") .orderBy( F.when( F.col("function").isNull(), 3 ).when( F.col("function") == "match", 2 ).otherwise(1), F.length(F.col("terms")).desc(), F.col("product").isNull() ) ) ) .filter("rn=1") .select( "campaign_name", F.coalesce("product", "empty_is_other").alias("prod"), F.coalesce("product_category", "empty_is_other").alias("prod_cat"), ) .na.fill("") ) print(":::::::::::transformations have been done finally::::::::::::") web_notif_data1 = web_notif_data # Just taking the backup of DF in case something goes wrong web_notif_data = web_notif_data.drop("campaign_name") web_notif_data = web_notif_data.withColumnRenamed("temp_campaign_name", "campaign_name") veryFinalDF = web_notif_data.join(output_df3, "campaign_name", "left_outer") # veryFinalDF.show(truncate=False) veryFinalDF.write.mode("overwrite").parquet(aggregatedPath) print("::::final data have been written successfully::::::")
其中fullLoad是从Redshift表读取的DataFrame。该代码在20万条记录的场景下运行正常,但生产环境中15天的数据集最少有2500万条记录,数据存储在Redshift表中,读取后再进行处理。通过Glue作业运行该代码时,卡在最后一步写入Parquet数据时出错。
已经尝试使用30个executor运行,从Redshift加载数据到fullLoad DataFrame需要约20分钟。
优化方案
- 清理重复持久化逻辑:代码中对
web_notif_data重复调用cache()和persist(),属于多余操作,cache()本身就是MEMORY_ONLY级别的持久化,每个DataFrame仅保留一次必要的持久化即可,避免浪费内存资源。 - 优化Redshift读取速度:使用Glue Redshift连接器时开启并行读取,指定
numPartitions参数,按照日期等低基数字段做分片切分,2500万数据正常可将读取时间压缩到5分钟以内。 - 优化join逻辑:
output_df3是仅包含去重campaign_name的小维度表,最后一步join时使用broadcast(output_df3)做广播join,避免大量shuffle操作,大幅降低资源消耗。 - 处理数据倾斜:代码中的窗口函数、explode操作容易产生数据倾斜,先统计
campaign_name的分布,若存在热点key,给热点key添加随机前缀拆分处理,避免单个executor负载过高崩溃。 - 优化写入配置:写入Parquet前对结果DataFrame做重分区,控制每个分区大小在128MB左右,同时开启Snappy压缩,添加Spark配置
spark.sql.parquet.compression.codec=snappy,减少写入数据量。 - 调整资源参数:如果使用G.1X规格的Glue executor,添加以下Spark配置避免内存溢出:
spark.executor.memory=12g spark.driver.memory=8g spark.memory.fraction=0.8
内容的提问来源于stack exchange,提问作者whatsinthename
相关产品推荐
相关产品推荐

