Sunday
Sunday
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
Sunday
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.
Monday
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).
10 hours ago
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.
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:
isEmpty() guard matters. Without it, a run with no new files would overwrite your table with an empty result._metadata.file_modification_time (or _metadata.file_path) before writing.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.
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.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.
[1] Ingest data from cloud object storage | Databricks on AWS โ https://docs.databricks.com/aws/en/ingestion/cloud-object-storage