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内部错误,仅返回堆栈跟踪信息。
操作步骤
- 通过JDBC创建临时视图关联RDS表,查询验证正常;
- 调用API获取JSON响应;
- 转换响应为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()) ]) )
- 创建临时视图
staging_test,查询验证正常; - 执行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,网络已打通。
排查步骤
- 检查字段映射一致性:对比DataFrame定义与MERGE语句,发现USING子句中引用的
amount字段在staging_test视图中不存在(原DataFrame字段为loanAmount),字段不存在会导致Spark SQL解析阶段出现异常,触发内部错误。 - 验证临时视图结构:执行
DESCRIBE staging_test确认视图的字段名和数据类型,确保MERGE语句中引用的字段均存在。 - 排查JDBC兼容性:确认Databricks自带的PostgreSQL JDBC驱动版本是否适配RDS PostgreSQL 11.19,跨区域网络延迟是否可能影响元数据获取。
- 简化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
相关产品推荐
相关产品推荐

