Delta Lake Theory - ๐ Delta Lake Beginner's Guide - Theory
Delta Lake Practical Guide Part 1 - ๐ Delta Lake Beginner's Guide - Practical Edition (Part 1. Local Environment)
Delta Lake in Practice Part 2 - [๐ Delta Lake Beginner's Guide - Practical Edition (Part 2. Utilizing the delta-spark library)
](https://velog.io/@newnew_daddy/data09)
In the previous post , ๐ Delta Lake Beginner's Guide - Theory, we covered the theoretical aspects of Delta Lake in detail. In this post, we will practice the features of Delta Lake in a Pyspark Docker Container environment.
In this practice session, we use PySpark as the tool for handling data, and files of type delta are stored and managed in a local directory.
hyunsoolee0506/pyspark-cloud:3.5.1For Docker containers , I recommend using an image I have created separately, considering that you will be integrating with a cloud environment for future practice . However, for this local environment practice, it is fine to proceed using Google Colab.
For the practice session, delta_data.zipwe used data from two CSV files located in the file below.
๐ย Talaria x3
/workspace/sparkConfigure the directory to be used as a volume on the user's computer to be mapped to the directory inside the container .
docker run -d \
--name pyspark \
-p 8888:8888 \
-p 4040:4040 \
-v [์ฌ์ฉ์ ๋๋ ํ ๋ฆฌ]:/workspace/spark \
hyunsoolee0506/pyspark-cloud:3.5.1
After executing the command above, 8888you can access the juypter lab development environment by connecting to the port.
Install the libraries required for the practice. hyunsoolee0506/pyspark-cloud:3.5.1Although they are already installed in the image, for Colab, you must install them by running the code below.
pip install pyspark==3.5.1 delta-spark==3.2.0 pyarrow findspark
from delta import *
from pyspark.sql import SparkSession
import pyspark.sql.functions as F
builder = SparkSession.builder.appName("DeltaLakeLocal") \
.enableHiveSupport() \
.config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
.config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog")
spark = configure_spark_with_delta_pip(builder).getOrCreate()
โช 1) Create Database
deltalake_dbCreates a new database named.
spark.sql("CREATE DATABASE IF NOT EXISTS deltalake_db")
spark.sql("SHOW DATABASES").show()
+------------+
| namespace|
+------------+
| default|
|deltalake_db|
+------------+
trainer_data.csv``trainerCreates a table to store the file's data .
query = f"""
CREATE TABLE IF NOT EXISTS deltalake_db.trainer (
id INT,
name STRING,
age INT,
hometown STRING,
prefer_type STRING,
badge_count INT,
level STRING
)
USING csv
OPTIONS (
path '[trainer_data.csv ํ์ผ ๊ฒฝ๋ก]',
header 'true',
inferSchema 'true',
delimiter ','
)
"""
spark.sql(query)
spark.sql("SHOW TABLES FROM deltalake_db").show()
+------------+---------+-----------+
| namespace|tableName|isTemporary|
+------------+---------+-----------+
|deltalake_db| trainer| false|
+------------+---------+-----------+
csvSince data from an existing file deltacannot be inserted directly into a table of the same type simultaneously with its creation, deltayou must create the table through two steps as follows.deltaCreate empty table of typedelta``csvInsert table data into the tabletrainer_delta``deltaCreates a table named .
/workspace/spark/deltalake/delta_local/trainer_delta/``delta
We have configured the settings so that table-related data is stored under the corresponding path .
query = f"""
CREATE TABLE IF NOT EXISTS deltalake_db.trainer_delta (
id INT,
name STRING,
age INT,
hometown STRING,
prefer_type STRING,
badge_count INT,
level STRING
)
USING delta
LOCATION '/workspace/spark/deltalake/delta_local/trainer_delta/'
"""
spark.sql(query)
Insert trainerthe data from the table created above into the table.trainer_delta
query = """
INSERT INTO deltalake_db.trainer_delta
SELECT * FROM deltalake_db.trainer;
"""
spark.sql(query)
Once data insertion is complete, _delta_log/you can verify that a new folder and parquet file have been created in the delta table directory.

