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.delta_data.zipwe used data from two CSV files located in the file below.
๐ delta_data.zip
/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
8888you can access the juypter lab development environment by connecting to the port.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 Fbuilder = 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()
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.csvtrainerCreates 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|
+------------+---------+-----------+
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 typedeltacsvInsert table data into the tabletrainer_deltadeltaCreates a table named ./workspace/spark/deltalake/delta_local/trainer_delta/deltaWe 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)
trainerthe data from the table created above into the table.trainer_deltaquery = """
INSERT INTO deltalake_db.trainer_delta
SELECT * FROM deltalake_db.trainer;
"""
spark.sql(query)
_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|
+------------+
# Advanced ์ ์ธํ dataframe ์์ฑ df_2 = df_1.filter(F.col('level') != 'Advanced')# ๊ธฐ์กด ๊ฒฝ๋ก์ ๋ฎ์ด์ฐ๊ธฐ
df_2.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|
| Master|
|Intermediate|
+------------+
_delta_log/metadata (files) are also added within the folder ..json
deltaView the change history for the table.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}'")
- ๊ธฐ์กด ์ปฌ๋ผ : ['id', 'name', 'age', 'hometown', 'prefer_type', 'badge_count', 'level'] - ๋ณ๊ฒฝ๋ ํ ์ด๋ธ ์ปฌ๋ผ : ['id', 'name', 'age', 'hometown', 'prefer_type', 'badge_count', 'level', 'dummy_col']
๐ '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'.
option("mergeSchema", "true")you must add options.df_diff.write \
.format('delta') \
.mode('overwrite') \
.option("mergeSchema", "true") \
.save(LOCAL_DELTA_PATH)
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 |
query = "OPTIMIZE deltalake_db.trainer_delta"
spark.sql(query)
query = """ OPTIMIZE deltalake_db.trainer_delta ZORDER BY (trainer_id, region) """
spark.sql(query)
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.# ์ค์ ํ์ธ 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.# ๊ธฐ๋ณธ 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)
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)
parquetdeltaThis 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)