Delta Lake is a table format: Parquet data files plus a transaction log, _delta_log/, that records which files make up each version of the table. For a team that has run Hive tables on HDFS for years, it replaces "the table is whatever files are in these directories" with "the table is whatever the log says". That one change brings atomic commits, consistent reads during writes, MERGE, time travel and schema enforcement.
The internals of the log (actions, checkpoints, optimistic concurrency and deletion vectors) are covered in Delta Lake architecture in depth. This page is for a Hadoop estate adopting Delta: whether commits are safe on HDFS and on object storage, how to run it from Spark on YARN, how to migrate existing Hive tables, which engines can read the result, how maintenance interacts with the NameNode, and how Delta compares with Iceberg and Hive ACID for this audience.
What changes compared with a Hive table
A classic Hive table on HDFS is a directory tree plus a metastore entry. The metastore knows the schema and the partition directories. Which files are in the table is decided by listing those directories at query time. That design causes the familiar problems:
| Problem with a directory table | What Delta does instead |
|---|---|
| A reader listing a partition while a job writes it sees half the files | Readers use a snapshot from the log; new files are invisible until the commit |
| A failed INSERT OVERWRITE leaves a partition empty or mixed | The old files stay in the table until the commit that removes them succeeds |
| Updates mean rewriting whole partitions by hand | MERGE, UPDATE and DELETE rewrite only the affected files, in one commit |
| Schema drift goes unnoticed until a query fails | Writes with an unexpected schema are rejected unless evolution is requested |
| No history; a bad load is fixed by rerunning upstream | DESCRIBE HISTORY, time travel and RESTORE to an earlier version |
What Delta does not change: the data is still Parquet, the files still live on HDFS or object storage, and compute is still Spark (or another engine with a Delta reader). It is not a database server. Nothing runs beside the files, which is why the storage layer's guarantees matter so much.
Commit safety on HDFS and object storage
A Delta commit writes new data files and then creates the next numbered log file, for example version 42. If two writers both try to create version 42, exactly one must win; the loser re-reads the log, checks whether its change conflicts, and retries as 43 or fails. Correctness rests on the storage system refusing a second create of the same file name.
HDFS provides this directly. Delta's documentation states that HDFS needs no extra configuration, because the Hadoop FileSystem API gives the atomic rename the commit protocol relies on. That is one of the easier parts of running Delta on an existing cluster.
S3 is different. It has no mutual exclusion on object creation, so Delta's default S3 mode is safe only when all writes to a table come from a single Spark driver. Two clusters writing the same table in that mode can overwrite each other's commit and lose data without any error. For multi-cluster writes Delta ships a DynamoDB-backed LogStore:
# extra artifact: io.delta:delta-storage-s3-dynamodb (same version as delta-spark)
spark.delta.logStore.s3a.impl=io.delta.storage.S3DynamoDBLogStore
spark.io.delta.storage.S3DynamoDBLogStore.ddb.tableName=delta_log
spark.io.delta.storage.S3DynamoDBLogStore.ddb.region=us-east-1The DynamoDB table uses tablePath as partition key and fileName as sort key. Every writer to the table must use the same configuration. One job that forgets it reopens the race.
Running Delta from Spark on YARN
On a YARN cluster, Delta is a library added to the Spark job plus two session settings. Use the delta-spark artifact that matches your Spark and Scala versions; the Delta 4.x line targets Spark 4. Check the compatibility table in the Delta release notes rather than guessing.
spark-submit \
--master yarn --deploy-mode cluster \
--packages io.delta:delta-spark_2.13:<delta-version> \
--conf spark.sql.extensions=io.delta.sql.DeltaSparkSessionExtension \
--conf spark.sql.catalog.spark_catalog=org.apache.spark.sql.delta.catalog.DeltaCatalog \
daily_clicks.pyOn clusters without internet access, --packages cannot download anything. Stage the jars in your artifact repository or on HDFS and pass them with --jars. Register tables in the existing Hive Metastore so names keep working: CREATE TABLE web.clicks USING DELTA LOCATION 'hdfs:///warehouse/web/clicks'. The metastore then holds a pointer to the location, and the schema and file list come from the log.
Worked example: migrating a Hive table
The worked example is web.clicks: a Hive table on HDFS, partitioned by dt, about 400 daily partitions, loaded nightly with INSERT OVERWRITE. How you migrate depends on its file format.
Parquet: convert in place. CONVERT TO DELTA scans the existing files and writes a first log version listing them. No data is rewritten, so it is fast and the files stay where they are. Partition columns must be declared, because they exist only in directory names:
-- 1. stop the nightly load and any other writer
-- 2. convert: records existing files in _delta_log, rewrites nothing
CONVERT TO DELTA parquet.`hdfs:///warehouse/web/clicks` PARTITIONED BY (dt STRING);
-- 3. point the metastore name at the Delta table
DROP TABLE web.clicks; -- external table: drops metadata only; check before running
CREATE TABLE web.clicks USING DELTA LOCATION 'hdfs:///warehouse/web/clicks';
-- 4. validate before re-enabling writers
SELECT dt, count(*) FROM web.clicks WHERE dt >= '2026-09-01' GROUP BY dt ORDER BY dt;Before running step 3, check the table type. Dropping a managed Hive table deletes its files. Compare the row counts in step 4 with counts taken from the old table before conversion. Then switch the loader from INSERT OVERWRITE to Delta writes and re-enable it.
ORC: rewrite. CONVERT TO DELTA accepts Parquet and Iceberg, not ORC, and many Hive warehouses are ORC. Those tables must be rewritten, for example CREATE TABLE web.clicks_delta USING DELTA PARTITIONED BY (dt) AS SELECT * FROM web.clicks, backfilled a range of partitions at a time for large tables, then swapped by name. Budget cluster time and twice the storage during the overlap. The same applies to Hive ACID tables, whose base and delta directory layout belongs to Hive's own transaction system.
From INSERT OVERWRITE to MERGE
With the table in Delta, the nightly overwrite becomes an upsert. Late-arriving clicks and corrections update existing rows instead of forcing a partition rebuild:
from delta.tables import DeltaTable
staged = (spark.read.parquet("hdfs:///landing/clicks/2026-10-02")
.dropDuplicates(["click_id"])) # MERGE fails if two source rows hit one target row
clicks = DeltaTable.forName(spark, "web.clicks")
(clicks.alias("t")
.merge(staged.alias("s"),
"t.click_id = s.click_id AND t.dt >= '2026-09-25'") # prune to recent partitions
.whenMatchedUpdateAll()
.whenNotMatchedInsertAll()
.execute())Two details matter. The source must have at most one row per key, which is why the dedup is there. And the extra t.dt predicate lets Delta skip old partitions instead of scanning all 400. It also reduces conflicts if another job writes older partitions concurrently, because Delta's conflict check looks at which files each transaction read.
Who can read a Delta table
This is where Hadoop estates get caught out. Spark with the Delta library reads the log. Trino has a Delta Lake connector. Hive and Impala, the engines many Hadoop users query with, do not read the Delta log natively (Impala reads Iceberg, not Delta). The dangerous case is an engine that reads the directory as plain Parquet:
- After an UPDATE, MERGE or OPTIMIZE, the old files are logically removed but physically present until VACUUM. A plain-Parquet reader sees old and new files together: duplicated and stale rows, with no error.
- A plain reader also sees files from in-progress or failed writes that were never committed.
So never leave a Hive or Impala external Parquet table pointing at a Delta directory. Options, from simplest: query through Spark SQL or Trino; publish a separate plain-Parquet copy for legacy readers; or enable UniForm, which makes Delta also write Iceberg metadata so Iceberg readers can query the same files. UniForm is set per table with 'delta.enableIcebergCompatV2' = 'true' and 'delta.universalFormat.enabledFormats' = 'iceberg'. It requires column mapping and, per the Delta documentation, the Hive Metastore as catalog. Test the exact reader version you run before relying on it.
Maintenance: OPTIMIZE, VACUUM and history
Delta adds two maintenance jobs, and both interact with HDFS limits.
Compaction. Streaming and frequent small MERGEs create many small files, and every file is a NameNode object (see the HDFS small-files problem). OPTIMIZE web.clicks WHERE dt >= '2026-09-25' rewrites small files into larger ones for recent partitions. The old files stay until VACUUM, so compaction temporarily increases file count and HDFS usage.
VACUUM. Removed files stay on disk for time travel and for readers still using old snapshots. VACUUM deletes files no longer referenced by any version inside the retention window, which defaults to 7 days. The log and table history are kept for 30 days by default (delta.logRetentionDuration). Schedule VACUUM daily or weekly per table, and watch HDFS quota and namespace usage between runs.
DESCRIBE HISTORY web.clicks LIMIT 5; -- who wrote what, with row and file metrics
SELECT count(*) FROM web.clicks VERSION AS OF 812; -- read an earlier version
RESTORE TABLE web.clicks TO VERSION AS OF 812; -- undo a bad load (a new commit)
VACUUM web.clicks; -- default 7-day retentionRESTORE is itself a data-changing commit, so streaming consumers of the table may reprocess data. HDFS snapshots still work on a Delta directory and are useful before a risky migration, but they protect files, not table semantics. Restoring a snapshot of only the data files, without the matching log, produces a broken table.
Failure modes
Failure modes seen in practice:
- Two writers on S3 in default mode. Silent lost commits. Use the DynamoDB LogStore, or route all writes for a table through one job.
- VACUUM with a short retention. Setting
spark.databricks.delta.retentionDurationCheck.enabled=falseand vacuuming with zero hours deletes files that running queries and streams still need, and they fail. Keep retention longer than your longest job. - Legacy readers on the directory. Duplicate rows in Hive or Impala reports after the first MERGE. Covered above; audit external tables during migration.
- Concurrent writers to the same partitions. ConcurrentAppendException and related errors. Partition the work so jobs touch disjoint partitions, add partition predicates, and retry.
- Partition explosion. Partitioning by a high-cardinality column multiplies files exactly as it did in Hive. Partition by date, and use OPTIMIZE with Z-ordering or clustering for other filter columns.
- Version skew. Enabling a new table feature upgrades the protocol, and older clients can no longer read or write the table. Upgrade every reader before enabling features.
Delta, Iceberg or Hive ACID
For a Hadoop estate the realistic choices are Delta, Iceberg, or staying on Hive (ACID or plain). The modern lakehouse overview covers the wider landscape, and Iceberg with Trino covers the main alternative stack.
| Delta Lake | Apache Iceberg | Hive ACID / plain Hive | |
|---|---|---|---|
| Best-supported engine | Spark | Spark, Trino, Flink, Impala, Hive | Hive (ACID); everything (plain) |
| In-place migration from Parquet | Yes (CONVERT TO DELTA) | Yes (migrate procedures) | n/a |
| ORC tables | Rewrite to Parquet | Supported as data format | Native |
| Hidden partitioning, partition evolution | No; generated columns and clustering instead | Yes | No |
| Readable by Impala | No (UniForm: via Iceberg metadata) | Yes | Yes |
| Fits when | Spark-centric pipelines, Databricks interop | Many engines over one table | Hive-only legacy, no new features |
Choose Delta when Spark does almost all the writing and the read side is Spark, Trino or Databricks. Choose Iceberg when Impala and Hive users must keep querying the same tables without extra machinery. Avoid running both formats for the same data unless UniForm covers your readers.
What to do next
- Inventory the Hive tables you want to move: format (Parquet or ORC), managed or external, every writer, and every engine that reads each one.
- Decide the format per reader set: Delta if readers are Spark or Trino, Iceberg if Impala or Hive must keep reading.
- For object storage with more than one writing cluster, configure the DynamoDB LogStore in every job before the first write.
- Convert one Parquet table with CONVERT TO DELTA in a quiet window, validate per-partition counts, and repoint the metastore name.
- Drop or rename any Hive or Impala external table that still points at the converted directory.
- Replace INSERT OVERWRITE loads with MERGE, using deduplicated sources and partition predicates.
- Schedule OPTIMIZE for recent partitions and VACUUM with default retention; alert on HDFS namespace and quota usage.
- Upgrade every reader before enabling new table features, and practise RESTORE once before you need it.