Databricks streaming job failed appending to existing delta table.

Lam, Hoanh 0 Reputation points
2026-09-17T09:28:09.5033333+00:00

We run a databricks notebook task taking a number of streaming data sources and appending to a delta table as part of our data processing pipeline.

The streaming job has been running for some time when a batch failed. The reason for the failure is that the stream tried to create instead of appending to the existing table.

Stack trace as follows:

File "/databricks/spark/python/lib/py4j-0.10.9.9-src.zip/py4j/clientserver.py", line 644, in _call_proxy return_value = getattr(self.pool[obj_id], method)(*params) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/databricks/spark/python/pyspark/sql/utils.py", line 176, in call raise e File "/databricks/spark/python/pyspark/sql/utils.py", line 173, in call self.func(DataFrame(jdf, wrapped_session_jdf), batch_id) File "/databricks/python/lib/python3.12/site-packages/i4/workflow/tasks/sink/sink.py", line 146, in <lambda> return lambda df, epoch_id: self.process_batch(df, epoch_id) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/Workspace/Repos/.internal/4af487ee06_commits/cc2d8aeb7cf39eeb69f92c7ee6fdbed6586349e8/src/workflow/tasks/importer/sink.py", line 23, in process_batch df.write.mode('append').format('delta').saveAsTable(<TABLE>) File "/databricks/spark/python/pyspark/databricks/instrumentation/instrumentation_utils.py", line 217, in wrapper res = func(*args, **kwargs) ^^^^^^^^^^^^^^^^^^^^^ File "/databricks/spark/python/pyspark/sql/readwriter.py", line 1876, in saveAsTable self._jwrite.saveAsTable(name) File "/databricks/spark/python/lib/py4j-0.10.9.9-src.zip/py4j/java_gateway.py", line 1362, in call return_value = get_return_value( ^^^^^^^^^^^^^^^^^ File "/databricks/spark/python/pyspark/errors/exceptions/captured.py", line 275, in deco raise converted from None pyspark.errors.exceptions.captured.AnalysisException: [TABLE_OR_VIEW_ALREADY_EXISTS] Cannot create table or view because it already exists. Choose a different name, drop the existing object, add the IF NOT EXISTS clause to tolerate pre-existing objects, add the OR REPLACE clause to replace the existing materialized view, or add the OR REFRESH clause to refresh the existing streaming table. SQLSTATE: 42P07 at py4j.Protocol.getReturnValue(Protocol.java:476) at py4j.reflection.PythonProxyHandler.invoke(PythonProxyHandler.java:108) at jdk.proxy7/jdk.proxy7.$Proxy182.call(Unknown Source) at org.apache.spark.sql.execution.streaming.sources.PythonForeachBatchHelper$.$anonfun$callForeachBatch$1(ForeachBatchSink.scala:456) at org.apache.spark.sql.execution.streaming.sources.PythonForeachBatchHelper$.$anonfun$callForeachBatch$1$adapted(ForeachBatchSink.scala:456) at org.apache.spark.sql.execution.streaming.sources.ForeachBatchSink.callBatchWriter(ForeachBatchSink.scala:181) at org.apache.spark.sql.execution.streaming.sources.ForeachBatchSink.addBatchOptimized(ForeachBatchSink.scala:309) at org.apache.spark.sql.execution.streaming.sources.ForeachBatchSink.$anonfun$addBatch$2(ForeachBatchSink.scala:106) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at org.apache.spark.util.Utils$.timeTakenMs(Utils.scala:593) at org.apache.spark.sql.execution.streaming.sources.ForeachBatchSink.addBatch(ForeachBatchSink.scala:97) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.addBatch(MicroBatchExecution.scala:1314) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runBatch$20(MicroBatchExecution.scala:1589) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at org.apache.spark.sql.execution.streaming.ProgressContext.reportTimeTaken(ProgressReporter.scala:328) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.markAndTimeCollectBatch(MicroBatchExecution.scala:1322) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runBatch$19(MicroBatchExecution.scala:1589) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId0$11(SQLExecution.scala:568) at com.databricks.sql.util.MemoryTrackerHelper.withMemoryTracking(MemoryTrackerHelper.scala:80) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId0$10(SQLExecution.scala:494) at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:895) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId0$1(SQLExecution.scala:423) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:1512) at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId0(SQLExecution.scala:282) at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:832) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runBatch$18(MicroBatchExecution.scala:1579) at org.apache.spark.sql.execution.streaming.ProgressContext.reportTimeTaken(ProgressReporter.scala:328) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runBatch(MicroBatchExecution.scala:1579) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$executeOneBatch$6(MicroBatchExecution.scala:829) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.handleDataSourceException(MicroBatchExecution.scala:2218) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$executeOneBatch$5(MicroBatchExecution.scala:829) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.withSchemaEvolution(MicroBatchExecution.scala:2163) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$executeOneBatch$4(MicroBatchExecution.scala:825) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at org.apache.spark.sql.execution.streaming.ProgressContext.reportTimeTaken(ProgressReporter.scala:328) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$executeOneBatch$3(MicroBatchExecution.scala:780) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at com.databricks.logging.AttributionContextTracing.$anonfun$withAttributionContext$1(AttributionContextTracing.scala:49) at com.databricks.logging.AttributionContext$.$anonfun$withValue$1(AttributionContext.scala:293) at scala.util.DynamicVariable.withValue(DynamicVariable.scala:62) at com.databricks.logging.AttributionContext$.withValue(AttributionContext.scala:289) at com.databricks.logging.AttributionContextTracing.withAttributionContext(AttributionContextTracing.scala:47) at com.databricks.logging.AttributionContextTracing.withAttributionContext$(AttributionContextTracing.scala:44) at com.databricks.spark.util.PublicDBLogging.withAttributionContext(DatabricksSparkUsageLogger.scala:29) at com.databricks.logging.AttributionContextTracing.withAttributionTags(AttributionContextTracing.scala:96) at com.databricks.logging.AttributionContextTracing.withAttributionTags$(AttributionContextTracing.scala:77) at com.databricks.spark.util.PublicDBLogging.withAttributionTags(DatabricksSparkUsageLogger.scala:29) at com.databricks.spark.util.PublicDBLogging.withAttributionTags0(DatabricksSparkUsageLogger.scala:108) at com.databricks.spark.util.DatabricksSparkUsageLogger.withAttributionTags(DatabricksSparkUsageLogger.scala:215) at com.databricks.spark.util.UsageLogging.$anonfun$withAttributionTags$1(UsageLogger.scala:668) at com.databricks.spark.util.UsageLogging$.withAttributionTags(UsageLogger.scala:780) at com.databricks.spark.util.UsageLogging$.withAttributionTags(UsageLogger.scala:789) at com.databricks.spark.util.UsageLogging.withAttributionTags(UsageLogger.scala:668) at com.databricks.spark.util.UsageLogging.withAttributionTags$(UsageLogger.scala:666) at org.apache.spark.sql.execution.streaming.StreamExecution.withAttributionTags(StreamExecution.scala:92) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.executeOneBatch(MicroBatchExecution.scala:774) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStreamWithListener$1(MicroBatchExecution.scala:735) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStreamWithListener$1$adapted(MicroBatchExecution.scala:735) at org.apache.spark.sql.execution.streaming.TriggerExecutor.runOneBatch(TriggerExecutor.scala:85) at org.apache.spark.sql.execution.streaming.TriggerExecutor.runOneBatch$(TriggerExecutor.scala:73) at org.apache.spark.sql.execution.streaming.ProcessingTimeExecutor.runOneBatch(TriggerExecutor.scala:130) at org.apache.spark.sql.execution.streaming.ProcessingTimeExecutor.execute(TriggerExecutor.scala:143) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runActivatedStreamWithListener(MicroBatchExecution.scala:735) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runActivatedStream(MicroBatchExecution.scala:485) at org.apache.spark.sql.execution.streaming.StreamExecution.$anonfun$runStream$2(StreamExecution.scala:463) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:1512) at org.apache.spark.sql.execution.streaming.StreamExecution.$anonfun$runStream$1(StreamExecution.scala:402) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at com.databricks.logging.AttributionContextTracing.$anonfun$withAttributionContext$1(AttributionContextTracing.scala:49) at com.databricks.logging.AttributionContext$.$anonfun$withValue$1(AttributionContext.scala:293) at scala.util.DynamicVariable.withValue(DynamicVariable.scala:62) at com.databricks.logging.AttributionContext$.withValue(AttributionContext.scala:289) at com.databricks.logging.AttributionContextTracing.withAttributionContext(AttributionContextTracing.scala:47) at com.databricks.logging.AttributionContextTracing.withAttributionContext$(AttributionContextTracing.scala:44) at com.databricks.spark.util.PublicDBLogging.withAttributionContext(DatabricksSparkUsageLogger.scala:29) at com.databricks.logging.AttributionContextTracing.withAttributionTags(AttributionContextTracing.scala:96) at com.databricks.logging.AttributionContextTracing.withAttributionTags$(AttributionContextTracing.scala:77) at com.databricks.spark.util.PublicDBLogging.withAttributionTags(DatabricksSparkUsageLogger.scala:29) at com.databricks.spark.util.PublicDBLogging.withAttributionTags0(DatabricksSparkUsageLogger.scala:108) at com.databricks.spark.util.DatabricksSparkUsageLogger.withAttributionTags(DatabricksSparkUsageLogger.scala:215) at com.databricks.spark.util.UsageLogging.$anonfun$withAttributionTags$1(UsageLogger.scala:668) at com.databricks.spark.util.UsageLogging$.withAttributionTags(UsageLogger.scala:780) at com.databricks.spark.util.UsageLogging$.withAttributionTags(UsageLogger.scala:789) at com.databricks.spark.util.UsageLogging.withAttributionTags(UsageLogger.scala:668) at com.databricks.spark.util.UsageLogging.withAttributionTags$(UsageLogger.scala:666) at org.apache.spark.sql.execution.streaming.StreamExecution.withAttributionTags(StreamExecution.scala:92) at org.apache.spark.sql.execution.streaming.StreamExecution.org$apache$spark$sql$execution$streaming$StreamExecution$$runStream(StreamExecution.scala:374) at org.apache.spark.sql.execution.streaming.StreamExecution$$anon$1.$anonfun$run$3(StreamExecution.scala:292) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at org.apache.spark.JobArtifactSet$.withActiveJobArtifactState(JobArtifactSet.scala:97) at org.apache.spark.sql.execution.streaming.StreamExecution$$anon$1.$anonfun$run$2(StreamExecution.scala:292) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at com.databricks.unity.UCSEphemeralState$Handle.runWith(UCSEphemeralState.scala:51) at com.databricks.unity.HandleImpl.runWith(UCSHandle.scala:104) at com.databricks.unity.HandleImpl.$anonfun$runWithAndClose$1(UCSHandle.scala:109) at scala.util.Using$.resource(Using.scala:269) at com.databricks.unity.HandleImpl.runWithAndClose(UCSHandle.scala:108) at org.apache.spark.sql.execution.streaming.StreamExecution$$anon$1.run(StreamExecution.scala:291)

