Delta Lake Testing: ACID Transactions, Time Travel, and Schema Evolution

Delta Lake Testing: ACID Transactions, Time Travel, and Schema Evolution

Delta Lake adds ACID transactions, time travel, and schema enforcement to Apache Spark data lakes. These features introduce behaviors that require specific testing: concurrent writes must not corrupt data, time travel must return correct historical snapshots, schema changes must follow defined rules. This guide covers how to test Delta Lake behavior that goes beyond standard DataFrame testing.

Setting Up Delta Lake for Testing

pip install delta-spark pyspark pytest
# conftest.py
import pytest
from pyspark.sql import SparkSession
import tempfile
import shutil

@pytest.fixture(scope="session")
def spark():
    session = SparkSession.builder \
        .appName("delta-tests") \
        .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
        .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
        .config("spark.jars.packages", "io.delta:delta-core_2.12:2.4.0") \
        .master("local[2]") \
        .getOrCreate()
    
    session.sparkContext.setLogLevel("ERROR")
    yield session
    session.stop()

@pytest.fixture
def delta_table_path(tmp_path):
    """Provide a clean temp directory for each test."""
    path = str(tmp_path / "delta-table")
    yield path
    # Cleanup happens automatically with tmp_path

Testing ACID Writes

# test_acid.py
from delta.tables import DeltaTable
from pyspark.sql import functions as F
import pytest

def test_atomic_write(spark, delta_table_path):
    """Write should be atomic — either all or nothing."""
    initial_data = spark.createDataFrame([
        (1, "Alice", 100),
        (2, "Bob", 200),
    ], ["id", "name", "amount"])
    
    initial_data.write.format("delta").save(delta_table_path)
    
    # Try to write invalid data that should fail
    # (e.g., schema mismatch when enforced)
    try:
        invalid_data = spark.createDataFrame([
            (3, "Charlie", "not-a-number"),  # Wrong type
        ], ["id", "name", "amount"])
        
        invalid_data.write \
            .format("delta") \
            .mode("append") \
            .option("mergeSchema", "false") \
            .save(delta_table_path)
    except Exception:
        pass  # Expected to fail
    
    # Table should be unchanged
    result = spark.read.format("delta").load(delta_table_path)
    assert result.count() == 2

def test_concurrent_appends_both_succeed(spark, delta_table_path):
    """Multiple concurrent writes should both succeed without data loss."""
    from concurrent.futures import ThreadPoolExecutor
    import threading
    
    initial = spark.createDataFrame([(0, "init")], ["id", "name"])
    initial.write.format("delta").save(delta_table_path)
    
    errors = []
    
    def append_batch(batch_id):
        try:
            batch = spark.createDataFrame(
                [(i, f"batch-{batch_id}-{i}") for i in range(1, 11)],
                ["id", "name"]
            )
            # Delta Lake handles concurrent appends with optimistic concurrency
            batch.write.format("delta").mode("append").save(delta_table_path)
        except Exception as e:
            errors.append(e)
    
    with ThreadPoolExecutor(max_workers=3) as executor:
        futures = [executor.submit(append_batch, i) for i in range(3)]
        for f in futures:
            f.result(timeout=60)
    
    assert len(errors) == 0, f"Concurrent writes failed: {errors}"
    
    result = spark.read.format("delta").load(delta_table_path)
    # 1 initial + 3 batches × 10 rows each = 31
    assert result.count() == 31

def test_overwrite_is_atomic(spark, delta_table_path):
    """Overwrite should replace all data atomically."""
    original = spark.createDataFrame([
        (1, "original-1"),
        (2, "original-2"),
        (3, "original-3"),
    ], ["id", "name"])
    original.write.format("delta").save(delta_table_path)
    
    replacement = spark.createDataFrame([
        (10, "new-1"),
        (20, "new-2"),
    ], ["id", "name"])
    replacement.write.format("delta").mode("overwrite").save(delta_table_path)
    
    result = spark.read.format("delta").load(delta_table_path)
    
    assert result.count() == 2
    assert set(result.select("id").rdd.flatMap(lambda x: x).collect()) == {10, 20}
    # Original data should be gone
    assert result.filter(F.col("id") < 5).count() == 0

Testing Time Travel

# test_time_travel.py
from delta.tables import DeltaTable

