在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}")
注意事项
- 批量插入场景:如果一次性插入多行数据,
558664返回的是第一个插入行的ID,而非最后一个。 - 权限与网络:确保Glue作业的IAM角色拥有访问目标RDS实例的权限,且RDS安全组允许Glue的IP访问。
- MySQL版本兼容:若使用MySQL 8.x,需显式指定驱动类为
com.mysql.cj.jdbc.Driver;MySQL 5.x则使用默认的com.mysql.jdbc.Driver即可。
内容的提问来源于stack exchange,提问作者Propsz
相关产品推荐
相关产品推荐

