cancel
Showing results forย 
Search instead forย 
Did you mean:ย 
Data Engineering
Join discussions on data engineering best practices, architectures, and optimization strategies within the Databricks Community. Exchange insights and solutions with fellow data engineers.
cancel
Showing results forย 
Search instead forย 
Did you mean:ย 

Databricks SDP (Spark declarative Pipelins) Overwrite table

pvrcloudtech
New Contributor
1) Consider I have orders folder which orders_1.csv file  and it is loaded to orders tables
 
orders/                 ============> load to order table 
     orders_1.csv 
 
2) 
next day new file (orders_2.csv) arrived
 
orders/                       =========> take only 2nd file ----->     overwrite the order table 
     orders_1.csv 
      orders_2.csv
 
 
Here streaming table only appending the data to orders table.  Materialized view is considering both files.
 
My requirement is overwrite the orders table. 
 
Can you help me how to do that ?? 
4 REPLIES 4

SumeshKashyap
New Contributor III

Streaming tables are append-only by design, so they can't overwrite. Use a materialized view that keeps only the rows from the newest file, via the _metadata column:

 

sql

CREATE OR REFRESH MATERIALIZED VIEW orders AS

SELECT * EXCEPT (file_time)

FROM (

  SELECT *,

         _metadata.file_modification_time AS file_time

  FROM read_files('/path/orders/', format => 'csv', header => true)

)

QUALIFY file_time = MAX(file_time) OVER ();

 

Each refresh then replaces the table contents with only the latest file (orders_2.csv, then orders_3.csv, and so on).

 

If the folder grows large, ingest with a streaming table (Auto Loader) into a bronze table, storing _metadata.file_name and _metadata.file_modification_time as columns, and build the same "latest file only" materialized view on top of it. That way all the files aren't re-read on every refresh.

Hope this helps

VK210287
New Contributor

A Streaming Table may not be the right fit for this requirement. It is designed to process new files incrementally, so when orders2.csv arrives, it will append the new data.

If each new file is a full snapshot and should completely replace the previous data, you can use a batch job to read only the latest file and overwrite the target Delta table using:

df.write.mode("overwrite").saveAsTable("orders")

Another option is to keep only the latest file in a  current/ folder and move older files to an archive/ folder.

So in this case, I would prefer a batch overwrite approach rather than a Streaming Table.

Vannurswamy Kuruba

AbhilashNagilla
Databricks Employee
Databricks Employee

If each file is a full snapshot, use AUTO CDC FROM SNAPSHOT with SCD Type 1; it processes files in numeric order and makes the target match the latest snapshot, including deleting absent keys (version-function example, snapshot API). Use this code in place of the current orders definition:

import re

from pyspark import pipelines as dp

ORDERS_DIR = "/Volumes/<catalog>/<schema>/<volume>/orders/"


def next_snapshot_and_version(latest_version):
    versions = sorted(
        int(m.group(1))
        for f in dbutils.fs.ls(ORDERS_DIR)
        if (m := re.fullmatch(r"orders_(\d+)\.csv", f.name))
    )
    newer = [v for v in versions if latest_version is None or v > latest_version]
    if not newer:
        return None
    df = spark.read.format("csv").option("header", True).load(f"{ORDERS_DIR}orders_{newer[0]}.csv")
    return df, newer[0]


dp.create_streaming_table("orders")
dp.create_auto_cdc_from_snapshot_flow(
    target="orders",
    source=next_snapshot_and_version,
    keys=["order_id"],
    stored_as_scd_type=1,
)

Set keys to the column or columns that uniquely identify a row (snapshot API). Without .schema(...), this CSV read returns every column as a string; add the schema before .load(...) for typed columns (CSV options, specify a schema). The flow requires serverless or the Pro or Advanced edition (requirements). For an existing streaming table, switch the code, then run a full refresh of only orders; for an existing materialized view, switch the code, run DROP MATERIALIZED VIEW, then run the same full refresh (type-change guidance, full-refresh action).

anuj_lathi
Databricks Employee
Databricks Employee

Streaming tables are append-only by design, and a materialized view recomputes over everything in the folder. Neither gives you "replace the table with just the newest file" out of the box, so you need a small pattern on top. Here are the two I'd consider.

Option 1: Auto Loader + foreachBatch overwrite (my preference)

Auto Loader tracks which files it has already processed and only picks up new ones on each run [1]. Instead of appending, you overwrite the target table with whatever each batch contains. That means the table always holds only the newly arrived file.

def overwrite_orders(batch_df, batch_id):
    if batch_df.isEmpty():
        return  # nothing new, leave the table untouched
    (batch_df.write
        .mode("overwrite")
        .saveAsTable("my_catalog.my_schema.orders"))

(spark.readStream
    .format("cloudFiles")
    .option("cloudFiles.format", "csv")
    .option("header", "true")
    .option("cloudFiles.schemaLocation", "/Volumes/my_catalog/my_schema/chk/orders_schema")
    .load("/Volumes/my_catalog/my_schema/landing/orders/")
  .writeStream
    .foreachBatch(overwrite_orders)
    .option("checkpointLocation", "/Volumes/my_catalog/my_schema/chk/orders")
    .trigger(availableNow=True)
    .start())

Notes:

  • Run this as a normal job or notebook task, scheduled daily or file-arrival triggered. Don't try to put it inside a declarative pipeline streaming table, because that is the append model you're already seeing.
  • The isEmpty() guard matters. Without it, a run with no new files would overwrite your table with an empty result.
  • If two files land between runs, one batch will contain both. If you need strictly "only the latest file", filter the batch on _metadata.file_modification_time (or _metadata.file_path) before writing.
  • Delta overwrite is atomic, so readers never see a half-empty table, and you keep time travel on the previous version.

Option 2: Materialized view that selects only the latest file

If you want to stay fully declarative, keep the materialized view but filter to the newest file using the file metadata column:

CREATE OR REFRESH MATERIALIZED VIEW orders AS
WITH src AS (
  SELECT *, _metadata.file_modification_time AS file_ts
  FROM read_files(
    '/Volumes/my_catalog/my_schema/landing/orders/',
    format => 'csv',
    header => true
  )
)
SELECT * EXCEPT (file_ts)
FROM src
WHERE file_ts = (SELECT max(file_ts) FROM src);

The trade-off is that the source folder is still scanned on refresh, so as old files pile up the cost grows. If you go this route, archive or delete processed files periodically.

What I'd avoid

  • TRUNCATE followed by COPY INTO. COPY INTO is idempotent and skips files it has already loaded [1], so it would load just the new file. But the two steps aren't atomic, and a run with no new file leaves you with an empty table.
  • Full refresh on the streaming table. It reprocesses everything in the folder, which is the opposite of what you want.

One more thing to check

If "overwrite" actually means that rows in the new file should replace matching rows (same order_id) rather than wipe the table, you don't want an overwrite at all. You want an upsert, either with MERGE INTO inside foreachBatch or with the declarative auto CDC flow in a pipeline. Let me know which semantics you need and I can sketch that version.

References

[1] Ingest data from cloud object storage | Databricks on AWS โ€” https://docs.databricks.com/aws/en/ingestion/cloud-object-storage

Anuj Lathi
Solutions Engineer @ Databricks