Post

Delta Lake Basics for Data Engineers: A Practical Guide

In this article, let us understand the basics of Delta Lake and build a small pipeline using Apache Spark. A normal data lake is easy to start with because we can write Parquet files directly to object storage. The difficulty starts when we need reliable updates, concurrent jobs, schema checks, or recovery from a failed write. Delta Lake adds these features while keeping the data in the data lake.

For our use case, we will create an orders table, append records, correct an order, and process a change batch. The examples use PySpark, but the SQL also works in Databricks and other Spark environments with Delta Lake configured.

What Delta Lake adds

A Delta table is mainly Parquet data files together with a transaction log stored in a _delta_log directory. Each successful operation records metadata and the files added or removed. Readers use this log to find a consistent table version instead of listing random Parquet files and hoping a write is complete.

This gives us ACID transactions, schema enforcement, controlled schema evolution, UPDATE, DELETE, MERGE, and table history. Delta Lake does not replace Spark or cloud storage. Spark performs the processing, while data still lives in S3, ADLS, GCS, or another supported file system.

Storage approachUpdates and deletesSchema checksTransaction logBest fit
Plain ParquetManual rewriteLimitedNoSimple immutable datasets
Delta LakeBuilt-in operationsYesYesPipelines with changing data
Data warehouseBuilt-inYesManagedBI and serving workloads

Create a Delta table

Let us start with a small DataFrame. For open-source Spark, include a compatible Delta Lake package and configure the Delta Spark extensions. Databricks runtimes already include this setup.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, LongType, StringType, DoubleType

spark = SparkSession.builder.appName("delta-basics").getOrCreate()

schema = StructType([
    StructField("order_id", LongType(), False),
    StructField("customer_id", StringType(), False),
    StructField("status", StringType(), False),
    StructField("amount", DoubleType(), False)
])

orders = spark.createDataFrame([
    (1001, "C101", "CREATED", 89.50),
    (1002, "C102", "CREATED", 45.00),
    (1003, "C101", "SHIPPED", 120.25)
], schema)

orders.write.format("delta").mode("overwrite").save("/data/delta/orders")

We can read it like Parquet, with the format changed to delta.

1
2
delta_orders = spark.read.format("delta").load("/data/delta/orders")
delta_orders.show()

For repeated access, I prefer registering the location as a table. The storage path remains the source of data, but SQL becomes easier.

1
2
3
4
5
6
7
CREATE TABLE IF NOT EXISTS orders
USING DELTA
LOCATION '/data/delta/orders';

SELECT status, COUNT(*) AS order_count
FROM orders
GROUP BY status;

In production, the path would normally be a cloud URI such as s3://company-lake/silver/orders, and the table would be registered in a catalog. We also need storage permissions for both data files and the transaction log.

Append data and enforce the schema

Appending a batch is straightforward.

1
2
3
4
5
new_orders = spark.createDataFrame([
    (1004, "C103", "CREATED", 67.90)
], schema)

new_orders.write.format("delta").mode("append").save("/data/delta/orders")

If the incoming DataFrame has an unexpected column or incompatible type, Delta rejects the write by default. This is useful because a bad source file does not silently change the table. We still need data quality checks; schema enforcement cannot tell us that an amount is negative or a status is invalid.

For a legitimate new column, we can allow schema evolution explicitly.

1
2
3
4
5
changed_batch.write \
    .format("delta") \
    .mode("append") \
    .option("mergeSchema", "true") \
    .save("/data/delta/orders")

I would not enable automatic schema merging for every job. A misspelled column could become part of the table. Production schema changes should be reviewed and preferably applied through a controlled deployment.

Update and delete records

A plain Parquet dataset has no direct row-level update. Usually, we read data, change it, and overwrite a partition or complete dataset. Delta exposes these operations through SQL.

1
2
3
4
5
6
UPDATE orders
SET status = 'SHIPPED'
WHERE order_id = 1002;

DELETE FROM orders
WHERE order_id = 1004 AND status = 'CANCELLED';

Delta does not edit a Parquet row in place. It writes replacement files and records old files as removed. An update touching a large part of a table can therefore be expensive. Partitioning and file layout matter even though the SQL looks like a database operation.

Upsert a change batch with MERGE

MERGE is one of the main reasons I use Delta Lake for ingestion pipelines. If a source sends new orders and changes to existing orders, we can update matching IDs and insert IDs that do not exist.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
from delta.tables import DeltaTable

changes = spark.createDataFrame([
    (1002, "C102", "DELIVERED", 45.00),
    (1005, "C104", "CREATED", 210.00)
], schema)

target = DeltaTable.forPath(spark, "/data/delta/orders")

(target.alias("target")
    .merge(changes.alias("source"),
           "target.order_id = source.order_id")
    .whenMatchedUpdateAll()
    .whenNotMatchedInsertAll()
    .execute())

For a real CDC pipeline, I would include the source event timestamp or sequence number. Otherwise, a delayed event could overwrite a newer value. Deduplicate each incoming batch before MERGE, because multiple source rows for one target key can make the result ambiguous or fail the operation.

Check table history

Every successful write creates a table version. We can inspect operations to understand what changed.

1
DESCRIBE HISTORY orders;

History includes the operation, timestamp, version, and operation metrics. This helps when troubleshooting because we can see whether a job appended, overwrote, or merged data. Version reads are useful for investigation, although they should not be the only backup strategy. Old files are eventually removed according to retention and cleanup settings.

1
SELECT * FROM orders VERSION AS OF 1;

Maintenance and practical caveats

Delta Lake gives table reliability, but it still needs maintenance. Frequent small writes create many small Parquet files, increasing file listing and scan overhead. Depending on the platform, use OPTIMIZE or a compaction job. Choose partition columns carefully. A high-cardinality column such as order_id is usually a poor partition choice because it creates too many directories and small files.

Removed files are not deleted immediately because older versions may reference them. VACUUM cleans files after the retention period. Reducing retention without understanding active readers can break long-running jobs and remove access to older versions.

There are a few more points to be careful about:

  • Match the Delta Lake version with the Spark version used by the job.
  • Do not write directly to Parquet files inside a Delta table path. All writers should use Delta.
  • Monitor table size, file count, merge duration, and failed transactions.
  • Use a catalog and clear ownership instead of scattering unmanaged paths across buckets.
  • Test concurrent writes with the actual partition and merge pattern.

For a demo, a local path and one Spark session are enough. In production, I would add a catalog, cloud access controls, checkpoints for streaming jobs, data quality checks, retention policies, and scheduled compaction. I would also make ingestion idempotent so retrying a failed batch does not create duplicate records.

Conclusion

Delta Lake is useful when a data lake needs more than append-only files. The transaction log gives consistent reads, schema enforcement blocks accidental changes, and MERGE makes incremental pipelines easier to build. The first table is simple, but file sizing, schema changes, retention, and late-arriving records still need deliberate handling in production.

This post is licensed under CC BY 4.0 by the author.