在Azure Synapse Notebook删除Delta Lake数据时遇java.lang.StackOverflowError
问题:Azure Synapse Delta表删除操作触发java.lang.StackOverflowError
问题背景
在Azure Synapse的Notebook中实现数据仓库事实表填充逻辑:处理源数据后,基于两列匹配源端与目标端的共同记录,存在则先删除目标端所有匹配记录,再插入新记录。该Notebook为多客户通用,每天至少由管道触发4次。此前出现该错误时重新运行即可解决,但此次已持续2天无法恢复。
错误信息
An error occurred while calling o5686.delete. : org.apache.spark.SparkException: Job aborted. at org.apache.spark.sql.errors.QueryExecutionErrors$.jobAbortedError(QueryExecutionErrors.scala:651) at org.apache.spark.sql.execution.datasources.FileFormatWriter$.write(FileFormatWriter.scala:284) at org.apache.spark.sql.delta.files.TransactionalWrite.$anonfun$writeFiles$1(TransactionalWrite.scala:456) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$6(SQLExecution.scala:111) at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:183) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:97) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:779) at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:66) at org.apache.spark.sql.delta.files.TransactionalWrite.writeFiles(TransactionalWrite.scala:391) at org.apache.spark.sql.delta.files.TransactionalWrite.writeFiles$(TransactionalWrite.scala:355) at org.apache.spark.sql.delta.OptimisticTransaction.writeFiles(OptimisticTransaction.scala:103) at org.apache.spark.sql.delta.files.TransactionalWrite.writeFiles(TransactionalWrite.scala:215) at org.apache.spark.sql.delta.files.TransactionalWrite.writeFiles$(TransactionalWrite.scala:212) at org.apache.spark.sql.delta.OptimisticTransaction.writeFiles(OptimisticTransaction.scala:103) at org.apache.spark.sql.delta.files.TransactionalWrite.writeFiles(TransactionalWrite.scala:232) at org.apache.spark.sql.delta.files.TransactionalWrite.writeFiles$(TransactionalWrite.scala:231) at org.apache.spark.sql.delta.OptimisticTransaction.writeFiles(OptimisticTransaction.scala:103) at org.apache.spark.sql.delta.commands.DeleteCommand.$anonfun$rewriteFiles$2(DeleteCommand.scala:395) at org.apache.spark.sql.delta.util.DeltaProgressReporter.withJobDescription(DeltaProgressReporter.scala:53) at org.apache.spark.sql.delta.util.DeltaProgressReporter.withStatusCode(DeltaProgressReporter.scala:32) at org.apache.spark.sql.delta.util.DeltaProgressReporter.withStatusCode$(DeltaProgressReporter.scala:27) at org.apache.spark.sql.delta.commands.DeleteCommand.withStatusCode(DeleteCommand.scala:99) at org.apache.spark.sql.delta.commands.DeleteCommand.rewriteFiles(DeleteCommand.scala:375) at org.apache.spark.sql.delta.commands.DeleteCommand.performDelete(DeleteCommand.scala:275) at org.apache.spark.sql.delta.commands.DeleteCommand.$anonfun$run$2(DeleteCommand.scala:115) at org.apache.spark.sql.delta.DeltaLog.withNewTransaction(DeltaLog.scala:262) at org.apache.spark.sql.delta.commands.DeleteCommand.$anonfun$run$1(DeleteCommand.scala:114) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at org.apache.spark.sql.delta.metering.DeltaLogging.recordFrameProfile(DeltaLogging.scala:140) at org.apache.spark.sql.delta.metering.DeltaLogging.recordFrameProfile$(DeltaLogging.scala:138) at org.apache.spark.sql.delta.commands.DeleteCommand.recordFrameProfile(DeleteCommand.scala:99) at org.apache.spark.sql.delta.metering.DeltaLogging.$anonfun$recordDeltaOperationInternal$1(DeltaLogging.scala:133) at com.microsoft.spark.telemetry.delta.SynapseLoggingShim.recordOperation(SynapseLoggingShim.scala:95) at com.microsoft.spark.telemetry.delta.SynapseLoggingShim.recordOperation$(SynapseLoggingShim.scala:81) at org.apache.spark.sql.delta.commands.DeleteCommand.recordOperation(DeleteCommand.scala:99) at org.apache.spark.sql.delta.metering.DeltaLogging.recordDeltaOperationInternal(DeltaLogging.scala:132) at org.apache.spark.sql.delta.metering.DeltaLogging.recordDeltaOperation(DeltaLogging.scala:122) at org.apache.spark.sql.delta.metering.DeltaLogging.recordDeltaOperation$(DeltaLogging.scala:110) at org.apache.spark.sql.delta.commands.DeleteCommand.recordDeltaOperation(DeleteCommand.scala:99) at org.apache.spark.sql.delta.commands.DeleteCommand.run(DeleteCommand.scala:112) at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult$lzycompute(commands.scala:75) at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult(commands.scala:73) at org.apache.spark.sql.execution.command.ExecutedCommandExec.executeCollect(commands.scala:84) at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.$anonfun$applyOrElse$1(QueryExecution.scala:152) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$6(SQLExecution.scala:111) at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:183) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:97) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:779) at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:66) at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.applyOrElse(QueryExecution.scala:152) at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.applyOrElse(QueryExecution.scala:145) at org.apache.spark.sql.catalyst.trees.TreeNode.$anonfun$transformDownWithPruning$1(TreeNode.scala:584) at org.apache.spark.sql.catalyst.trees.CurrentOrigin$.withOrigin(TreeNode.scala:176) at org.apache.spark.sql.catalyst.trees.TreeNode.transformDownWithPruning(TreeNode.scala:584) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.org$apache$spark$sql$catalyst$plans$logical$AnalysisHelper$$super$transformDownWithPruning(LogicalPlan.scala:31) at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDownWithPruning(AnalysisHelper.scala:267) at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDownWithPruning$(AnalysisHelper.scala:263) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDownWithPruning(LogicalPlan.scala:31) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDownWithPruning(LogicalPlan.scala:31) at org.apache.spark.sql.catalyst.trees.TreeNode.transformDown(TreeNode.scala:560) at org.apache.spark.sql.execution.QueryExecution.eagerlyExecuteCommands(QueryExecution.scala:145) at org.apache.spark.sql.execution.QueryExecution.commandExecuted$lzycompute(QueryExecution.scala:129) at org.apache.spark.sql.execution.QueryExecution.commandExecuted(QueryExecution.scala:123) at org.apache.spark.sql.Dataset.<init>(Dataset.scala:231) at org.apache.spark.sql.Dataset$.$anonfun$ofRows$1(Dataset.scala:93) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:779) at org.apache.spark.sql.Dataset$.ofRows(Dataset.scala:90) at org.apache.spark.sql.delta.util.AnalysisHelper.toDataset(AnalysisHelper.scala:88) at org.apache.spark.sql.delta.util.AnalysisHelper.toDataset$(AnalysisHelper.scala:87) at io.delta.tables.DeltaTable.toDataset(DeltaTable.scala:44) at io.delta.tables.execution.DeltaTableOperations.$anonfun$executeDelete$1(DeltaTableOperations.scala:45) at org.apache.spark.sql.delta.util.AnalysisHelper.improveUnsupportedOpError(AnalysisHelper.scala:105) at org.apache.spark.sql.delta.util.AnalysisHelper.improveUnsupportedOpError$(AnalysisHelper.scala:91) at io.delta.tables.DeltaTable.improveUnsupportedOpError(DeltaTable.scala:44) at io.delta.tables.execution.DeltaTableOperations.executeDelete(DeltaTableOperations.scala:41) at io.delta.tables.execution.DeltaTableOperations.executeDelete$(DeltaTableOperations.scala:41) at io.delta.tables.DeltaTable.executeDelete(DeltaTable.scala:44) at io.delta.tables.DeltaTable.delete(DeltaTable.scala:234) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357) at py4j.Gateway.invoke(Gateway.java:282) at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) at py4j.commands.CallCommand.execute(CallCommand.java:79) at py4j.GatewayConnection.run(GatewayConnection.java:238) at java.lang.Thread.run(Thread.java:750) Caused by: java.lang.StackOverflowError at org.codehaus.janino.CodeContext.flowAnalysis(CodeContext.java:535) at org.codehaus.janino.CodeContext.flowAnalysis(CodeContext.java:600) at org.codehaus.janino.CodeContext.flowAnalysis(CodeContext.java:600) at org.codehaus.janino.CodeContext.flowAnalysis(CodeContext.java:600) at org.codehaus.janino.CodeContext.flowAnalysis(CodeContext.java:600) at org.codehaus.janino.CodeContext.flowAnalysis(CodeContext.java:600) at org.codehaus.janino.CodeContext.flowAnalysis(CodeContext.java:600) at org.codehaus.janino.CodeContext.flowAnalysis(CodeContext.java:600) at org.codehaus.janino.CodeContext.flowAnalysis(CodeContext.java:600) at org.codehaus.janino.CodeContext.flowAnalysis(CodeContext.java:600) at org.codehaus.janino.CodeContext.flowAnalysis(CodeContext.java:600) at org.codehaus.janino.CodeContext.flowAnalysis(CodeContext.java:600) . . .
已尝试的解决方案
- 缓存数据
- 添加检查点(处理时间从约20分钟增至1小时45分钟)
- 将Spark集群池的智能缓存大小设为0%
- 减少单次删除批处理的过滤条件数量,从5000降至1500再到100
上述方案单独或组合尝试后,仍触发相同错误。
核心删除代码
通过循环获取参数调用以下方法执行删除:
def deleteRecordDelta(self, sql_f_str): deltaTable = DeltaTable.forPath(self.spark, self.path) deltaTable.delete(sql_f_str)
补充说明
预期逻辑为:源端记录与目标端多条记录匹配时,先删除目标端匹配记录,再插入新记录。除删除操作外,其他环节均正常,但删除操作偶尔失败,且每次失败均为该java.lang.StackOverflowError。
由于用于关联的列是多对多类型,使用Merge操作会触发Spark冲突错误,因此必须采用先删后插的方式。
内容的提问来源于stack exchange,提问作者Yash Varshney
相关产品推荐
相关产品推荐

