Glue 3.0中Spark SQL执行MERGE INTO报错,求解决方案
问题
使用支持Spark 3.1和Python 3的Glue 3.0,尝试通过Spark SQL的MERGE INTO target USING source实现UPSERT功能时,触发报错:
An error occurred while calling o91.sql. MERGE INTO TABLE is not supported temporarily.
场景说明:
- 目标表
target是通过Spark DataFrame Reader直接读取PostgreSQL AuroraDB得到的DataFrame - 源表
source是读取Parquet文件得到的DataFrame - 未使用Delta Table,更换Glue版本后问题依旧,查询资料发现多指向Iceberg和DeltaTable,想确认当前方案是否可行并寻求指导
附上代码:
def changeDataCapture(inputDf, currDf, spark): inputDf.createOrReplaceTempView('inputDf') currDf.createOrReplaceTempView('currDf') currDf = spark.sql(""" MERGE INTO currDf USING inputDf ON currDf.REG_NB = inputDf.registerNumber AND currDf.ANN_RTN_DT = inputDf.annual_return_date WHEN MATCHED THEN UPDATE SET currDf.LAST_SEEN_DT = inputDf.LAST_SEEN_DT, currDf.TO_DB_DT = inputDf.TO_DB_DT, currDf.TO_DB_TM = inputDf.TO_DB_TM, currDf.BATCH_ID = inputDf.BATCH_ID, currDf.DATA_PROC_ID = inputDf.DATA_PROC_ID, currDf.FIRST_SEEN_DT = CASE WHEN currDf.CO_REG_DEBT = inputDf.registered_indebtedness AND currDf.HLDR_LIST_CD = inputDf.holder_list_indicator AND currDf.HLDR_LEGAL_STAT = inputDf.holder_legal_status AND currDf.HLDR_REFRESH_CD = inputDf.holder_refresh_flag AND currDf.HLDR_SUPRESS_IN = inputDf.HLDR_SUPRESS_IN AND currDf.BULK_LIST_ID = inputDf.Bulk_List_In THEN currDf.FIRST_SEEN_DT ELSE inputDf.FIRST_SEEN_DT END, currDf.SUPERSEDED_DT = CASE WHEN currDf.CO_REG_DEBT = inputDf.registered_indebtedness AND currDf.HLDR_LIST_CD = inputDf.holder_list_indicator AND currDf.HLDR_LEGAL_STAT = inputDf.holder_legal_status AND currDf.HLDR_REFRESH_CD = inputDf.holder_refresh_flag AND currDf.HLDR_SUPRESS_IN = inputDf.HLDR_SUPRESS_IN AND currDf.BULK_LIST_ID = inputDf.Bulk_List_In THEN currDf.SUPERSEDED_DT ELSE inputDf.SUPERSEDED_DT END WHEN NOT MATCHED THEN INSERT (REG_NB, ANN_RTN_DT, SUPERSEDED_DT, TO_DB_DT, TO_DB_TM, FIRST_SEEN_DT, LAST_SEEN_DT, BATCH_ID, DATA_PROC_ID, CO_REG_DEBT, HLDR_LIST_CD, HLDR_LIST_DT, HLDR_LEGAL_STAT, HLDR_REFRESH_CD, HLDR_SUPRESS_IN, BULK_LIST_ID, DOC_TYPE_CD) VALUES (registerNumber, annual_return_date, SUPERSEDED_DT, TO_DB_DT, TO_DB_TM, FIRST_SEEN_DT, LAST_SEEN_DT, BATCH_ID, DATA_PROC_ID, registered_indebtedness, holder_list_indicator, holder_list_date, holder_legal_status, holder_refresh_flag, HLDR_SUPRESS_IN, Bulk_List_In, DOC_TYPE_CD) """) return currDf
结论与指导
当前方案不可行
Spark SQL原生的MERGE INTO语法在Glue 3.0中仅支持Delta Lake、Iceberg或Hudi这类支持ACID事务的湖仓表,直接对JDBC数据源(PostgreSQL AuroraDB)的临时视图执行MERGE INTO是不被支持的,这就是触发报错的核心原因。替代方案
- 方案一:引入湖仓表作为中间层
将PostgreSQL的目标数据同步到Delta Lake/Iceberg表中,对湖仓表和Parquet源表执行MERGE INTO完成UPSERT后,再将结果写回PostgreSQL。 - 方案二:手动实现UPSERT逻辑
通过Spark DataFrame API拆分操作:- 基于匹配条件找出需要更新的记录,生成更新数据集
- 找出源表中不存在于目标表的记录,生成插入数据集
- 分别对PostgreSQL执行更新和插入操作
- 方案三:利用PostgreSQL原生UPSERT
将Parquet源数据加载到PostgreSQL临时表,使用PostgreSQL原生的INSERT ... ON CONFLICT ... DO UPDATE语法完成UPSERT,性能更贴近数据库原生能力。
内容的提问来源于stack exchange,提问作者SDS
相关产品推荐
相关产品推荐