def test_time_travel_by_version(spark, delta_table_path):
    """Time travel should return data as of a specific version."""
    # Version 0 — initial data
    v0_data = spark.createDataFrame([(1, "v0-row-1")], ["id", "name"])
    v0_data.write.format("delta").save(delta_table_path)
    
    # Version 1 — append
    v1_data = spark.createDataFrame([(2, "v1-row-2")], ["id", "name"])
    v1_data.write.format("delta").mode("append").save(delta_table_path)
    
    # Version 2 — overwrite
    v2_data = spark.createDataFrame([(3, "v2-row-3"), (4, "v2-row-4")], ["id", "name"])
    v2_data.write.format("delta").mode("overwrite").save(delta_table_path)
    
    # Current (v2): 2 rows
    current = spark.read.format("delta").load(delta_table_path)
    assert current.count() == 2
    
    # Time travel to v0: 1 row
    v0 = spark.read.format("delta").option("versionAsOf", 0).load(delta_table_path)
    assert v0.count() == 1
    assert v0.collect()[0]["name"] == "v0-row-1"
    
    # Time travel to v1: 2 rows
    v1 = spark.read.format("delta").option("versionAsOf", 1).load(delta_table_path)
    assert v1.count() == 2
    
    # Time travel to v2: 2 rows (after overwrite)
    v2 = spark.read.format("delta").option("versionAsOf", 2).load(delta_table_path)
    assert v2.count() == 2
    assert set(v2.select("id").rdd.flatMap(lambda x: x).collect()) == {3, 4}

def test_time_travel_by_timestamp(spark, delta_table_path):
    """Time travel by timestamp should return correct snapshot."""
    import time
    
    # Write version 0
    v0 = spark.createDataFrame([(1, "version-zero")], ["id", "name"])
    v0.write.format("delta").save(delta_table_path)
    
    # Record timestamp after v0
    ts_after_v0 = time.time()
    time.sleep(1)  # Ensure timestamp difference
    
    # Write version 1
    v1 = spark.createDataFrame([(2, "version-one")], ["id", "name"])
    v1.write.format("delta").mode("append").save(delta_table_path)
    
    # Time travel to ts_after_v0 — should see only v0
    from datetime import datetime
    ts_str = datetime.fromtimestamp(ts_after_v0).strftime('%Y-%m-%d %H:%M:%S')
    
    historical = spark.read.format("delta") \
        .option("timestampAsOf", ts_str) \
        .load(delta_table_path)
    
    assert historical.count() == 1
    assert historical.collect()[0]["name"] == "version-zero"

def test_history_logged(spark, delta_table_path):
    """Delta should log all operations in history."""
    data = spark.createDataFrame([(1, "row1")], ["id", "name"])
    data.write.format("delta").save(delta_table_path)
    
    data2 = spark.createDataFrame([(2, "row2")], ["id", "name"])
    data2.write.format("delta").mode("append").save(delta_table_path)
    
    delta_table = DeltaTable.forPath(spark, delta_table_path)
    history = delta_table.history()
    
    assert history.count() >= 2
    
    operations = history.select("operation").rdd.flatMap(lambda x: x).collect()
    assert "WRITE" in operations

Testing Schema Evolution

# test_schema_evolution.py

def test_schema_enforcement_blocks_incompatible_types(spark, delta_table_path):
    """Schema enforcement should reject type mismatches."""
    initial = spark.createDataFrame([(1, 100)], ["id", "amount"])
    initial.write.format("delta").save(delta_table_path)
    
    # Try to write wrong type for 'amount'
    wrong_type = spark.createDataFrame([(2, "not-a-number")], ["id", "amount"])
    
    with pytest.raises(Exception):
        wrong_type.write \
            .format("delta") \
            .mode("append") \
            .option("mergeSchema", "false") \
            .save(delta_table_path)

def test_schema_evolution_adds_column(spark, delta_table_path):
    """mergeSchema should allow adding new columns."""
    v1 = spark.createDataFrame([(1, "Alice")], ["id", "name"])
    v1.write.format("delta").save(delta_table_path)
    
    # Add new column in v2
    v2 = spark.createDataFrame([(2, "Bob", "NY")], ["id", "name", "city"])
    v2.write \
        .format("delta") \
        .mode("append") \
        .option("mergeSchema", "true") \
        .save(delta_table_path)
    
    result = spark.read.format("delta").load(delta_table_path)
    
    assert "city" in result.columns
    assert result.count() == 2
    
    # Original row should have null for new column
    alice = result.filter(result.id == 1).collect()[0]
    assert alice["city"] is None
    
    # New row has city value
    bob = result.filter(result.id == 2).collect()[0]
    assert bob["city"] == "NY"

def test_cannot_drop_column_without_overwrite(spark, delta_table_path):
    """You can't drop a column via append — must use overwrite."""
    v1 = spark.createDataFrame([(1, "Alice", "NY")], ["id", "name", "city"])
    v1.write.format("delta").save(delta_table_path)
    
    # Append without 'city' column — should fail (schema mismatch if strict)
    no_city = spark.createDataFrame([(2, "Bob")], ["id", "name"])
    
    # This will either fail or pad with null — test your specific behavior
    try:
        no_city.write \
            .format("delta") \
            .mode("append") \
            .option("mergeSchema", "false") \
            .save(delta_table_path)
        
        # If it succeeds, row 2 should have null city
        result = spark.read.format("delta").load(delta_table_path)
        bob = result.filter(result.id == 2).collect()[0]
        assert bob["city"] is None
    except Exception:
        pass  # Schema enforcement blocked the write — also acceptable

