使用AWS Glue创建Redshift空表失败求助:表结构异常且报错
解决AWS Glue创建Redshift表仅生成dummy列及SQLException问题
问题说明
尝试通过AWS Glue创建Redshift表flight.flight_details时,抛出错误exception: java.sql.SQLException: Exception thrown in awaitResult。尽管表最终被创建,但仅包含dummy列,完全不符合自定义CREATE TABLE语句中定义的表结构。
原因分析
- 提前调用
job.commit():原代码在表创建逻辑执行前就调用了job.commit(),导致Glue作业提前终止,后续代码执行不完整,引发异常。 - 空DataFrame触发schema覆盖:使用仅含
dummy列的空DataFrame执行写入操作时,Spark会根据DataFrame的schema修改目标表结构,忽略preactions中定义的自定义表结构。 - preactions与写入schema不匹配:
preactions创建了正确结构的表后,后续写入空DataFrame时因schema不匹配触发异常,最终导致表结构被错误覆盖。
解决方案
1. 移除提前的job.commit()
删除表创建逻辑前的job.commit()调用,确保作业能完整执行所有代码。
2. 直接执行CREATE TABLE语句
无需通过写入空DataFrame来触发preactions,直接使用Spark JDBC连接执行自定义的CREATE TABLE语句,避免schema覆盖问题。
3. 确保后续数据写入的schema匹配
如果需要写入数据,保证读取的DynamicFrame/DataFrame的schema与目标表结构一致。
修正后的完整代码
import sys import os from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job from awsglue.dynamicframe import DynamicFrame ## @params: [JOB_NAME] job_name = os.environ.get('AWS_GLUE_JOB_NAME', 'testing1') args = getResolvedOptions(sys.argv, ['JOB_NAME','TempDir']) temp_dir = args['TempDir'] sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args['JOB_NAME'], args) redshift_url = "jdbc:redshift://flighttest.ckbw5q2qccz6.eu-north-1.redshift.amazonaws.com:5439/dev" redshift_user = "awsuser" redshift_password = "Prateek1997" redshift_table = "flight.flight_details" create_table_sql = f""" CREATE TABLE IF NOT EXISTS {redshift_table}( Carrier varchar(100), OriginAirportID INT, DestAirportID INT, DepDelay INT, ArrDelay INT ) DISTSTYLE KEY DISTKEY(OriginAirportID); """ try: # 直接通过JDBC执行CREATE TABLE语句 spark.read.format("jdbc") \ .option("url", redshift_url) \ .option("dbtable", f"({create_table_sql}) AS temp") \ .option("user", redshift_user) \ .option("password", redshift_password) \ .option("driver", "com.amazon.redshift.jdbc42.Driver") \ .load() print("表结构创建成功") # 读取数据并验证(可选) database_name = "flight" table_name = "prateekproject1" df = glueContext.create_dynamic_frame.from_catalog( database = database_name, table_name = table_name ) df.printSchema() dataframe = df.toDF() print("Sample Data:") dataframe.show(10) # 如果需要写入数据,确保schema匹配后执行 # glueContext.write_dynamic_frame.from_options( # frame=df, # connection_type="redshift", # connection_options={ # "url": redshift_url, # "user": redshift_user, # "password": redshift_password, # "dbtable": redshift_table, # "redshiftTmpDir": temp_dir # } # ) except Exception as e: print(f"错误信息: {e}") job.commit()
额外说明
- 如果需要写入数据,取消注释代码中的写入逻辑,并确保读取的DataFrame列名、数据类型与Redshift表完全匹配。
- 确保Glue作业角色拥有Redshift的访问权限以及临时S3目录的读写权限。
内容的提问来源于stack exchange,提问作者Prateek Goel
相关产品推荐
相关产品推荐

