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: 

Change data feed from a materialized view

cdn_yyz_yul
Contributor III

Hi everyone,

Trying the CDF on materialized view (declarative pipeline) feature that is currently in Beta.

- Configuration:
** Verified the MVs (materialized veiws) have: 'delta.enableRowTracking'
SHOW TBLPROPERTIES my_mv ('delta.enableRowTracking');  returns true

** Enabled external metadata flag,

 In the declarative pipeline setting has included :
"configuration": {
"spark.databricks.acl.needAdminPermissionToViewLogs": "false",
"pipelines.externalMetadata.enabled": "true"
},

For the MV that used for testing, also added 

@dp.materialized_view(...) has:
    table_properties={
        "quality": "gold",
        "pipelines.reset.allowed": "true",
        "pipelines.externalMetadata.enabled": "true",
        "delta.enableRowTracking": "true",
    },

** DLT settings has channel: PREVIEW

- Read from CDF:
Dropped the MV used for the test: my_mv, and run full refresh of the pipeline. 

Open a notebook and connected to serverless, environment 5. 
Trying to read CDF from the MV:

spark.sql("SELECT * FROM table_changes('my_mv', 0)"😞

got error:
Spark version: 4.2.0 Error: (com.databricks.sql.transaction.tahoe.DeltaFileNotFoundException) [DELTA_LOG_FILE_NOT_FOUND] Unable to retrieve the delta log files to construct table version 0 starting from checkpoint version -1 at abfss://gold@myhostname123.dfs.core.windows.net/__unitystorage/schemas/7e714b9b-79cd-4afe-a95d-28f0026f6ef4/tables/c0072537-f031-4f89-9ab8-7bd4d2af0a70/_external_metadata/_delta_log.

---------

Can you help to point out what is missing?
Thanks,


 

8 REPLIES 8

Louis_Frolio
Databricks Employee
Databricks Employee

Hello @cdn_yyz_yul, I took a look at both internal and external documentation and here is what I found.

Your setup checklist matches the docs. Row tracking on, pipelines.externalMetadata.enabled set (pipeline and table level, the table property wins), PREVIEW channel, recent runtime. For MV change feeds delta.enableChangeDataFeed isn't the switch anyway, it runs on automatic CDF, so you haven't missed a flag.

Look at where the error points: .../_external_metadata/_delta_log. The MV change feed doesn't read the MV's own Delta log, it reads a separate external metadata log the pipeline writes when that flag is on. Your error says that log doesn't exist yet, so table_changes() can't build version 0. Configured, but never generated (or the update that should've generated it didn't stick).

Here's the order I'd try things in:

  1. Force the metadata to generate. Run this from serverless or a standard cluster on DBR 17.3 or above, as a principal with MODIFY on the MV:

    REPAIR TABLE <catalog>.<schema>.my_mv SYNC METADATA;
    
  2. Retry with the three-part name, and don't assume version 0 exists. Since you dropped and recreated the MV, the earliest available version may be later than 0. Try a timestamp from just before the full refresh, or set spark.databricks.delta.changeDataFeed.timestampOutOfRange.enabled = true so an out-of-range start returns empty instead of erroring.

  3. Confirm the full refresh actually finished green in the pipeline UI. A failed or partial update won't write the external metadata.

  4. Check for catalog commits. Catalog commits (Beta) isn't compatible with external metadata and has to be disabled first. Run DESCRIBE DETAIL <catalog>.<schema>.my_mv and look at the table features.

  5. Check the metastore-level "external data access" toggle in the account console. The MV CDF docs only mention the dataset-level flag, so I'm not certain it's required here, but if it's off I'd flip it and rerun the pipeline to rule it out.

  6. Rule out the reader. The MV CDF page says DBR 18 LTS or above, the general automatic CDF page says DBR 19 or above. Serverless environment 5 on Spark 4.2 should clear both, but if nothing above moves the needle, try the same query from a serverless SQL warehouse.

Two things to know as you build this out. You can't read an MV's change feed from inside the pipeline that creates it; use a separate PREVIEW pipeline for that pattern. And for pipeline-managed MVs, the supported reset is a full refresh or removing the MV from the definition, not a direct DROP. It worked here, but my not work in other situations.

References:

Let us know how it goes.

Regards, Louis.

Thanks @Louis_Frolio 

 

What I have done following your instructions:

- REPAIR TABLE: Run it, successful

- run pipeline full refresh, then run it  again (not selecting full refresh). No new or changed rows in between the two run. I just wanted to make sure it runs successfully. 

- read CDF fron that same notebook

SELECT * FROM table_changes('my_catalog.my_schema.my_mv', CURRENT_TIMESTAMP() - INTERVAL 4 HOUR)

[DELTA_MISSING_CHANGE_DATA] Error getting change data for range [3 , 34] as change data was not
recorded for version [3]. If you've enabled change data feed on this table,
use `DESCRIBE HISTORY` to see when it was first enabled.
Otherwise, to start recording change data, use `ALTER TABLE table_name SET TBLPROPERTIES
(delta.enableChangeDataFeed=true)`. SQLSTATE: KD002

The error is different this time. It seems to say CDF is supported but the changes are not recorded.

Questions: How to see when it was first enabled on a MV? Why asking:


Your point 4) Check for catalog commits. Catalog commits (Beta) isn't compatible with external metadata and has to be disabled first.

