Delta Lake Change Data Feed for Downstream Processing

By Anika Schneider • • 9 min read

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:

StrategyImpactWhen to Apply
Partition pruning in CDF readsReduces data scanned by orders of magnitudeAlways, when source table is partitioned
Version-based reads over timestampAvoids timestamp-to-version resolution overheadWhen exact version tracking is feasible
Filtering _change_type earlyEliminates unnecessary preimage readsWhen downstream does not need before-state
Small file compaction on CDF tablesPrevents file proliferation from frequent writesTables with high write frequency
Z-ORDER on join keysSpeeds up MERGE operations that trigger CDFTables 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_data directory 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.