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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 11:10:39