However, DESCRIBE DETAIL or DESCRIBE HISTORY can not be run again MV.
for example:
[EXPECT_TABLE_NOT_VIEW.NO_ALTERNATIVE] 'DESCRIBE DETAIL' expects a table but `my_catlog`.`my_schema`.`my_mv` is a view. SQLSTATE: 42809


Your point 2) spark.databricks.delta.changeDataFeed.timestampOutOfRange.enabled = true
This can not be set on serverless:

[CONFIG_NOT_AVAILABLE.SERVERLESS_DELTA_CHANGE_DATA_FEED_TIMESTAMP_OUT_OF_RANGE_ENABLED]
Configuration spark.databricks.delta.changeDataFeed.timestampOutOfRange.enabled is not available. SQLSTATE: 42K0I

----------------


Regarding the metastore level settings:

your point 5) Check the metastore-level "external data access" toggle in the account console.
Unfortunately, I do not have account admin permission in our development workspace. But, if it is not enabled, should I get some sort of permission denies error?
Or it would not say anything?


thanks again,

Na.

Louis_Frolio
Databricks Employee
Databricks Employee

@cdn_yyz_yul , Thanks for the detailed follow-up.

Read the new error carefully: "range [3, 34]." That means the external metadata log now exists and has versions 3 through 34 in it, so REPAIR TABLE did its job. The reader resolved "4 hours ago" to version 3, then couldn't produce change data for that one version. Everything from 4 onward may well be fine.

My best guess at why: version 3 is the first version in that freshly generated metadata log, essentially a snapshot with nothing before it to diff against, so there are no row-level changes to compute for it. The error text you got is the generic message Delta throws when it can't produce changes for a version, and it's written for legacy CDF, which is why it talks about delta.enableChangeDataFeed. Ignore that suggestion, it doesn't apply to MVs.

Two quick probes to confirm:

SELECT * FROM table_changes('my_catalog.my_schema.my_mv', 4);
SELECT * FROM table_changes('my_catalog.my_schema.my_mv', 34);

If either returns rows, the feature's working and the only issue is where you start reading. Two ways to handle that:

  1. Batch: start from the first version after the metadata was generated, not from a timestamp that lands on the bootstrap version.
  2. Streaming, which is what the docs recommend anyway: spark.readStream.option("readChangeFeed", "true").table("my_catalog.my_schema.my_mv") with no starting version. It takes the current snapshot as inserts and tracks changes forward from there, so you never have to guess at a version.

One more thing about your test. Your two pipeline runs had no row changes between them, so even a working read would show nothing new. To really validate it, change a row upstream, run a regular (not full refresh) update, then read from the version before that update. You should see the delete and insert pair the docs describe.

