Lab 15: Edge Cases & Concurrency Control
Before we write code, let's understand our Goal: Everything works perfectly when you are testing on your laptop. But in production, you have 50 different jobs trying to insert, update, and delete data into the exact same table at the exact same millisecond. We need to prevent the database from corrupting.
The Tools: Optimistic Concurrency Control (OCC) and automated Compaction.
1. Optimistic Concurrency Control (OCC)
Real-World Analogy Mapping: Imagine a Google Doc.
- The Long Way (Pessimistic Locking): Traditional databases use row-level locks. If you are typing on Page 1 of the Google doc, I am completely locked out of the document until you finish. This makes distributed processing incredibly slow.
- The Smart Way (Optimistic Concurrency Control): Delta Lake lets us both edit the Google Doc at the same time (Optimistic). But when we click "Save", the system checks if our edits clash. If you deleted a sentence, and I tried to bold that same sentence, Delta Lake stops me, throws a
ConcurrentAppendException, and says: "Wait, the document changed while you were typing. Please read the new version and try again."
2. Handling Concurrency Exceptions
Problems Faced:
When Delta Lake throws a ConcurrentAppendException, your PySpark job completely crashes and turns red.
How Present Technology Solves It:
You must wrap complex PySpark MERGE statements in a "Retry Block". If the job crashes because someone else was writing to the table, your code catches the error, waits 5 seconds, and tries again automatically!
import time
from pyspark.sql.utils import AnalysisException
max_retries = 3
for attempt in range(max_retries):
try:
# Attempt the highly concurrent MERGE
spark.sql("""
MERGE INTO target_table t
USING updates_table u
ON t.id = u.id
WHEN MATCHED THEN UPDATE SET *
""")
print("Merge successful!")
break # Exit the loop, it worked!
except AnalysisException as e:
if "Concurrent" in str(e):
print(f"Collision detected! Another job is writing. Retrying {attempt + 1}/{max_retries}...")
time.sleep(5) # Wait for the other job to finish
else:
raise e # If it's a different error (like a syntax error), crash the job normally
3. The "Small File" Streaming Problem
Follow the Data: Imagine you have a Structured Streaming job processing 1,000 JSON files per minute. It writes these files into a Delta table. After 24 hours, your Delta table consists of 1.4 million tiny 1KB files.
When a business analyst runs SELECT * FROM table, it takes 15 minutes! Why? Because opening 1.4 million individual files has massive network overhead.
The Fix: You must constantly compact the small files into large, efficient 1GB files. Instead of writing a manual script, you just flip a switch on the table!
-- Tell Databricks to automatically compact small files in the background!
ALTER TABLE target_table SET TBLPROPERTIES (
'delta.autoOptimize.autoCompact' = true,
'delta.autoOptimize.optimizeWrite' = true
);
← Previous: Lab 14: Infrastructure as Code with Terraform | End of Lab Series