AWS Glue中PySpark DataFrame过滤后操作耗时过长的优化请求
AWS Glue小数据集过滤后操作耗时过长的优化方案
问题背景
我有一个仅含4条记录的base_df DataFrame,数据如下:
+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+---------------+--------------------+-----------------+ | RECORDS_001| RECORDS_002| RECORDS_003| RECORDS_004| RECORDS_005| RECORDS_006| RECORDS_007| RECORDS_HEADER| input_file_name| ROW_ID|unique_batch_id| FAILED_VALIDATIONS|VALIDATION_STATUS| +--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+---------------+--------------------+-----------------+ |{{01/01/2023, 06/...|{[{07/19/2022, 07...|{[{0.00, 99213, D...|{[{0.00, 0.00, 0....|{{100.00, 29495.0...|{{{B}, 06/30/2023...|{[{H0354, 0.00, 0...|{{1, 7, 16, 7, 1,...|s3://gov-solution...|3f92edbf-4af9-48a...| 0|MEMBER_ID is miss...| FAILED| |{{01/01/2023, 11/...|{[{08/29/2023, 08...|{[{0.00, 96401, C...|{[{900.00, 0.00, ...|{{1235.77, 4700.0...|{{B, 11/30/2024, ...| null|{{1, 3, 6, 3, 1, ...|s3://gov-solution...|3f92edbf-4af9-48a...| 0| null| PASSED| |{{01/01/2023, 06/...|{[{07/19/2022, 07...|{[{0.00, 99213, D...|{[{0.00, 0.00, 0....|{{100.00, 29495.0...|{{{B}, 06/30/2023...|{[{H0354, 0.00, 0...|{{1, 7, 16, 7, 1,...|s3://gov-solution...|3f92edbf-4af9-48a...| 0| null| PASSED| |{{01/01/2023, 06/...|{[{07/19/2022, 07...|{[{0.00, 99213, D...|{[{0.00, 0.00, 0....|{{100.00, 29495.0...|{{{B}, 06/30/2023...|{[{H0354, 0.00, 0...|{{1, 7, 16, 7, 1,...|s3://gov-solution...|3f92edbf-4af9-48a...| 0|MEMBER_ID is missing| FAILED| +--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+---------------+--------------------+-----------------+
执行过滤操作后,后续操作耗时极长,使用的代码如下:
print("before filter") base_df = base_df.repartition(16) print(base_df.rdd.getNumPartitions()) print(base_df.show()) # 输出为上述数据 base_df = base_df.filter(base_df.VALIDATION_STATUS == 'PASSED') print("After filtering") # 此步骤前的作业在1分钟内完成,但后续打印语句耗时极长 print(base_df.rdd.getNumPartitions()) base_df.select('FAILED_VALIDATIONS').show(truncate=False)
运行环境:AWS Glue Job(版本4.0),工作节点类型G.8X。
优化方案
1. 移除不必要的重分区操作
仅4条记录的数据集被重分区到16个分区,会导致任务被拆分成大量小任务,调度、启动开销远大于实际计算开销。直接删除base_df = base_df.repartition(16),让Spark自动处理小数据集的分区。
2. 提前缓存数据集
对小数据集进行缓存,避免重复扫描源数据和计算。在过滤前添加缓存触发逻辑:
base_df.cache() base_df.count() # 触发缓存动作,将数据加载到内存
缓存后,后续的过滤、查询操作都会直接从内存读取数据,大幅降低耗时。
3. 调整Glue Job资源配置
G.8X属于大资源节点,对于极小数据集来说资源过剩,调度开销反而更大:
- 降低节点类型,改用G.2X甚至更小的节点
- 将节点数量设置为1个,避免多节点间的调度协调开销
4. 简化不必要的调用
print(base_df.show())会额外打印show()的返回值None,且触发一次DataFrame扫描,简化为base_df.show()即可。
优化后的示例代码
print("before filter") # 移除重分区,添加缓存 base_df.cache() base_df.count() print(base_df.rdd.getNumPartitions()) base_df.show() # 直接打印,无需嵌套print base_df = base_df.filter(base_df.VALIDATION_STATUS == 'PASSED') print("After filtering") print(base_df.rdd.getNumPartitions()) base_df.select('FAILED_VALIDATIONS').show(truncate=False)
内容的提问来源于stack exchange,提问作者Suraj Pandey
相关产品推荐
相关产品推荐

