在AWS Glue中读写Aurora MySQL及执行更新语句的问题
解决AWS Glue执行Aurora MySQL更新语句及VPC连接问题
问题背景
现有AWS Glue作业已实现从S3读取CSV、通过Spark SQL生成DataFrame后转Dynamic Frame写入Aurora MySQL的功能,但无法支持UPDATE/REPLACE这类DML操作。尝试自定义JDBC驱动代码连接同账号VPC内的无服务器Aurora MySQL时,连接失败。
原写入代码:
spark.read.option("header","true").option("delimiter",",").option("quote",'"').option("escape",'"').option("multiline","true").csv(s3_input_path).createOrReplaceTempView("transformed_data") sql_query = """SELECT column1,column2 FROM transformed_data """ result_df = spark.sql(sql_query) logging.info(f"result_dyf") result_dyf = DynamicFrame.fromDF(result_df, glueContext, "result_dyf") glueContext.write_dynamic_frame.from_jdbc_conf( result_dyf, "my_saved_connection", connection_options={"dbtable": "TARGETTABLE", "database": "DBNAME"} )
失败的自定义JDBC连接代码:
mysql_jar_s3_location = "s3://your-bucket/mysql-connector-java-8.0.27.jar" # Set the subnet ID and security group ID subnet_id = "subnet-12345678" # Replace with your actual subnet ID security_group_id = "sg-12345678" # Replace with your actual security group ID # Set up Spark and Glue contexts sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session # Include the MySQL JDBC driver JAR file from S3 as a script library spark._jsc.addJar(mysql_jar_s3_location) # Specify the VPC, subnet, and security group details vpc_config = { "SubnetIds": [subnet_id], "SecurityGroupIds": [security_group_id] } # Configure your Aurora MySQL connection options aurora_connection_options = { "url": "jdbc:mysql://your-aurora-hostname:3306/your-database-name", "user": "your-username", "password": "your-password", "driver": "com.mysql.cj.jdbc.Driver", "vpcConfig": vpc_config } # Example SQL UPDATE query update_query = "UPDATE your-table SET column1 = 'new-value' WHERE condition" # Execute the UPDATE query result_df = spark.sql(update_query)
核心问题分析
- Dynamic Frame局限性:仅支持数据写入(append/overwrite模式),无法执行DML语句。
- VPC连接配置错误:代码中传入
vpcConfig无效,Glue作业需通过控制台配置VPC访问权限,而非代码参数。 - DML执行方式错误:
spark.sql()仅操作Spark临时表,无法直接执行JDBC数据库的DML语句。
解决方案
1. 配置Glue作业的VPC访问权限
无服务器Aurora MySQL在VPC内,必须让Glue作业加入对应VPC才能建立连接:
- 进入Glue作业编辑页面 → 安全配置、脚本库和作业参数 → 网络
- 勾选「使用VPC」,选择Aurora所在的子网和安全组
- 确保安全组规则允许Glue作业的IP访问Aurora的3306端口(或直接允许安全组之间的入站3306流量)
2. 正确执行DML语句的代码实现
方式一:利用Spark JDBC直接执行DML(推荐)
from pyspark.context import SparkContext from awsglue.context import GlueContext # 初始化上下文 sc = SparkContext.getOrCreate() glueContext = GlueContext(sc) spark = glueContext.spark_session # JDBC连接参数(无需在代码中配置VPC,作业已通过控制台配置) jdbc_url = "jdbc:mysql://your-aurora-hostname:3306/your-database-name" jdbc_props = { "user": "your-username", "password": "your-password", "driver": "com.mysql.cj.jdbc.Driver" } # 定义DML执行函数 def run_dml(query): # 获取JDBC连接 conn = spark._jvm.java.sql.DriverManager.getConnection( jdbc_url, jdbc_props["user"], jdbc_props["password"] ) try: stmt = conn.createStatement() affected_rows = stmt.executeUpdate(query) print(f"执行成功,影响行数:{affected_rows}") finally: conn.close() # 执行UPDATE示例 update_sql = "UPDATE your-table SET column1 = 'new-value' WHERE condition" run_dml(update_sql)
方式二:复用已保存的Glue Connection
如果已有保存的Glue Connection(如my_saved_connection),可直接提取连接参数,避免硬编码:
from pyspark.context import SparkContext from awsglue.context import GlueContext sc = SparkContext.getOrCreate() glueContext = GlueContext(sc) spark = glueContext.spark_session # 提取已保存的Connection参数 conn_name = "my_saved_connection" conn_params = glueContext.extract_jdbc_conf(conn_name) jdbc_url = conn_params["url"] jdbc_props = { "user": conn_params["user"], "password": conn_params["password"], "driver": "com.mysql.cj.jdbc.Driver" } # 执行DML def run_dml(query): conn = spark._jvm.java.sql.DriverManager.getConnection(jdbc_url, jdbc_props["user"], jdbc_props["password"]) try: stmt = conn.createStatement() stmt.executeUpdate(query) print(f"DML语句执行完成:{query}") finally: conn.close() run_dml("UPDATE TARGETTABLE SET column2 = 'updated' WHERE column1 = 'some-value'")
3. 驱动配置注意事项
- Glue默认已包含MySQL JDBC驱动,无需手动用
spark._jsc.addJar()加载 - 若需自定义驱动版本,直接在Glue作业的「脚本库」中添加S3路径的JAR包即可,作业会自动加载
验证要点
- 确认Glue作业的VPC、子网与Aurora集群完全匹配
- 检查Aurora安全组允许Glue作业所在安全组的3306端口入站流量
- 验证JDBC URL中的主机名正确(无服务器Aurora端点格式:
cluster-name.cluster-xxxxxx.region.rds.amazonaws.com) - 确保数据库账号拥有对应表的UPDATE权限
内容的提问来源于stack exchange,提问作者Santhosh
相关产品推荐
相关产品推荐