On your questions:

  • History on an MV: you can't, and that's by design. DESCRIBE HISTORY, DESCRIBE DETAIL, and time travel all reject MVs. The practical substitute is the _commit_version and _commit_timestamp columns in a table_changes result that does work, plus the pipeline event log for update timing.
  • Catalog commits: SHOW TBLPROPERTIES does work on MVs, so SHOW TBLPROPERTIES my_catalog.my_schema.my_mv ('delta.feature.catalogManaged') is the check. That said, since REPAIR TABLE succeeded and the metadata log got written, I'd cross catalog commits off the suspect list. Incompatible features would've blocked that step.
  • The timestampOutOfRange config: you're right, it isn't available on serverless. My mistake. Explicit versions are the way to go.

If versions 4 and 34 both fail with the same error, then automatic CDF isn't kicking in on the reader at all. Try the same query from a serverless SQL warehouse to rule out the notebook compute, and if that fails too, open a support ticket. You now have a very crisp report for them: REPAIR succeeded, external metadata log has versions 3 to 34, DELTA_MISSING_CHANGE_DATA on every version. That's exactly the kind of thing Beta support cases are for.

References:

Regards, Louis.

Thanks again @Louis_Frolio for all the information. 

1) 

SELECT * FROM table_changes('my_catalog.my_schema.my_mv', 4);
SELECT * FROM table_changes('my_catalog.my_schema.my_mv', 34);

I have already tried both and several other version numbers from both notebook and SQL warehouse.

example output:
-------------------------------
[DELTA_MISSING_CHANGE_DATA] Error getting change data for range [4 , 34] as change data was not
recorded for version [4]. If you've enabled change data feed on this table,
use `DESCRIBE HISTORY` to see when it was first enabled.
Otherwise, to start recording change data, use `ALTER TABLE table_name SET TBLPROPERTIES
(delta.enableChangeDataFeed=true)`.

-------------------------------
[DELTA_MISSING_CHANGE_DATA] Error getting change data for range [34 , 34] as change data was not
recorded for version [34]. If you've enabled change data feed on this table,
use `DESCRIBE HISTORY` to see when it was first enabled.
Otherwise, to start recording change data, use `ALTER TABLE table_name SET TBLPROPERTIES
(delta.enableChangeDataFeed=true)`.

-----------------------------
[DELTA_MISSING_CHANGE_DATA] Error getting change data for range [10 , 34] as change data was not
recorded for version [10]. If you've enabled change data feed on this table,
use `DESCRIBE HISTORY` to see when it was first enabled.
Otherwise, to start recording change data, use `ALTER TABLE table_name SET TBLPROPERTIES
(delta.enableChangeDataFeed=true)`.

2) introduce changes in between two pipeline runs.
I thought to do so once the read CDF returns no changes. However, I will try it. Maybe the above message is expected when no changes. 

cdn_yyz_yul
Contributor III

Update: Introduce row level changes to my_mv.

 Before adding new data to my_mv, CDF has enabled, the error can be found in my previous message.
At this stage:
- added new rows to the delta table where my_mv gets its input.
- run DLT pipeline, not selecting full refresh. ==> DLT automatically did "Full recompute".

- after the run finished successfully, check the number of rows in my_mv. the number of rows is changed from 83,330 to 85,095

- Run: 

SELECT * FROM table_changes('my_catalog.my_schema.my_mv', CURRENT_TIMESTAMP() - INTERVAL 5 HOUR)
fails with same error:
[DELTA_MISSING_CHANGE_DATA] Error getting change data for range [28 , 35] as change data was not recorded for version [28]. If you've enabled change data feed on this table, use `DESCRIBE HISTORY` to see when it was first enabled. Otherwise, to start recording change data, use `ALTER TABLE table_name SET TBLPROPERTIES (delta.enableChangeDataFeed=true)`. SQLSTATE: KD002

 

Louis_Frolio
Databricks Employee
Databricks Employee

Thanks for running the experiment, @cdn_yyz_yul. This one's useful, even though it failed.

First, I owe you a correction. Last time I guessed version 3 failed because it was the first version in the freshly generated metadata log. Now version 28 fails the same way, and 28 is nowhere near the start. Two different versions, same error, so my original theory was not right. Something's missing at every version you've landed on, not just the first one.

