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

AWS Glue作业503错误解决后运行过久,求代码优化方案

AWS Glue ETL作业性能优化方案

问题背景

在AWS Glue上调度ETL作业,已解决503请求过多错误,但200MB压缩文件的作业运行5小时仍未完成,需优化代码缩短执行时长。

核心优化措施

1. 移除低效的手动分批循环

原代码通过collect()获取分区值后循环分批,再用limit()截取数据,存在三个致命问题:

  • collect()将数据拉取到Driver节点,易引发内存溢出
  • 循环分批触发多次Spark Job,产生大量重复计算
  • limit()+orderBy()无法实现正确分页(每次limit都会从头读取数据,而非跳过已处理部分)

优化:直接利用Spark的分区写入+分桶功能,无需手动分批,Spark会自动处理数据分片与文件生成。

2. 调整数据扁平化顺序,减少Join数据量

原代码先执行Join再做explode操作,导致Join后数据先膨胀再扁平化,大幅增加计算开销。应先完全扁平化in_network数据,再执行Join,缩小Join时的数据规模。

3. 优化Join类型,避免Full Outer Join

Full Outer Join是开销最大的Join类型,若业务逻辑允许,优先替换为Inner Join或Left Join;若必须保留Full Outer Join,需确保关联字段有索引或数据分布均匀。

4. 简化分区策略,避免重复Shuffle

原代码先按billing_code repartition,之后又按billing_code_type和billing_code partitionBy,重复分区会增加Shuffle开销。只需保留最终的partitionBy+bucketBy即可,Spark会自动处理数据分布。

5. 移除不必要的Count操作

原代码中partition_data.count()会触发额外的Spark Job,增加运行时间,手动分批无需统计总行数,直接交给Spark处理即可。

修改后的完整代码

import sys
from pyspark.context import SparkContext
from pyspark.sql import SparkSession
from pyspark.sql.functions import explode, broadcast

sc = SparkContext.getOrCreate()
spark = SparkSession(sc)

# 读取基础数据
base_df = spark.read.json('s3://yourbucket03/Atena/preprocessed/base/')

# 关联provider_references与base数据(broadcast小表减少Shuffle)
prvd_df = spark.read.json('s3://yourbucket03/Atena/preprocessed/provider_references/')
prvd_df = prvd_df.crossJoin(broadcast(base_df))

# 先完全扁平化in_network数据,缩小Join数据规模
in_ntwrk_df = spark.read.json('s3://yourbucket03/Atena/preprocessed/in_network/')

# 第一层展开:negotiated_rates
in_ntwrk_df2 = in_ntwrk_df.select(
    'billing_code', 'billing_code_type', 'billing_code_type_version',
    'description', 'name',
    explode('negotiated_rates').alias('exploded_negotiated_rates'),
    'negotiation_arrangement'
)

# 第二层展开:provider_references,同时提取negotiated_prices
in_ntwrk_df3 = in_ntwrk_df2.select(
    'billing_code', 'billing_code_type', 'billing_code_type_version',
    'description', 'name',
    'exploded_negotiated_rates.negotiated_prices',
    explode('exploded_negotiated_rates.provider_references').alias('provider_ref_id'),
    'negotiation_arrangement'
)

# 第三层展开:negotiated_prices
in_ntwrk_df4 = in_ntwrk_df3.select(
    'billing_code', 'billing_code_type', 'billing_code_type_version',
    'description', 'name',
    explode('negotiated_prices').alias('negotiated_price'),
    'provider_ref_id',
    'negotiation_arrangement'
)

# 提取negotiated_price中的字段,完成扁平化
in_ntwrk_flat = in_ntwrk_df4.select(
    'billing_code', 'billing_code_type', 'billing_code_type_version',
    'description', 'name',
    'provider_ref_id',
    'negotiation_arrangement',
    'negotiated_price.billing_class',
    'negotiated_price.billing_code_modifier',
    'negotiated_price.expiration_date',
    'negotiated_price.negotiated_rate',
    'negotiated_price.negotiated_type',
    'negotiated_price.service_code'
)

# 优化Join类型:根据业务需求替换为'inner'或'left',此处以left join为例
jdf = prvd_df.join(
    in_ntwrk_flat,
    prvd_df.provider_group_id == in_ntwrk_flat.provider_ref_id,
    'left'
)

# 扁平化provider_groups和npi字段
jdf2 = jdf.select(
    'reporting_entity_name', 'reporting_entity_type', 'last_updated_on', 'version',
    'provider_group_id',
    explode('provider_groups').alias('exploded_provider_groups'),
    'billing_code', 'billing_code_type', 'billing_code_type_version',
    'description', 'name',
    'billing_class', 'billing_code_modifier',
    'expiration_date', 'negotiated_rate',
    'negotiated_type', 'service_code',
    'negotiation_arrangement'
)

# 最终扁平化npi,并重命名字段
jdf_final = jdf2.select(
    'reporting_entity_name', 'reporting_entity_type', 'last_updated_on', 'version',
    'provider_group_id',
    explode('exploded_provider_groups.npi').alias('npi'),
    'exploded_provider_groups.tin.type',
    'exploded_provider_groups.tin.value',
    'billing_code', 'billing_code_type', 'billing_code_type_version',
    'description', 'name',
    'billing_class', 'billing_code_modifier',
    'expiration_date', 'negotiated_rate',
    'negotiated_type', 'service_code',
    'negotiation_arrangement'
).withColumnRenamed('type', 'tin_type').withColumnRenamed('value', 'tin')

datasink_path = "s3://yourbucket03/Atena/processed/billing_code_npi/parquet/Atena_batched/"

# 直接写入,利用Spark自动分区、分桶,无需手动分批
jdf_final.write.format('parquet') \
    .mode("overwrite") \
    .partitionBy('billing_code_type', 'billing_code') \
    .bucketBy(4, 'npi') \
    .saveAsTable('Atena_table', path=datasink_path)

AWS Glue作业配置优化

  • Worker配置:使用G.1X或G.2X类型Worker,数量设置为5-10个,替代默认标准Worker提升计算性能
  • Spark参数调整:在作业配置中添加以下参数:
    • --conf spark.sql.shuffle.partitions=200(调整Shuffle分区数,根据数据量灵活调整)
    • --conf spark.sql.autoBroadcastJoinThreshold=104857600(设置100MB自动广播阈值,小表自动广播减少Shuffle)
    • --conf spark.dynamicAllocation.enabled=true(开启动态资源分配,根据负载自动调整Worker数量)
  • 预处理优化:若预处理后的JSON为大量小文件,先合并为大文件再处理,减少Spark读取文件的开销

内容的提问来源于stack exchange,提问作者Harish Kanta

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 06:57:01