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

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

结论与指导

  1. 当前方案不可行
    Spark SQL原生的MERGE INTO语法在Glue 3.0中仅支持Delta Lake、Iceberg或Hudi这类支持ACID事务的湖仓表,直接对JDBC数据源(PostgreSQL AuroraDB)的临时视图执行MERGE INTO是不被支持的,这就是触发报错的核心原因。

  2. 替代方案

  • 方案一:引入湖仓表作为中间层
    将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 17:10:36