Options
- Mark as New
- Bookmark
- Subscribe
- Mute
- Subscribe to RSS Feed
- Permalink
- Report Inappropriate Content
01-19-2022 09:27 AM
def upsertToDelta(microBatchOutputDF, batchId):
microBatchOutputDF.createOrReplaceTempView("updates")
microBatchOutputDF._jdf.sparkSession().sql("""
MERGE INTO old o
USING updates u
ON u.id = o.id
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *
""")
stream_new_df = spark.readStream.format("delta").load(new_data_frame_path)
stream_old_df = spark.readStream.format("delta").load(old_data_frame_path)
stream_old_df.createOrReplaceTempView("old")
stream_new_df.writeStream.format("delta") \
.option("checkpointLocation", "") \
.option("mergeSchema", "true") \
.option("path", "") \
.foreachBatch(upsertToDelta) \
.trigger(once=True) \
.outputMode("update") \
.table("")I'm trying to execute this code but I get the following error:
Data source com.databricks.sql.transaction.tahoe.sources.DeltaDataSource does not support Update output mode