LOCAL_DELTA_PATH = '/workspace/spark/deltalake/delta_local/trainer_delta'
df = spark.read.format("delta").load(LOCAL_DELTA_PATH)
df.show(5)
---
+---+------+---+--------+-----------+-----------+------------+
| id| name|age|hometown|prefer_type|badge_count| level|
+---+------+---+--------+-----------+-----------+------------+
| 1| Brian| 28| Seoul| Electric| 8| Master|
| 3| Susan| 18| Gwangju| Rock| 7| Expert|
| 6| Vicki| 17| Daejeon| Ice| 4|Intermediate|
| 9|Olivia| 45| Incheon| Psychic| 3|Intermediate|
| 10| Mark| 16| Gangwon| Fire| 4|Intermediate|
+---+------+---+--------+-----------+-----------+------------+
only showing top 5 rows
delta.Read as โช 2)query = f"SELECT * FROM delta.`{LOCAL_DELTA_PATH}`"
spark.sql(query)
spark.table('deltalake_db.trainer_delta')
# Beginner ์ ์ธํ dataframe ์์ฑ
df_1 = df.filter(F.col('level') != 'Beginner')
# ๊ธฐ์กด ๊ฒฝ๋ก์ ๋ฎ์ด์ฐ๊ธฐ
df_1.write \
.format('delta') \
.mode('overwrite') \
.save(LOCAL_DELTA_PATH)
# ๋ฐ์ดํฐ ํ์ธ
df = spark.read.format("delta").load(LOCAL_DELTA_PATH)
df.select('level').distinct().show()
---
+------------+
| level|
+------------+
| Expert|
| Advanced|
| Master|
|Intermediate|
+------------+

