Redshift含IDENTITY列表通过Glue脚本插入时报表不存在错误
Redshift表Glue插入报错“Table not found”问题排查
问题背景
在Redshift中创建了包含IDENTITY列的schema_name.employee表,通过DBeaver执行创建和插入语句均正常,但使用Glue脚本执行插入时,出现“Table not found: schema_name.employee”错误,且脚本可正常读取该表数据。
创建表及DBeaver插入语句
create table schema_name.employee ( surrogate_key bigint IDENTITY(1,1), first_name varchar(200), last_name varchar(200), phone_number varchar(200), creditcard_number bigint ) insert into schema_name.employee(first_name,last_name,phone_number,creditcard_number) values('gaurang', 'shah', '356-776-4456', 4716973408090483)
Glue插入脚本
insert_query = """insert into schema_name.employee(first_name,last_name,phone_number,creditcard_number) values('gaurang', 'shah', '356-776-4456', 4716973408090483)""" spark.sql(insert_query).write \ .format("com.amazonaws.redshift.spark") \ .option("url", url) \ .option("dbtable","schema_name.employee") \ .option("user","user") \ .option("password", "password") \ .option("driver", driver) \ .mode("overwrite") \ .save()
错误详情
"name": "AnalysisException", "message": "Table not found: schema_name.employee; line 1 pos 12;\n'InsertIntoStatement 'UnresolvedRelation [schema_name, employee], [], false, [first_name, last_name, phone_number, creditcard_number], false, false\n+- LocalRelation [col1#147, col2#148, col3#149, col4#150L]\n", Table not found: schema_name.employee; line 1 pos 12;\n'InsertIntoStatement 'UnresolvedRelation [schema_name, employee], [], false, [first_name, last_name, phone_number, creditcard_number], false, false\n+- LocalRelation [col1#147, col2#148, col3#149, col4#150L]\n" }
问题原因
你的Glue脚本写法逻辑错误:
spark.sql(insert_query)是让Spark SQL引擎解析执行这条插入语句,此时Spark会在自己的Catalog中查找schema_name.employee表,但你并没有将Redshift的这张表注册到Spark Catalog中,所以Spark找不到该表,抛出“Table not found”错误。- 而读取表时你应该是用了
read.format("com.amazonaws.redshift.spark")的方式,这种方式是直接通过Redshift连接器读取Redshift集群中的表,不需要依赖Spark Catalog,所以能正常读取。
解决思路
方案1:构造DataFrame后写入Redshift(推荐)
Redshift Spark连接器的write接口是用来将DataFrame数据写入Redshift表的,正确写法如下:
# 构造要插入的数据为DataFrame data = [("gaurang", "shah", "356-776-4456", 4716973408090483)] columns = ["first_name", "last_name", "phone_number", "creditcard_number"] df = spark.createDataFrame(data, columns) # 将DataFrame写入Redshift表 df.write \ .format("com.amazonaws.redshift.spark") \ .option("url", url) \ .option("dbtable", "schema_name.employee") \ .option("user", "user") \ .option("password", "password") \ .option("driver", driver) \ .mode("append") # 注意用append模式,overwrite会删除原表重建,丢失IDENTITY列配置 .save()
方案2:直接通过JDBC执行SQL语句
如果一定要执行原生SQL插入,可以用JDBC连接直接操作Redshift,示例代码如下:
import pyodbc # 建立JDBC连接 conn = pyodbc.connect( f"Driver={driver};Server={server};Database={db};UID={user};PWD={password};Port={port}" ) cursor = conn.cursor() # 执行插入语句 insert_query = """insert into schema_name.employee(first_name,last_name,phone_number,creditcard_number) values('gaurang', 'shah', '356-776-4456', 4716973408090483)""" cursor.execute(insert_query) conn.commit() # 关闭连接 cursor.close() conn.close()
内容的提问来源于stack exchange,提问作者Bab
相关产品推荐
相关产品推荐