Here's what I think is going on. When REPAIR TABLE ran, it wrote versions 3 through 34 in one burst, a backfill of the MV's history. Version 35 is the only version written by the pipeline itself with all the flags on. Both of your timestamp reads resolved into the backfill range (the commit timestamps in that log all cluster around when REPAIR ran, so "N hours ago" keeps landing there). My suspicion is the backfilled versions don't carry the row-tracking metadata that automatic CDF needs, and 35 might.

One query settles it:

SELECT * FROM table_changes('my_catalog.my_schema.my_mv', 35);

If that returns rows, you're in business. Expect a lot of them, since the full recompute rewrote every row and the docs say a full rewrite shows up as deletes and inserts for unchanged rows too. From here on, read from explicit versions or use a streaming read, and skip timestamps until the backfill range is well behind you.

If 35 fails with the same error, then automatic CDF isn't engaging on this MV at all, and I don't think it's anything you can configure around. Three things to gather before opening a support ticket:

  1. The full property dump, SHOW TBLPROPERTIES my_catalog.my_schema.my_mv with no filter. Paste it here too. I want to see whether delta.feature.rowTracking shows as supported alongside the delta.enableRowTracking = true you already confirmed.
  2. Whether the pipeline runs on serverless or classic compute. The docs say serverless MVs get row tracking by default; on classic you're relying on the table property doing the job.
  3. The timeline as you've described it here: REPAIR succeeded, log has versions 3 to 35, DELTA_MISSING_CHANGE_DATA at 3, 28, and 35.

For what it's worth, your reader's fine. Spark 4.2.0 is Databricks Runtime 19, which clears both the "18 LTS" bar on the MV CDF page and the "19 or above" bar on the general automatic CDF page. Cross that one off.

References:

Fingers crossed on 35.

Regards, Louis.

Good morning @Louis_Frolio 

Re: > My suspicion is the backfilled versions don't carry the row-tracking metadata that automatic CDF needs, and 35 might.

1) SELECT * FROM table_changes('my_catalog.my_schema.my_mv', <an earlier enought timestamp that allows to show all versions after inception>)

It gives:
Error getting change data for range [3 , 35]

2) had a small loop to select from table_changes at each verstion from 3 to 36.

Got the "[DELTA_MISSING_CHANGE_DATA] Error getting change data for range" error that has posted many times in our discussion.

the version 36, got an expected error:
36
IllegalArgumentException, [DELTA_CDC_START_VERSION_AFTER_LATEST] Start version 36 for change data feed exceeds the latest table version 35.


Run another repair, then above tests. the tests had the same results.

 

- SHOW TBLPROPERTIES:

key,value
clusteringColumns,"[[""part_number""],[""serial_number""]]"
delta.enableChangeDataFeed,true
delta.enableDeletionVectors,true
delta.enableRowTracking,true
delta.feature.appendOnly,supported
delta.feature.changeDataFeed,supported
delta.feature.clustering,supported
delta.feature.deletionVectors,supported
delta.feature.domainMetadata,supported
delta.feature.invariants,supported
delta.feature.rowTracking,supported
delta.feature.v2Checkpoint,supported
delta.minReaderVersion,3
delta.minWriterVersion,7
delta.parquet.compression.codec,zstd
delta.rowTracking.materializedRowCommitVersionColumnName,_row-commit-version-col-d3dbf3a4-cc8b-448b-9c66-b92571ab0419
delta.rowTracking.materializedRowIdColumnName,_row-id-col-e84e791d-6d19-4b93-ad58-da9930f33298
pipeline_internal.catalogType,UNITY_CATALOG
pipeline_internal.enzymeMode,Advanced
pipelines.pipelineId,f7c4ab9f-c635-4a84-95bf-e106316e0fef
pipelines.reset.allowed,true
quality,gold



- Serverless compute (19.x	PREVIEW)


The timeline as you've described: REPAIR succeeded, log has versions 3 to 35, DELTA_MISSING_CHANGE_DATA from versions 3 to 35.



Regarding stream reader:
Streaming, which is what the docs recommend anyway: spark.readStream.option("readChangeFeed", "true").table("my_catalog.my_schema.my_mv") with no starting version. It takes the current snapshot as inserts and tracks changes forward from there, so you never have to guess at a version.

---
Streaming from Materialized Views is not supported.