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

Databricks中Spark SQL MERGE语句触发内部错误求助

Spark SQL MERGE INTO 触发INTERNAL_ERROR的排查与解决(Databricks + AWS RDS PostgreSQL场景)

问题场景

通过Python requests获取API数据,转换为PySpark DataFrame后,使用Spark SQL的MERGE语句向跨VPC的AWS RDS PostgreSQL表执行Upsert操作时,触发Spark内部错误,仅返回堆栈跟踪信息。

操作步骤

  1. 通过JDBC创建临时视图关联RDS表,查询验证正常;
  2. 调用API获取JSON响应;
  3. 转换响应为PySpark DataFrame:
df = spark.createDataFrame(data=[(response['ID'],int(response['clientID']),float(response['amount']))],
    schema=StructType(fields=[
    StructField('ID',IntegerType()),
    StructField('clientID',IntegerType()),
    StructField('loanAmount',FloatType())
])
    )
  1. 创建临时视图staging_test,查询验证正常;
  2. 执行MERGE语句时触发错误:
MERGE INTO test whs
USING (SELECT ID as loan_id,clientID as account_id,amount as amount FROM staging_test) staging
ON whs.loan_id=staging.loan_id
WHEN MATCHED THEN
  UPDATE SET whs.amount=staging.amount
WHEN NOT MATCHED THEN
 INSERT (loan_id,account_id,amount)
  VALUES (staging.loan_id,staging.account_id,staging.amount)

错误信息

Error in SQL statement: SparkException: [INTERNAL_ERROR] The Spark SQL phase planning failed with an internal error. You hit a bug in Spark or the Spark plugins you use. Please, report this bug to the corresponding communities or vendors, and provide the full stack trace.
com.databricks.backend.common.rpc.SparkDriverExceptions$SQLExecutionException: org.apache.spark.SparkException: [INTERNAL_ERROR] The Spark SQL phase planning failed with an internal error. You hit a bug in Spark or the Spark plugins you use. Please, report this bug to the corresponding communities or vendors, and provide the full stack trace.
    at org.apache.spark.SparkException$.internalError(SparkException.scala:100)
    at org.apache.spark.sql.execution.QueryExecution$.toInternalError(QueryExecution.scala:925)
    at org.apache.spark.sql.execution.QueryExecution$.withInternalError(QueryExecution.scala:937)
    at org.apache.spark.sql.execution.QueryExecution.$anonfun$executePhase$2(QueryExecution.scala:431)
    at com.databricks.util.LexicalThreadLocal$Handle.runWith(LexicalThreadLocal.scala:63)
    at org.apache.spark.sql.execution.QueryExecution.$anonfun$executePhase$1(QueryExecution.scala:427)
    at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:1113)
    at org.apache.spark.sql.execution.QueryExecution.executePhase(QueryExecution.scala:427)
    at org.apache.spark.sql.execution.QueryExecution.sparkPlan$lzycompute(QueryExecution.scala:365)
    at org.apache.spark.sql.execution.QueryExecution.sparkPlan(QueryExecution.scala:358)
    at org.apache.spark.sql.execution.QueryExecution.$anonfun$executedPlan$1(QueryExecution.scala:379)
    at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:1113)
    at org.apache.spark.sql.execution.QueryExecution.executedPlan$lzycompute(QueryExecution.scala:379)
    at org.apache.spark.sql.execution.QueryExecution.executedPlan(QueryExecution.scala:374)
    at org.apache.spark.sql.execution.QueryExecution.simpleString(QueryExecution.scala:485)
    at org.apache.spark.sql.execution.QueryExecution.org$apache$spark$sql$execution$QueryExecution$$explainString(QueryExecution.scala:550)
    at org.apache.spark.sql.execution.QueryExecution.explainStringLocal(QueryExecution.scala:512)
    at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withCustomExecutionEnv$8(SQLExecution.scala:225)
    at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:498)

环境信息

  • Databricks:单节点m6i.large,版本13.3 LTS(Spark 3.4.1、Scala 2.12),未启用Photon;
  • AWS RDS:跨区域跨VPC PostgreSQL 11.19,实例db.m5.large,网络已打通。

排查步骤

  1. 检查字段映射一致性:对比DataFrame定义与MERGE语句,发现USING子句中引用的amount字段在staging_test视图中不存在(原DataFrame字段为loanAmount),字段不存在会导致Spark SQL解析阶段出现异常,触发内部错误。
  2. 验证临时视图结构:执行DESCRIBE staging_test确认视图的字段名和数据类型,确保MERGE语句中引用的字段均存在。
  3. 排查JDBC兼容性:确认Databricks自带的PostgreSQL JDBC驱动版本是否适配RDS PostgreSQL 11.19,跨区域网络延迟是否可能影响元数据获取。
  4. 简化MERGE语句测试:移除USING子句中的子查询,直接使用staging_test视图,观察是否仍触发错误,定位是否为子查询解析问题。

解决方法

1. 修复字段映射错误

将MERGE语句中USING子句的amount替换为loanAmount,修正后的语句如下:

MERGE INTO test whs
USING (SELECT ID as loan_id, clientID as account_id, loanAmount as amount FROM staging_test) staging
ON whs.loan_id = staging.loan_id
WHEN MATCHED THEN
  UPDATE SET whs.amount = staging.amount
WHEN NOT MATCHED THEN
 INSERT (loan_id, account_id, amount)
  VALUES (staging.loan_id, staging.account_id, staging.amount)

2. 用PySpark DataFrame API替代MERGE

若MERGE语句仍触发内部错误,可使用DataFrame的join逻辑实现Upsert,绕开Spark SQL的MERGE解析问题:

# 读取目标RDS表
jdbc_url = "jdbc:postgresql://<rds-endpoint>:5432/<db-name>"
jdbc_properties = {
    "user": "<username>",
    "password": "<password>",
    "driver": "org.postgresql.Driver"
}
target_df = spark.read.jdbc(url=jdbc_url, table="test", properties=jdbc_properties)

# 执行全外连接合并数据,优先保留新数据
merged_df = target_df.join(df, target_df.loan_id == df.ID, "full_outer") \
    .selectExpr(
        "coalesce(target.loan_id, source.ID) as loan_id",
        "coalesce(target.account_id, source.clientID) as account_id",
        "coalesce(source.loanAmount, target.amount) as amount"
    )

# 写入目标表(根据业务需求选择mode,overwrite会覆盖全表,需谨慎)
merged_df.write.jdbc(url=jdbc_url, table="test", mode="overwrite", properties=jdbc_properties)

3. 升级Databricks版本

若确认是Spark 3.4.1的已知bug,可升级至Databricks 14.3 LTS及以上版本,新版本通常会修复此类内部解析错误。

4. 检查JDBC驱动配置

若跨区域网络存在不稳定情况,可添加JDBC连接参数(如connectTimeout、socketTimeout),并确认驱动版本与RDS PostgreSQL版本匹配。

内容的提问来源于stack exchange,提问作者Eugenio.Gastelum96

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 05:51:03