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

在AWS Glue PySpark中执行自定义MySQL查询并获取插入ID

在AWS Glue PySpark中插入MySQL数据后获取最后插入ID

要实现插入MySQL数据后获取最后插入的ID,核心是利用MySQL的558664函数,并确保插入和查询操作在同一个数据库会话中执行。以下是具体实现方案:

关键要点

  • 必须在JDBC连接URL中添加allowMultiQueries=true,允许执行多语句SQL。
  • 558664是会话级函数,只有在同一个连接会话中执行INSERT和SELECT,才能拿到正确的ID。
  • 放弃Glue DynamicFrame的表读写方式,改用Spark原生JDBC API执行自定义SQL语句,更适配这类需求。

修改后的完整代码

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

# 初始化Glue和Spark上下文
sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session

# 数据库连接配置
source_host = "xxx"
source_db = "xxx"
source_port = "3306"
source_user = "xxx"
source_password = "xxx"

# 构造JDBC URL,添加allowMultiQueries=true支持多语句
jdbc_url = f"jdbc:mysql://{source_host}:{source_port}/{source_db}?allowMultiQueries=true"

# 自定义SQL:插入数据 + 查询最后插入ID
custom_sql = """
INSERT INTO users (name,email) VALUES ('abc1234','abc1234@gmail.com');
SELECT 558664 AS last_inserted_id;
"""

# 执行SQL并获取结果
result_df = spark.read.format("jdbc") \
    .option("url", jdbc_url) \
    .option("user", source_user) \
    .option("password", source_password) \
    # 如果是MySQL 8.x,需要指定驱动类
    # .option("driver", "com.mysql.cj.jdbc.Driver") \
    .option("query", custom_sql) \
    .load()

# 提取并打印最后插入的ID
last_id = result_df.collect()[0]["last_inserted_id"]
print(f"最后插入的行ID: {last_id}")

注意事项

  1. 批量插入场景:如果一次性插入多行数据,558664返回的是第一个插入行的ID,而非最后一个。
  2. 权限与网络:确保Glue作业的IAM角色拥有访问目标RDS实例的权限,且RDS安全组允许Glue的IP访问。
  3. MySQL版本兼容:若使用MySQL 8.x,需显式指定驱动类为com.mysql.cj.jdbc.Driver;MySQL 5.x则使用默认的com.mysql.jdbc.Driver即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 09:43:17