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

AWS EMR环境下Iceberg表MERGE INTO操作报错求助

解决AWS EMR 6.15.0中Iceberg表MERGE INTO不支持的问题

问题背景

在AWS EMR 6.15.0环境的Jupyter Notebook中,对Iceberg表执行MERGE INTO操作时触发错误:MERGE INTO TABLE is not supported temporarily.
环境版本:

  • Spark 3.4.1-amzn-2
  • Hadoop 3.3.6
  • Hive 3.1.3
  • Scala 2.12.15
  • Iceberg 1.4.0

以下是具体的解决步骤:


1. 统一Iceberg依赖版本,消除冲突

EMR自带的Iceberg包与手动指定的依赖存在版本冲突,需统一依赖来源:

  • 移除spark.jars中本地Iceberg包路径,改用spark.jars.packages统一拉取匹配版本的依赖
  • 修正Scala版本格式,仅需指定主版本(2.12),无需带小版本(15)

修改后的Cell 1配置代码:

%%configure -f
{
    "conf": {
        "spark.jars.packages": "com.microsoft.azure:spark-mssql-connector_2.12:1.2.0,org.apache.iceberg:iceberg-spark-runtime-3.4_2.12:1.4.0",
        "spark.submit.pyFiles": "s3://dev-datalake-us-west-2-apps/ci-artifacts/utility_files/utils.py,s3://dev-datalake-us-west-2-apps/ci-artifacts/utility_files/eventlogs_utils.py,s3://dev-datalake-us-west-2-apps/ci-artifacts/utility_files/flowlogs_utils.py"
    }
}

from pyspark.sql import SparkSession
from pyspark.sql.functions import current_timestamp, unix_timestamp
from pyspark.sql import Row
from pyspark.sql import functions as F
import utils as u
from datetime import datetime, timedelta

spark = SparkSession.builder \
    .appName("spark job") \
    .config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.iceberg.spark.SparkSessionCatalog") \
    .config("spark.sql.catalog.spark_catalog.type", "hive") \
    .config("spark.sql.catalog.local", "org.apache.iceberg.spark.SparkCatalog") \
    .config("spark.sql.catalog.local.type", "hadoop") \
    .config("spark.sql.catalog.local.warehouse", "hdfs:///tmp") \
    .config("spark.hadoop.hive.cli.print.header", "true") \
    .config("spark.sql.caseSensitive", "true") \
    .config("spark.sql.legacy.timeParserPolicy", "LEGACY") \
    .config("spark.sql.parquet.int96RebaseModeInWrite", "LEGACY") \
    .config("spark.sql.sources.partitionOverwriteMode", "dynamic") \
    .enableHiveSupport() \
    .getOrCreate()

2. 确认Iceberg扩展已正确加载

spark.sql.extensions配置必须指向Iceberg的扩展类,这是启用Iceberg专属SQL语法(包括MERGE INTO)的核心。如果扩展未加载,Spark会用自身不支持MERGE的逻辑处理Iceberg表,导致报错。

3. 调整MERGE INTO语法兼容性

Iceberg 1.4.0支持MERGE INTO,但需确保语法符合规范:

  • 若表名无特殊字符,可去掉反引号简化语句;若保留反引号需确保包裹完整
  • Iceberg 1.4.0已支持通配符SET *和INSERT *,无需逐个指定字段

简化后的MERGE INTO语句:

MERGE INTO local.default.iceberg_table t
USING records_data_view u
ON t.cid = u.cid
AND t.sid = u.sid
AND t.satid = u.satid
AND t.t_str = u.t_str
AND t.fznum = u.fznum
AND t.fzid = u.fzid
AND t.time = u.time
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *

4. 验证表的Iceberg格式

执行以下查询确认目标表为Iceberg格式:

DESCRIBE EXTENDED local.default.iceberg_table

查看输出中的Table Type字段,若值为ICEBERG则格式正确;若不是,需重新创建表。

5. 重启Jupyter Notebook会话

修改Spark配置后,必须重启Notebook会话,确保新依赖和配置完全生效,避免旧配置残留导致问题。


内容的提问来源于stack exchange,提问作者shady xv

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 09:42:33