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