Testing MERGE (Upsert) Operations

# test_merge.py

def test_merge_upsert(spark, delta_table_path):
    """MERGE should update existing records and insert new ones."""
    # Initial table
    initial = spark.createDataFrame([
        (1, "Alice", 100),
        (2, "Bob", 200),
        (3, "Charlie", 300),
    ], ["id", "name", "amount"])
    initial.write.format("delta").save(delta_table_path)
    
    delta_table = DeltaTable.forPath(spark, delta_table_path)
    
    # Upsert — update existing ids, insert new
    updates = spark.createDataFrame([
        (2, "Bob Updated", 250),   # Update Bob
        (4, "Dave", 400),          # Insert Dave
    ], ["id", "name", "amount"])
    
    delta_table.alias("target").merge(
        updates.alias("source"),
        "target.id = source.id"
    ).whenMatchedUpdateAll().whenNotMatchedInsertAll().execute()
    
    result = spark.read.format("delta").load(delta_table_path)
    
    assert result.count() == 4  # Original 3 + new Dave
    
    bob = result.filter(result.id == 2).collect()[0]
    assert bob["name"] == "Bob Updated"
    assert bob["amount"] == 250
    
    dave = result.filter(result.id == 4).collect()[0]
    assert dave["name"] == "Dave"
    
    charlie = result.filter(result.id == 3).collect()[0]
    assert charlie["amount"] == 300  # Unchanged

def test_merge_delete_matching(spark, delta_table_path):
    """MERGE with whenMatchedDelete removes matching records."""
    initial = spark.createDataFrame([
        (1, "keep-me"),
        (2, "delete-me"),
        (3, "keep-me-too"),
    ], ["id", "name"])
    initial.write.format("delta").save(delta_table_path)
    
    delta_table = DeltaTable.forPath(spark, delta_table_path)
    
    to_delete = spark.createDataFrame([(2,)], ["id"])
    
    delta_table.alias("target").merge(
        to_delete.alias("source"),
        "target.id = source.id"
    ).whenMatchedDelete().execute()
    
    result = spark.read.format("delta").load(delta_table_path)
    
    assert result.count() == 2
    assert result.filter(result.id == 2).count() == 0

Testing Table Statistics and OPTIMIZE

def test_optimize_reduces_files(spark, delta_table_path):
    """OPTIMIZE should compact small files."""
    # Write many small batches (creates many small files)
    for i in range(10):
        batch = spark.createDataFrame([(i, f"row-{i}")], ["id", "name"])
        batch.write.format("delta").mode("append").save(delta_table_path)
    
    delta_table = DeltaTable.forPath(spark, delta_table_path)
    
    # Count files before optimize
    files_before = delta_table.detail().select("numFiles").collect()[0]["numFiles"]
    
    # Run OPTIMIZE
    spark.sql(f"OPTIMIZE delta.`{delta_table_path}`")
    
    # Count files after
    files_after = delta_table.detail().select("numFiles").collect()[0]["numFiles"]
    
    # Should have fewer files after compaction
    assert files_after <= files_before
    
    # Data should be unchanged
    result = spark.read.format("delta").load(delta_table_path)
    assert result.count() == 10

def test_vacuum_removes_old_versions(spark, delta_table_path):
    """VACUUM should remove files from old versions."""
    initial = spark.createDataFrame([(1, "v1")], ["id", "name"])
    initial.write.format("delta").save(delta_table_path)
    
    update = spark.createDataFrame([(1, "v2")], ["id", "name"])
    update.write.format("delta").mode("overwrite").save(delta_table_path)
    
    # VACUUM with 0 hours retention (test only — never use in production)
    spark.conf.set("spark.databricks.delta.retentionDurationCheck.enabled", "false")
    spark.sql(f"VACUUM delta.`{delta_table_path}` RETAIN 0 HOURS")
    
    # Time travel to v0 should now fail (files vacuumed)
    with pytest.raises(Exception, match="not found"):
        spark.read.format("delta") \
            .option("versionAsOf", 0) \
            .load(delta_table_path) \
            .count()
    
    # Current data should still be accessible
    current = spark.read.format("delta").load(delta_table_path)
    assert current.count() == 1

Delta Lake's guarantees — atomicity, time travel, schema enforcement — are only meaningful if you test them. These tests catch bugs in your write pipeline, schema evolution configuration, and merge logic before they corrupt production data that's hard to recover.

Start now free