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
相关产品推荐
相关产品推荐

