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

AWS Glue自定义视觉脚本无限运行问题求助

自定义Glue Visual Transform实现MySQL表截断时任务无限运行的问题

需求

通过Custom Visual Transform在加载数据前截断MySQL表,不修改Glue自动生成的脚本。

问题现象

任务持续运行无终止,仅输出以下日志:

23/05/14 04:25:00 INFO MultipartUploadOutputStream: close closed:false s3://aws-glue-assets-849950158560-ap-south-1/sparkHistoryLogs/spark-application-1684037765713.inprogress

可正常运行的简化代码

仅保留过滤逻辑时,任务可正常执行:

from awsglue import DynamicFrame

def truncate_mysql_table(self, database_name, table_name, connection_name):
    return self.filter(lambda row: row['age'] == '21')

DynamicFrame.truncate_mysql_table = truncate_mysql_table

故障复现的完整代码

包含MySQL截断逻辑的完整代码导致任务无限运行:

import pymysql
import boto3
import json

from awsglue import DynamicFrame

def truncate_mysql_table(self, database_name, table_name, connection_name):
    client = boto3.client('glue')
    response = client.get_connection(Name=connection_name, HidePassword=False)
    connection_props = response.get("Connection").get("ConnectionProperties")
    host_name = connection_props.get("JDBC_CONNECTION_URL").rsplit(":", 1)[0].split("//")[1]
    port = int(connection_props.get("JDBC_CONNECTION_URL").rsplit(":", 1)[1].split("/", 1)[0])
    secret_id = connection_props.get("SECRET_ID")

    client = boto3.client('secretsmanager')
    response = client.get_secret_value(SecretId=secret_id)

    secret_data = json.loads(response.get("SecretString"))
    username = secret_data.get("username")
    password = secret_data.get("password")

    con = pymysql.connect(host=host_name,
                          user=username,
                          passwd=password,
                          db=database_name,
                          port=port,
                          connect_timeout=60)

    with con.cursor() as cur:
        cur.execute(f"TRUNCATE TABLE {database_name.strip()}.{table_name.strip()}")
        con.commit()
        con.close()

    # print("Table Truncated")
    return self

DynamicFrame.truncate_mysql_table = truncate_mysql_table

环境说明

  • Glue Connection与MySQL RDS处于同一VPC
  • 已配置S3和Secrets Manager的VPC端点
  • 简化代码可正常执行,排除基础环境问题

解决方案

1. 定位卡点:添加日志排查

取消代码中print("Table Truncated")的注释,或在关键步骤添加日志输出,比如:

# 获取连接属性后
print(f"Got DB host: {host_name}, port: {port}")
# 获取密钥后
print(f"Fetched username: {username}")
# 连接数据库后
print("Connected to MySQL successfully")
# 执行截断后
print("Truncate command executed")

通过日志确认任务卡在哪个环节(比如数据库连接、密钥获取还是截断执行)。

2. 修复数据库连接资源管理

原代码中手动关闭连接的方式可能存在资源泄漏或阻塞风险,改用上下文管理器自动管理连接生命周期:

# 替换原有的pymysql连接及执行代码
with pymysql.connect(host=host_name,
                     user=username,
                     passwd=password,
                     db=database_name,
                     port=port,
                     connect_timeout=60) as con:
    with con.cursor() as cur:
        cur.execute(f"TRUNCATE TABLE {database_name.strip()}.{table_name.strip()}")
        con.commit()
# 无需手动调用con.close(),上下文管理器会自动处理

3. 改用Spark原生JDBC执行截断操作

避免使用pymysql的同步连接,改用Glue/Spark原生的JDBC方式执行DDL,更适配Spark执行模型:

# 替换原有的pymysql相关代码
from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()
jdbc_url = f"jdbc:mysql://{host_name}:{port}/{database_name}"
jdbc_props = {
    "user": username,
    "password": password,
    "driver": "com.mysql.cj.jdbc.Driver"
}

# 通过Spark JDBC执行截断
spark._jvm.java.sql.DriverManager.getConnection(
    jdbc_url, 
    jdbc_props["user"], 
    jdbc_props["password"]
).createStatement().execute(f"TRUNCATE TABLE {database_name.strip()}.{table_name.strip()}")

4. 验证IAM权限细节

确保Glue任务的IAM角色拥有以下权限:

  • glue:GetConnection:获取Glue连接属性
  • secretsmanager:GetSecretValue:读取数据库密钥
  • MySQL数据库的TRUNCATE权限:确认数据库用户拥有目标表的截断权限

5. 检查网络与安全组配置

即使处于同一VPC,仍需验证:

  • MySQL RDS的安全组允许Glue作业所在子网的访问(3306端口)
  • VPC端点的策略允许Glue访问Secrets Manager和S3

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 02:32:41