You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.29 13:36:02