deltaView the change history for the table.
So far, the table has a structure where a total of three write operations have occurred since the initial creation (CREATE). Therefore, the versions also exist as 0, 1, 2, and 3.
trainer_deltaTable creation status, No datatrainerInitial state after data is inserted into the tablequery = "DESCRIBE HISTORY deltalake_db.trainer_delta"
spark.sql(query).show(vertical=True, truncate=False)
-RECORD 0--------------------------------------------------------------------------------------------------------------
version | 3
timestamp | 2025-03-21 02:09:39.085
userId | NULL
userName | NULL
operation | WRITE
operationParameters | {mode -> Overwrite, partitionBy -> []}
job | NULL
notebook | NULL
clusterId | NULL
readVersion | 2
isolationLevel | Serializable
isBlindAppend | false
operationMetrics | {numFiles -> 1, numOutputRows -> 42, numOutputBytes -> 3125}
userMetadata | NULL
engineInfo | Apache-Spark/3.5.1 Delta-Lake/3.2.0
-RECORD 1--------------------------------------------------------------------------------------------------------------
version | 2
timestamp | 2025-03-21 02:08:32.646
userId | NULL
userName | NULL
operation | WRITE
operationParameters | {mode -> Overwrite, partitionBy -> []}
job | NULL
notebook | NULL
clusterId | NULL
readVersion | 1
isolationLevel | Serializable
isBlindAppend | false
operationMetrics | {numFiles -> 1, numOutputRows -> 85, numOutputBytes -> 3868}
userMetadata | NULL
engineInfo | Apache-Spark/3.5.1 Delta-Lake/3.2.0
-RECORD 2--------------------------------------------------------------------------------------------------------------
version | 1
timestamp | 2025-03-21 01:46:21.446
userId | NULL
userName | NULL
operation | WRITE
operationParameters | {mode -> Append, partitionBy -> []}
job | NULL
notebook | NULL
clusterId | NULL
readVersion | 0
isolationLevel | Serializable
isBlindAppend | true
operationMetrics | {numFiles -> 1, numOutputRows -> 90, numOutputBytes -> 3980}
userMetadata | NULL
engineInfo | Apache-Spark/3.5.1 Delta-Lake/3.2.0
-RECORD 3--------------------------------------------------------------------------------------------------------------
version | 0
timestamp | 2025-03-21 01:43:25.471
userId | NULL
userName | NULL
operation | CREATE TABLE
operationParameters | {partitionBy -> [], clusterBy -> [], description -> NULL, isManaged -> false, properties -> {}}
job | NULL
notebook | NULL
clusterId | NULL
readVersion | NULL
isolationLevel | Serializable
isBlindAppend | true
operationMetrics | {}
userMetadata | NULL
engineInfo | Apache-Spark/3.5.1 Delta-Lake/3.2.0
๐ Load initial version (version 0) table
df_pre = spark.read \
.format("delta") \
.option("versionAsof", 0) \
.load(LOCAL_DELTA_PATH)
df_pre.select('level').distinct().show()
---
+-----+
|level|
+-----+
+-----+
๐ Load version 2 table
df_pre = spark.read \
.format("delta") \
.option("versionAsof", 2) \
.load(LOCAL_DELTA_PATH)
df_pre.select('level').distinct().show()
---
+------------+
| level|
+------------+
| Expert|
| Advanced|
| Master|
|Intermediate|
+------------+
๐ Querying Time Travel with SQL
df_pre = spark.sql("SELECT * FROM deltalake_db.trainer_delta VERSION AS OF 3")
df_pre.select('level').distinct().show()
---
+------------+
| level|
+------------+
| Expert|
| Master|
|Intermediate|
+------------+
๐ Load table status for a specified time period
TABLE_TIMESTAMP = "2025-03-21T02:09:00"
spark.read.format("delta") \
.option("timestampAsOf", TABLE_TIMESTAMP) \
.table("deltalake_db.trainer_delta")
๐ Retrieving tables for a specific time zone using SQL
TABLE_TIMESTAMP = "2025-03-21T02:09:00"
spark.sql(f"SELECT * FROM deltALake_db.trainer_delta TIMESTAMP AS OF '{TABLE_TIMESTAMP}'")
Delta Lake uses the 'Strict Schema Enforcement' option, which causes an error if you attempt to write data that differs from the existing table schema. Therefore, if the schema of an existing table changes, you must add specific options to enable overwriting.
The details of the practical training are as follows.
๐ 'dummy_col' ์ด๋ผ๋ ์ปฌ๋ผ์ด ์ถ๊ฐ๋์ด ์คํค๋ง๊ฐ ๋ณ๊ฒฝ๋ ํ ์ด๋ธ ๋ฎ์ด์ฐ๊ธฐ
LOCAL_DELTA_PATH = '/workspace/spark/deltalake/delta_local/trainer_delta'
# ํ
์ด๋ธ ๋ถ๋ฌ์ค๊ธฐ
df = spark.table("deltalake_db.trainer_delta")
# 'dummy_col' ์ปฌ๋ผ ์ถ๊ฐ
df_diff = df.withColumn('dummy_col', F.lit(1))
# ์คํค๋ง ํฉ์น๊ธฐ ์๋
df_diff.write \
.format('delta') \
.mode('overwrite') \
.save(LOCAL_DELTA_PATH)
# ๋ฎ์ด์ฐ๋ ค๋ ํ
์ด๋ธ์ ์คํค๋ง๊ฐ ๋ฌ๋ผ ์๋์ ์๋ฌ ๋ฐ์
# ๐๐๐๐๐
---
AnalysisException: [_LEGACY_ERROR_TEMP_DELTA_0007] A schema mismatch detected when writing to the Delta table (Table ID: 31dbae5e-d042-467b-9454-e483fdad97bb).
To enable schema migration using DataFrameWriter or DataStreamWriter, please set:
'.option("mergeSchema", "true")'.
For other operations, set the session configuration
spark.databricks.delta.schema.autoMerge.enabled to "true". See the documentation
specific to the operation for details.
Table schema:
root
-- id: integer (nullable = true)
-- name: string (nullable = true)
-- age: integer (nullable = true)
-- hometown: string (nullable = true)
-- prefer_type: string (nullable = true)
-- badge_count: integer (nullable = true)
-- level: string (nullable = true)
Data schema:
root
-- id: integer (nullable = true)
-- name: string (nullable = true)
-- age: integer (nullable = true)
-- hometown: string (nullable = true)
-- prefer_type: string (nullable = true)
-- badge_count: integer (nullable = true)
-- level: string (nullable = true)
-- dummy_col: integer (nullable = true)
To overwrite your schema or change partitioning, please set:
'.option("overwriteSchema", "true")'.
Note that the schema can't be overwritten when using
'replaceWhere'.
To overwrite a table with a different schema, option("mergeSchema", "true")you must add options.
df_diff.write \
.format('delta') \
.mode('overwrite') \
.option("mergeSchema", "true") \
.save(LOCAL_DELTA_PATH)
Optimization method
explanation
Basic Optimize
Merging small files improves read performance
Z-Ordering
Sort by frequently filtered columns โ Reduce scans and improve query performance
Partition-based Optimize
Selective optimization of only frequently accessed partitions, such as specific dates/regions
It applies the standard optimization method that Delta Lake performs by default.
query = "OPTIMIZE deltalake_db.trainer_delta"
spark.sql(query)
A function that optimizes the physical storage order of data based on specific columns.
Applying Z-Order based on frequently filtered columns can reduce unnecessary file scans during queries.
query = """
OPTIMIZE deltalake_db.trainer_delta
ZORDER BY (trainer_id, region)
"""
spark.sql(query)
Instead of optimizing for the entire table, it merges only a specific range of partitioned data.
query = """
OPTIMIZE deltalake_db.trainer_delta
WHERE level = 'Master'
"""
spark.sql(query)
deltaFor data of a specific type, historical file data should not be deleted directly from the directory but must be deleted VACUUMvia a command to avoid compromising table consistency and ensure smooth future operations.VACUUMWhen the operation occurs, it becomes impossible to Time Travel to a date prior to the deleted date to retrieve data.You must disable the default settings in Spark to manage the retention period customly.
spark.conf.get("spark.databricks.delta.retentionDurationCheck.enabled")
-> 'true'
spark.conf.set("spark.databricks.delta.retentionDurationCheck.enabled", "false")
VACUUMExecute the command to delete parquet files created before the user-specified period .
DRY RUNThe option is a setting that prevents the operation from actually happening.
spark.sql("VACUUM deltalake_db.trainer_delta").show(truncate=False)
spark.sql("VACUUM deltalake_db.trainer_delta RETAIN 0 HOURS DRY RUN").show(truncate=False)
spark.sql("VACUUM deltalake_db.trainer_delta RETAIN 2 DAYS DRY RUN").show(truncate=False)
deltaSets the default retention period when creating a type table.
query = f"""
CREATE TABLE IF NOT EXISTS deltalake_db.trainer_delta_2 (
id INT,
name STRING,
age INT,
hometown STRING,
prefer_type STRING,
badge_count INT,
level STRING
)
USING delta
LOCATION '/workspace/spark/deltalake/delta_local/trainer_delta_2/'
TBLPROPERTIES ('delta.deletedFileRetentionDuration' = 'interval 2 days');
"""
spark.sql(query)
parquet``deltaThis is a function that converts data stored in a specific format into table data of a specific type.๐ fish_data.csvSaves data as parquet.
# csv ํ์ผ ์ฝ์ด์ค๊ธฐ
fish = spark.read.option('header', 'true').csv('fish_data.csv')
# ๋ก์ปฌ ๋๋ ํ ๋ฆฌ ์ ์ฅ + Catalog ์ ์ฅ
fish.write \
.mode('overwrite') \
.format('parquet') \
.option('path', '/workspace/spark/deltalake/delta_local/fish_parquet/') \
.saveAsTable('deltalake_db.fish_parquet')
Convert parquetdata stored as ๐ todelta
query = """
CONVERT TO DELTA
parquet.`/workspace/spark/deltalake/delta_local/fish_parquet/`
"""
spark.sql(query)
๐ SpeciesWrite column-partitioned Parquet data
fish_df.write.mode('overwrite')\
.format('parquet') \
.partitionBy('Species') \
.option('path', '/workspace/spark/deltalake/delta_local/fish_parquet_partitioned/') \
.saveAsTable('deltalake_db.fish_parquet_partitioned')
delta๐ Convert partitioned parquet data to
query = """
CONVERT TO DELTA
parquet.`/workspace/spark/deltalake/delta_local/fish_parquet_partitioned/`
PARTITIONED BY (Species STRING)
"""
spark.sql(query)