Azure Databricks
Azure Databricks

An Apache Spark-based analytics platform optimized for Azure.

0 comments No comments

1 answer

Sort by: Most helpful
  1. Allan Solomon Mejia 10,225 Reputation points
    2026-09-21T21:21:59.3766667+00:00

    Hi @Lam, Hoanh

    The stack trace is useful because the failure is occurring specifically at df.write.mode("append").format("delta").saveAsTable(<TABLE>) inside your foreachBatch callback.

    Normally, when <TABLE> already exists and resolves to a compatible Delta table, mode("append").saveAsTable() should append to it rather than attempt to create it. So don't drop or recreate the table based on this exception.

    First, check whether the table name resolves consistently inside the streaming job. If you're using an unqualified name such as:

    saveAsTable("my_table")
    

    change the test to the fully qualified Unity Catalog name:

    target = "catalog.schema.my_table"
    df.write \
      .format("delta") \
      .mode("append") \
      .saveAsTable(target)
    

    Then, immediately before the write inside process_batch, log what Spark sees:

    target = "catalog.schema.my_table"
    print("catalog =", spark.catalog.currentCatalog())
    print("database =", spark.catalog.currentDatabase())
    print("tableExists =", spark.catalog.tableExists(target))
    spark.sql(f"DESCRIBE DETAIL {target}").show(truncate=False)
    df.write \
      .format("delta") \
      .mode("append") \
      .saveAsTable(target)
    

    This matters because your exception isn't a typical schema-mismatch error. Databricks is reporting:

    [TABLE_OR_VIEW_ALREADY_EXISTS]

    Cannot create table or view because it already exists.

    SQLSTATE: 42P07

    Also, check the table history around the exact time the micro-batch failed:

    DESCRIBE HISTORY catalog.schema.my_table;
    

    Look for concurrent DDL or lifecycle operations such as CREATE, CREATE OR REPLACE, DROP, schema changes, or another job/pipeline manipulating the same target. If another process changed the catalog object while the stream was running, that would be important evidence.

    Also confirm what object Spark believes <TABLE> actually is:

    DESCRIBE EXTENDED catalog.schema.my_table;
    

    In particular, verify that it is the expected Delta table, not a view, streaming table, materialized view, or an object whose ownership/lifecycle is managed by another pipeline.

    Don't add IF NOT EXISTS, OR REPLACE, or automatically drop the table just because those suggestions appear in the generic exception text. Those are remedies for explicit table-creation operations and don't explain why an established append stream unexpectedly entered a create path.

    Since you mentioned that the streaming job had been running successfully for some time before one batch failed, could you also provide:

    • Databricks Runtime version
    • Dedicated/shared/serverless compute
    • Unity Catalog or Hive metastore
    • Whether <TABLE> is fully qualified
    • Whether any other job/pipeline writes to or modifies the same table

    If tableExists() returns True, DESCRIBE DETAIL identifies the expected Delta table, there was no concurrent DDL, and the same saveAsTable(..., mode="append") code intermittently enters the create path after previously succeeding, then capture the failing job/run ID, cluster/compute ID, UTC timestamp, full exception, target table name, and Delta history and raise it with Azure Databricks Support rather than work around it by recreating the table.

    One more point: because this is inside foreachBatch, make sure the batch write is idempotent before manually replaying/restarting failed batches. Databricks notes that foreachBatch provides at-least-once write guarantees by default, so a restarted/replayed batch needs appropriate handling if duplicate writes would be a problem.

    References:

    Databricks - Use foreachBatch to write to arbitrary data sinks

    Databricks - Delta table history


    Help make this community better for everyone: If this answer helped or resolved your issue, please accept it or upvote it. If not, share more details in a comment so we can continue the discussion and find the right solution. Thank you.

    Was this answer helpful?

    0 comments No comments

Your answer

Answers can be marked as 'Accepted' by the question author and 'Recommended' by moderators, which helps users know the answer solved the author's problem.