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_pathTesting 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() == 0Testing 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 operationsTesting 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 acceptableTesting 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() == 0Testing 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() == 1Delta 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.