Delta Lake Change Data Feed for Downstream Processing
In a modern data lakehouse, tables are rarely consumed by a single pipeline. A fact table in your bronze or silver layer might feed a dozen downstream consumers: aggregate tables, machine learning feature stores, reverse ETL systems, and operational dashboards. The naive approach to propagating changes is full table scans, where every downstream job rereads the entire source table to produce its output. At petabyte scale, this becomes untenable. Delta Lake Change Data Feed (CDF) solves this by exposing row-level changes between table versions, enabling truly incremental downstream processing.
CDF is not a replacement for source-system CDC tools like Debezium or AWS DMS. Those tools capture changes from database transaction logs. CDF operates at the lakehouse storage layer, recording what changed within a Delta table across commits. This distinction matters: CDF extends the incremental processing pattern deeper into your architecture, from the point where data lands in Delta all the way through to the serving layer.
Understanding the Change Data Feed Mechanism
When CDF is enabled on a Delta table, every write operation produces additional metadata that records the nature of each row-level change. Delta Lake classifies changes into three types:
- insert: A new row was added to the table.
- update_preimage: The state of a row before it was modified.
- update_postimage: The state of a row after it was modified.
- delete: A row was removed from the table.
These change records are stored in a _change_data directory within the Delta table's storage location, separate from the regular data files. Each change record includes three additional metadata columns: _change_type, _commit_version, and _commit_timestamp. These columns allow consumers to filter, order, and checkpoint their processing against specific table versions.
The storage mechanism is efficient because insert-only operations do not duplicate data. When a commit contains only inserts, CDF simply marks those data files as containing insert changes without writing separate change files. Only updates and deletes, which require capturing the before and after states, generate additional storage in the _change_data directory.
Enabling and Configuring CDF
CDF can be enabled on new or existing Delta tables. For new tables, set the table property at creation time:
-- Enable CDF on a new table
CREATE TABLE silver.customer_orders (
order_id BIGINT,
customer_id BIGINT,
product_id BIGINT,
quantity INT,
total_amount DECIMAL(12, 2),
order_status STRING,
updated_at TIMESTAMP
)
USING DELTA
TBLPROPERTIES (delta.enableChangeDataFeed = true);
For existing tables, alter the table properties. Note that CDF only captures changes from commits made after enablement; it does not retroactively generate change data for prior versions:
-- Enable CDF on an existing table
ALTER TABLE silver.customer_orders
SET TBLPROPERTIES (delta.enableChangeDataFeed = true);
-- Verify the property
DESCRIBE DETAIL silver.customer_orders;
-- Look for delta.enableChangeDataFeed = true in properties
You can also enable CDF as a default for all new Delta tables in a Spark session:
spark.conf.set(
"spark.databricks.delta.properties.defaults.enableChangeDataFeed",
"true"
)
The retention period for CDF data follows the same delta.logRetentionDuration setting as the Delta transaction log, defaulting to 30 days. Running VACUUM removes CDF files for versions beyond the retention window, so downstream consumers must process changes within that period or risk missing data.
Reading Change Data in Batch Mode
The simplest way to consume CDF is in batch mode, reading all changes between two versions or timestamps. This is the pattern used by most ETL pipelines that run on a schedule:
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.config("spark.sql.extensions",
"io.delta.sql.DeltaSparkSessionExtension") \
.config("spark.sql.catalog.spark_catalog",
"org.apache.spark.sql.delta.catalog.DeltaCatalog") \
.getOrCreate()
# Read changes between specific versions
changes_df = spark.read.format("delta") \
.option("readChangeFeed", "true") \
.option("startingVersion", 15) \
.option("endingVersion", 22) \
.table("silver.customer_orders")
# Examine the change metadata columns
changes_df.select(
"order_id",
"customer_id",
"order_status",
"_change_type",
"_commit_version",
"_commit_timestamp"
).show(truncate=False)
# Filter to only see updates (both pre and post images)
updates = changes_df.filter(
changes_df._change_type.isin(
"update_preimage", "update_postimage"
)
)
You can also specify time ranges instead of version numbers, which is often more convenient for scheduled pipelines:
# Read changes from the last 24 hours
from datetime import datetime, timedelta
yesterday = (datetime.now() - timedelta(days=1)).strftime(
"%Y-%m-%d %H:%M:%S"
)
changes_df = spark.read.format("delta") \
.option("readChangeFeed", "true") \
.option("startingTimestamp", yesterday) \
.table("silver.customer_orders")
A practical pattern for incremental ETL is to maintain a checkpoint that records the last processed version. Each pipeline run reads changes from that version forward, processes them, and updates the checkpoint:
# Read the last processed version from a checkpoint table
last_version = spark.sql("""
SELECT max(processed_version) as v
FROM etl_metadata.checkpoints
WHERE table_name = 'customer_orders'
""").collect()[0]["v"] or 0
# Get the current version of the source table
current_version = spark.sql(
"DESCRIBE HISTORY silver.customer_orders LIMIT 1"
).collect()[0]["version"]
if current_version > last_version:
changes = spark.read.format("delta") \
.option("readChangeFeed", "true") \
.option("startingVersion", last_version + 1) \
.option("endingVersion", current_version) \
.table("silver.customer_orders")
# Process changes and write to gold layer
process_incremental_changes(changes)
# Update checkpoint
spark.sql(f"""
INSERT INTO etl_metadata.checkpoints
VALUES ('customer_orders', {current_version},
current_timestamp())
""")
Streaming Integration with Structured Streaming
For near-real-time downstream processing, CDF integrates natively with Spark Structured Streaming. Instead of polling for batch changes on a schedule, you configure a streaming source that continuously picks up new commits as they happen:
# Continuous streaming from CDF
stream = spark.readStream.format("delta") \
.option("readChangeFeed", "true") \
.option("startingVersion", 0) \
.table("silver.customer_orders")
# Apply transformations on the change stream
processed = stream \
.filter("_change_type != 'update_preimage'") \
.withColumn("is_new",
when(col("_change_type") == "insert", True)
.otherwise(False)) \
.select(
"order_id",
"customer_id",
"total_amount",
"order_status",
"is_new",
"_commit_timestamp"
)
# Write to a downstream Delta table with checkpointing
query = processed.writeStream \
.format("delta") \
.outputMode("append") \
.option("checkpointLocation",
"s3://lake/checkpoints/order_changes") \
.trigger(processingTime="30 seconds") \
.toTable("gold.order_change_events")
The checkpoint mechanism in Structured Streaming tracks exactly which Delta versions have been processed, providing exactly-once semantics. If the stream fails and restarts, it resumes from the last committed checkpoint without duplicating or missing changes.
This streaming approach pairs naturally with stream processing frameworks like Flink and Kafka Streams for organizations that need to bridge their lakehouse changes into event-driven architectures. The Delta CDF stream can publish to Kafka topics, which downstream Flink jobs then consume.
Architecture Patterns for CDF Pipelines
CDF enables several architectural patterns that are difficult or inefficient to implement with full-table scans.
Medallion Architecture with Incremental Propagation
The medallion architecture (bronze, silver, gold) becomes dramatically more efficient when each layer consumes CDF from the layer below instead of rescanning entire tables. A change in bronze propagates incrementally through silver and gold, reducing processing time from hours to minutes for large datasets.
# Silver layer: incrementally process bronze changes
bronze_changes = spark.readStream.format("delta") \
.option("readChangeFeed", "true") \
.option("startingVersion", "latest") \
.table("bronze.raw_orders")
# Clean, validate, and deduplicate
silver_updates = bronze_changes \
.filter("_change_type != 'update_preimage'") \
.withColumn("clean_amount",
col("total_amount").cast("decimal(12,2)")) \
.dropDuplicates(["order_id", "_commit_version"])
# Merge into silver table
def upsert_to_silver(batch_df, batch_id):
batch_df.createOrReplaceTempView("changes")
spark.sql("""
MERGE INTO silver.customer_orders AS target
USING changes AS source
ON target.order_id = source.order_id
WHEN MATCHED AND source._change_type = 'delete'
THEN DELETE
WHEN MATCHED
THEN UPDATE SET *
WHEN NOT MATCHED AND source._change_type != 'delete'
THEN INSERT *
""")
silver_updates.writeStream \
.foreachBatch(upsert_to_silver) \
.option("checkpointLocation",
"s3://lake/checkpoints/bronze_to_silver") \
.trigger(processingTime="1 minute") \
.start()
Reverse ETL and Data Synchronization
CDF is ideal for reverse ETL, the practice of pushing analytics results back into operational systems. Instead of exporting full tables, you extract only the changed rows and sync them to downstream APIs, CRMs, or caches:
# Extract only changed customer scores for API sync
score_changes = spark.read.format("delta") \
.option("readChangeFeed", "true") \
.option("startingVersion", last_sync_version) \
.table("gold.customer_risk_scores")
# Only sync rows where the score actually changed
meaningful_changes = score_changes \
.filter("_change_type IN ('insert', 'update_postimage')") \
.select("customer_id", "risk_score", "risk_tier")
# Push to operational API in batches
for batch in meaningful_changes.toLocalIterator():
api_client.update_customer_score(
customer_id=batch.customer_id,
score=batch.risk_score,
tier=batch.risk_tier
)
Audit and Compliance Logging
CDF provides a built-in audit trail showing exactly what changed, when, and in which commit. This is valuable for regulatory compliance requirements that mandate data lineage and change tracking:
# Build an immutable audit log from CDF
audit_entries = spark.read.format("delta") \
.option("readChangeFeed", "true") \
.option("startingVersion", 0) \
.table("silver.customer_pii") \
.select(
"customer_id",
"_change_type",
"_commit_version",
"_commit_timestamp",
"email",
"phone_number"
)
# Write to an append-only audit table
audit_entries.write.format("delta") \
.mode("append") \
.saveAsTable("compliance.pii_change_audit")
Performance Optimization and Best Practices
Several strategies maximize CDF performance in production:
| Strategy | Impact | When to Apply |
|---|---|---|
| Partition pruning in CDF reads | Reduces data scanned by orders of magnitude | Always, when source table is partitioned |
| Version-based reads over timestamp | Avoids timestamp-to-version resolution overhead | When exact version tracking is feasible |
| Filtering _change_type early | Eliminates unnecessary preimage reads | When downstream does not need before-state |
| Small file compaction on CDF tables | Prevents file proliferation from frequent writes | Tables with high write frequency |
| Z-ORDER on join keys | Speeds up MERGE operations that trigger CDF | Tables with upsert-heavy patterns |
One important consideration is the interaction between CDF and OPTIMIZE and VACUUM operations. Running OPTIMIZE on a CDF-enabled table compacts data files but does not generate CDF records, since the logical content of the table is unchanged. However, VACUUM removes old CDF files beyond the retention period, which means consumers that fall behind will encounter errors when trying to read expired versions.
Set your CDF retention period to at least twice the maximum expected pipeline delay. If your slowest downstream job runs daily, retain CDF data for at least 48 hours to account for job failures and retries.
When tuning the upstream writes that generate CDF data, Spark's Adaptive Query Execution plays a critical role. AQE's dynamic partition coalescing ensures that MERGE and UPDATE operations produce well-sized output files, which directly affects CDF read performance.
Monitoring and Observability
Monitoring CDF pipelines requires tracking metrics that go beyond standard Spark job monitoring. Key metrics to capture include:
- Version lag: The difference between the source table's current version and the last version processed by each consumer. Growing lag indicates a consumer that cannot keep up with the write rate.
- Change volume per version: Spikes in change volume can indicate upstream data quality issues or unexpected bulk operations that may overwhelm downstream systems.
- CDF storage size: Track the
_change_datadirectory size separately from the main table data. Unexpectedly large CDF storage suggests inefficient update patterns that generate excessive preimage/postimage pairs. - Consumer checkpoint freshness: For batch consumers, track how recently each consumer updated its checkpoint. Stale checkpoints risk exceeding the retention window.
# Monitor CDF lag for all consumers
consumer_lag = spark.sql("""
SELECT
c.consumer_name,
c.table_name,
c.processed_version,
h.version as current_version,
h.version - c.processed_version as version_lag,
h.timestamp as latest_commit_time
FROM etl_metadata.checkpoints c
CROSS JOIN (
SELECT version, timestamp
FROM (DESCRIBE HISTORY silver.customer_orders)
LIMIT 1
) h
ORDER BY version_lag DESC
""")
consumer_lag.show()
For teams building comprehensive data observability platforms, CDF version lag and change volume metrics are essential signals that should feed into alerting and incident response workflows. A consumer that falls behind its retention window will silently produce incomplete data unless monitored.
Delta Lake CDF transforms the economics of data propagation in a lakehouse. By replacing full-table scans with targeted change reads, pipelines that once required hours of compute can run in minutes. The key to success is thoughtful architecture: enable CDF early, design consumers around version-based checkpointing, and monitor lag metrics to catch processing gaps before they become data quality incidents.