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

如何优雅终止AWS Glue PySpark作业并完成自定义资源清理

AWS Glue v3.0 PySpark 作业优雅终止实现方案

你之前采用signal.signal()捕获信号的方案失效,是因为Glue PySpark作业的终止信号首先会发给YARN ResourceManager/Spark Driver进程,而非直接投递到Python脚本主进程,且Glue最终强制终止时发送的SIGKILL信号本身无法被任何进程捕获,因此该方案不可用。

以下是三种可落地的实现方案:


1. 注册Spark Driver 关闭钩子(最通用方案)

Spark Driver提供了JVM级别的关闭钩子管理能力,只要Driver没有被强制杀死,无论是正常执行完成、抛出异常还是收到集群终止指令,都会按优先级执行注册的钩子,适合绝大多数场景。

from pyspark.context import SparkContext
from awsglue.context import GlueContext

sc = SparkContext()
glueContext = GlueContext(sc)

# 自定义资源清理逻辑
def custom_cleanup():
    # 此处写入你的清理逻辑,例如关闭自定义数据库连接、删除临时中间文件、上报作业状态等
    print("开始执行自定义资源清理")
    # 示例:关闭自定义连接
    # custom_mysql_conn.close()
    print("自定义资源清理完成")

# 注册JVM关闭钩子
shutdown_hook_manager = sc._gateway.jvm.org.apache.spark.util.ShutdownHookManager
# 将Python函数封装为JVM可识别的Runnable实例
runnable = sc._gateway.jvm.Runnable(custom_cleanup)
# 第二个参数为优先级,数值越大越先执行
shutdown_hook_manager.get().addShutdownHook(10, runnable)

# 后续正常写作业核心逻辑

2. 轮询检测Glue作业运行状态(可控度最高方案)

Glue作业收到终止指令后,会先进入STOPPING状态,最多等待2分钟才会强制杀死所有进程。你可以在作业核心逻辑的关键节点(比如批量处理的每轮循环前)调用Glue API查询自身作业状态,检测到终止状态后主动执行清理再退出。

注意:需要给作业绑定的IAM角色添加glue:GetJobRun权限

import boto3
import sys
from awsglue.utils import getResolvedOptions

args = getResolvedOptions(sys.argv, ['JOB_NAME', 'JOB_RUN_ID'])
glue_client = boto3.client('glue')

def check_job_stopping():
    resp = glue_client.get_job_run(
        JobName=args['JOB_NAME'],
        RunId=args['JOB_RUN_ID']
    )
    return resp['JobRun']['JobRunState'] == 'STOPPING'

# 示例:批量处理场景每轮循环前检测状态
for batch in all_batch_list:
    # 先检测是否收到终止指令
    if check_job_stopping():
        custom_cleanup() # 执行自定义清理
        sc.stop()
        sys.exit(0)
    # 处理当前批次数据
    process_single_batch(batch)

3. try-finally块包裹全量逻辑(轻量场景方案)

如果你的作业逻辑没有长周期循环,清理逻辑执行耗时很短,可以把所有核心逻辑放在try块中,finally块写入清理逻辑,只要Python进程没有被SIGKILL强制杀死,finally块的代码一定会执行。

try:
    # 所有作业核心逻辑写在此处
    run_main_job_logic()
except Exception as e:
    # 自定义异常处理逻辑
    print(f"作业执行出错:{str(e)}")
    raise e
finally:
    # 无论正常结束、执行报错、收到终止指令(进程未被强制杀死)都会执行此处
    custom_cleanup()
    sc.stop()

注意事项

  • 所有清理逻辑的执行耗时请控制在2分钟以内,超过2分钟Glue会强制杀死所有进程,导致清理逻辑中断
  • 如果是Glue流式作业,建议配合StreamingQuery.awaitTermination()方法实现优雅停止,流式作业的终止等待窗口期更长
  • 不要尝试捕获SIGKILL信号,该信号由操作系统内核直接发送,任何进程都无法拦截或捕获

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 18:09:00