๐ŸŒŠ Delta Lake Beginner's Guide - Practical Edition (Part 1. Local Environment)

Abdul whabยท3์ผ ์ „

post-thumbnail

  • 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)

    0. INTRO

    • 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

    1๏ธโƒฃ Setting up the Practice Environment

    โ–ช 1) Create a Docker Container

    • /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.


    โ–ช 2) Install Library

    • 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

  • โ–ช 3) Pyspark Delta Lake Environment Setup

    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()

2๏ธโƒฃ Create Database and Table

โ–ช 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|
    +------------+

ย 2) Create CSV type table

  • trainer_data.csv``trainerCreates a table to store the file's data .

    • ํ…Œ์ด๋ธ” ์ด๋ฆ„ : trainer
    • ์Šคํ‚ค๋งˆ :
      • id โ†’ INT
      • name โ†’ STRING
      • age โ†’ INT
      • hometown โ†’ STRING
      • prefer_type โ†’ STRING
      • badge_count โ†’ INT
      • level โ†’ STRING

    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|
    +------------+---------+-----------+

3๏ธโƒฃ Create delta type table

  • 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.
    1. deltaCreate empty table of type
    2. delta``csvInsert table data into the table

โ–ช 1) Create table

  • trainer_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)

โ–ช 2) Insert data

  • 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.

4๏ธโƒฃ Reading delta type tables

  • In PySpark, you can read tables stored in slightly different ways.

โ–ช 1) Basic reading in Spark

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)

โ–ช 3) Read from Hive Catalog

spark.table('deltalake_db.trainer_delta')

5๏ธโƒฃ Save after editing the table

  • This time, we will modify the table contents and overwrite the existing directory. This step is intended to prepare for future table change history inquiries and to practice Time Travel queries, a core feature of Delta Lake.

โ–ช 1) Save after excluding 'Beginner'

# 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|
+------------+

6๏ธโƒฃ Change History Lookup and Time Travel Query

โ–ช 1) View History

  • 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.

    • VERSION 0 โ†’ trainer_deltaTable creation status, No data
    • VERSION 1 โ†’ trainerInitial state after data is inserted into the table
    • VERSION 2 โ†’ Saved state after excluding 'Beginner' rows
    • VERSION 3 โ†’ Saved state after excluding 'Advanced' row

    query = "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

โ–ช 2) Time Travel - Version

  • Whenever the table is modified, the contents of the table corresponding to a specific version are retrieved based on the assigned version number.

๐Ÿ‘‰ 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|
+------------+

โ–ช 3) Time Travel - Timestamp

  • This is a query method based on table modification time, retrieving the status (version) of the Delta table that existed at a specified time (TIMESTAMP).

๐Ÿ‘‰ 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}'")

7๏ธโƒฃ Schema Change Task

  • 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.

    • ๊ธฐ์กด ์ปฌ๋Ÿผ : ['id', 'name', 'age', 'hometown', 'prefer_type', 'badge_count', 'level']
    • ๋ณ€๊ฒฝ๋œ ํ…Œ์ด๋ธ” ์ปฌ๋Ÿผ : ['id', 'name', 'age', 'hometown', 'prefer_type', 'badge_count', 'level', 'dummy_col']

    ๐Ÿ‘‰ 'dummy_col' ์ด๋ผ๋Š” ์ปฌ๋Ÿผ์ด ์ถ”๊ฐ€๋˜์–ด ์Šคํ‚ค๋งˆ๊ฐ€ ๋ณ€๊ฒฝ๋œ ํ…Œ์ด๋ธ” ๋ฎ์–ด์“ฐ๊ธฐ

โ–ช 1) Standard write - Operation failed

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'.

ย 2) Use with the schema merging option

  • 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)

Optimize file status

  • https://docs.databricks.com/aws/en/sql/language-manual/delta-optimize
  • Since Delta tables are fundamentally structured as continuously loaded files, a large number of small files accumulate over time. This leads to degraded query performance and increased read overhead. In this case, performance can be improved by merging data into a larger file using OPTIMIZE.

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

ย 1) Standard Optimization

  • It applies the standard optimization method that Delta Lake performs by default.

    query = "OPTIMIZE deltalake_db.trainer_delta"

    spark.sql(query)

2) Z-Ordering Optimization

  • 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)

โ–ช 3) Partition Optimization

  • 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)


9๏ธโƒฃ Delete past data (VACUUM)

  • https://docs.databricks.com/aws/en/sql/language-manual/delta-vacuum
  • In Delta Lake, past parquet files remain even if data is modified or deleted.
  • Since it is wasteful to keep past versions of files that are no longer used, you can use the VACUUM function to delete data older than a specific period.
  • 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.
  • The default retention period for files is 168 hours (7 days), and you can adjust the retention period by modifying the spark config.

โ–ช 1) Disable default retention period

  • 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")

โ–ช 2) Execute the VACUUM command

  • 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.

    ๊ธฐ๋ณธ VACUUM ๋ช…๋ น (168์‹œ๊ฐ„ ์ด์ „์˜ ํŒŒ์ผ ์‚ญ์ œ)

    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)

    2์ผ ์ด์ „์˜ ํŒŒ์ผ๋“ค ์‚ญ์ œ

    spark.sql("VACUUM deltalake_db.trainer_delta RETAIN 2 DAYS DRY RUN").show(truncate=False)

3) Set retention period when creating table

  • 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 to Delta Conversion

โ–ช 1) General Parquet Data Conversion

๐Ÿ‘‰ 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)

โ–ช 2) Converting partitioned parquet data

๐Ÿ‘‰ 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)

Reference materials

profile
12546765678900897

0๊ฐœ์˜ ๋Œ“๊ธ€