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

在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)

核心问题分析

  1. Dynamic Frame局限性:仅支持数据写入(append/overwrite模式),无法执行DML语句。
  2. VPC连接配置错误:代码中传入vpcConfig无效,Glue作业需通过控制台配置VPC访问权限,而非代码参数。
  3. 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包即可,作业会自动加载

验证要点

  1. 确认Glue作业的VPC、子网与Aurora集群完全匹配
  2. 检查Aurora安全组允许Glue作业所在安全组的3306端口入站流量
  3. 验证JDBC URL中的主机名正确(无服务器Aurora端点格式:cluster-name.cluster-xxxxxx.region.rds.amazonaws.com)
  4. 确保数据库账号拥有对应表的UPDATE权限

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 12:51:10