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 ?? 
3 REPLIES 3

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).