DeltaTable alapok (PySpark)
1. Alapelvek (Principles)
- Delta Lake nem egy új fájlformátum: Nem egyedi bináris struktúra, hanem standard, oszlopalapú Parquet fájlok és egy tranzakciós napló (
_delta_log) összessége. - ACID tranzakciók garantálása:
- Atomicitás: A módosítás vagy teljes egészében befejeződik, vagy egyáltalán nem (nincs félbehagyott állapot).
- Konzisztencia: Nincsenek piszkos olvasások (dirty reads), a sémaszabályok mindig érvényesülnek.
- Izoláció: Párhuzamos írások és olvasások nem zavarják egymást (optimista párhuzamosság-vezérlés).
- Tartósság: A sikeres tranzakció azonnal perzisztenssé válik a felhőtárhelyen.
- Parquet megváltoztathatatlansága (Immutability): A meglévő Parquet adatfájlok tartalma közvetlenül sosem írható felül.
2. Működési Minták (Patterns)
- Írási műveletek:
- INSERT: Létrehoz egy új
part-*.parquetfájlt, majd bejegyzi a tranzakciót egy sorszámozott.jsonnaplófájlba. - UPDATE / DELETE (Copy-on-Write minta): A módosítandó rekordokat tartalmazó Parquet fájl lemásolódik, létrejön az új verzió a módosításokkal, a régi fájl pedig a naplóban logikailag töröltnek (
remove) minősül.
- Checkpointing minta:
- A motor 10 tranzakciónként összefogja az addigi JSON naplókat, és generál egy
.checkpoint.parquetállapotmentést. - Ezzel elkerüli a több ezer JSON fájl egyenkénti újraolvasásának (replay) teljesítményveszteségét.
- Karbantartási minták:
- OPTIMIZE (Kompakció): Sok kisméretű fájlt egyesít kevesebb, optimális méretű Parquet fájlba a jobb lekérdezési sebességért.
- VACUUM (Fizikai törlés): Véglegesen törli a felhőtárhelyről a logikailag már eltávolított fájlokat a megadott megőrzési idő (retention period) után.
3. Fájltípusok és Rendszerezett Magyarázatok
| Fájlkiterjesztés / Mappa | Tartalom | Cél és Szerepkör |
sales_delta/ | Gyökérkönyvtár | A Delta tábla teljes tárolási mappája. |
part-*.parquet | Nyers adatrekordok | A tényleges adatokat oszlopalapon tároló, tömörített fájlok. |
_delta_log/ | Metaadat-tárhely | A tábla központi tranzakciós naplója. |
*.json | Tranzakciós napló | Egy-egy konkrét tranzakció műveletei (add/remove fájlok, commitInfo, protocol, séma). |
*.checkpoint.parquet | Összesített állapot | A tábla adott verziójának összesített metaadat-pillanatképe a gyors betöltéshez. |
*.crc | Ellenőrző összegek | Belső integritás-ellenőrzésre szolgáló fájlok. |
4. Kapcsolatok és Összefüggések (Relationships)
- Adatfájlok és Metaadat kapcsolata:
- A tábla aktuális állapota = Az érvényes Checkpoint + a rákövetkező JSON logok alapján aktívnak jelölt Parquet fájlok összessége.
- Tranzakciós napló és Time Travel kapcsolata:
- A Time Travel nem másolatokat készít a teljes tábláról, hanem a JSON naplók korábbi bejegyzései (
VERSION AS OF/TIMESTAMP AS OF) alapján rekonstruálja, hogy az adott pillanatban mely Parquet fájlok voltak érvényesek.
- OPTIMIZE és Táblaméret kapcsolata:
- Az
OPTIMIZEfuttatása ideiglenesen megnöveli a tárolt adatméretet, mivel az új, összevont fájlok mellett a régi kis fájlok is a tárhelyen maradnak a logikai törlés után.
5. Kivételek, Korlátok és Kockázatok (Exceptions & Caveats)
- Logikai vs. Fizikai törlés:
- A
DELETEparancs nem szabadít fel lemezterületet azonnal, csupán a metaadatban jelöli a fájlt töröltnek.
- Time Travel veszteség a VACUUM miatt:
- A
VACUUMfuttatása után a megőrzési időn kívül eső állapotok nem állíthatók vissza Time Travellel, mert a mögöttes Parquet fájlok fizikailag megsemmisülnek.
- Megőrzési idő (Retention period) kockázata:
- Alapértelmezett érték: 7 nap.
- Ennek csökkentése (pl. 0 napra kényszerítés) adatvesztést vagy futó olvasási tranzakciók összeomlását okozhatja.
- Séma evolúció korlátai:
- A Delta támogatja a séma kiterjesztését (
Schema Evolution), de az alapértelmezett viselkedés a szigorú kikényszerítés (Schema Enforcement), ami megakadályozza a nem illeszkedő adatok beírását.
6. Részletes Esettanulmány: Adatmódosítási Életciklus
Forgatókönyv: Egyetlen rekord frissítése (UPDATE)
1. Elemzési fázis: A lekérdező motor a metaadatok és partíciós statisztikák alapján azonosítja a célrekordot tartalmazó fájlt (pl.
part-0001.parquet).2. Új fájl generálása: A Spark beolvassa a
part-0001.parquettartalmát a memóriába, elvégzi a rekord módosítását, majd kiírja a teljes adathalmazt egy új fájlba (pl.part-0005.parquet).3. Tranzakció commitálása: Létrejön a következő naplófájl (pl.
00000000000000000011.json), amely két kulcsfontosságú műveletet tartalmaz:remove:part-0001.parquetadd:part-0005.parquet
4. Állapotfrissítés: A következő olvasási művelet már a
part-0005.parquetfájlt tekinti érvényesnek; apart-0001.parquetinaktívvá válik, de a lemezen marad.
7. Haladó Teljesítményoptimalizálási Minták
Z-Ordering (Többdimenziós rendezés):
Az
OPTIMIZE table_name ZORDER BY (col_a, col_b)parancs úgy rendezi át a Parquet fájlokon belüli adatokat, hogy minimalizálja az olvasandó fájlok számát szűréskor (Data Skipping).
Data Skipping és Fájl statisztikák:
Minden
.jsoncommit automatikusan tárolja az oszlopok minimum/maximum értékeit és null-érték számait az első 32 oszlopra.A motor ezek alapján képes teljes Parquet fájlok beolvasását kihagyni a lekérdezés futtatásakor.
Change Data Feed (CDF):
Lehetővé teszi a sor-szintű változások (beszúrás, törlés, frissítés előtti/utáni állapot) lekérdezését verziók között, csökkentve a downstream ETL folyamatok terhelését.
8. Architektúra-üzemeltetési és Hiba-elhárítási Szabályok
| Probléma / Jelenség | Kiváltó Ok | Megoldás / Bevált Gyakorlat |
| Small File Problem | Túl sok streaming mikrobatch vagy apró írási művelet. | Rendszeres OPTIMIZE futtatás vagy Auto-Compaction bekapcsolása. |
| Tárhelyköltség növekedése | Felhalmozódó inaktív fájlok a gyakori UPDATE/DELETE miatt. | Ütemezett VACUUM futtatása a megőrzési idő figyelembevételével. |
| Párhuzamossági ütközés (Concurrent Append / Write Conflict) | Több párhuzamos folyamat ugyanazt a partíciót vagy fájlt módosítja. | Particionálás optimalizálása, műveletek szétválasztása vagy Retry logika beépítése. |
| Lassú metaadat-olvasás | Túl sok egyedi .json fájl a _delta_log könyvtárban. | A Checkpoint mechanizmus ellenőrzése és a naplók tömörítése. |
1. PySpark és SQL Alapműveletek Delta Lake-en
Tábla inicializálása és Adatírás
SQL szintaxis:
CREATE TABLE sales_delta (
order_id BIGINT,
customer_id STRING,
amount DOUBLE,
country STRING,
order_date DATE
)
USING DELTA
PARTITIONED BY (country)
LOCATION '/mnt/lakehouse/sales_delta';
PySpark DataFrame szintaxis:
df.write \
.format("delta") \
.mode("append") \
.partitionBy("country") \
.save("/mnt/lakehouse/sales_delta")
Módosítás és Törlés (UPDATE / DELETE)
SQL szintaxis:
-- Rekord frissítése
UPDATE sales_delta
SET amount = amount * 1.1
WHERE country = 'US' AND order_date >= '2026-01-01';
-- Logikai törlés
DELETE FROM sales_delta
WHERE country = 'DE' AND amount <= 0;
PySpark (
DeltaTableAPI) szintaxis:
from delta.tables import DeltaTable
delta_table = DeltaTable.forPath(spark, "/mnt/lakehouse/sales_delta")
# UPDATE művelet
delta_table.update(
condition = "country = 'US' AND order_date >= '2026-01-01'",
set = { "amount": "amount * 1.1" }
)
# DELETE művelet
delta_table.delete(condition = "country = 'DE' AND amount <= 0")
Upsert / Merge Művelet (CDC és dimenziófrissítések mintája)
SQL szintaxis:
MERGE INTO sales_delta AS target
USING staging_sales AS source
ON target.order_id = source.order_id
WHEN MATCHED THEN
UPDATE SET
target.amount = source.amount,
target.order_date = source.order_date
WHEN NOT MATCHED THEN
INSERT (order_id, customer_id, amount, country, order_date)
VALUES (source.order_id, source.customer_id, source.amount, source.country, source.order_date);
PySpark szintaxis:
delta_table.alias("target") \
.merge(
staging_df.alias("source"),
"target.order_id = source.order_id"
) \
.whenMatchedUpdate(set = {
"amount": "source.amount",
"order_date": "source.order_date"
}) \
.whenNotMatchedInsert(values = {
"order_id": "source.order_id",
"customer_id": "source.customer_id",
"amount": "source.amount",
"country": "source.country",
"order_date": "source.order_date"
}) \
.execute()
Időutazás (Time Travel) és Visszaállítás
SQL szintaxis:
-- Verziószám alapján
SELECT * FROM sales_delta VERSION AS OF 3;
-- Időbélyeg alapján
SELECT * FROM sales_delta TIMESTAMP AS OF '2026-08-20 10:00:00';
-- Tábla visszaállítása egy korábbi verzióra
RESTORE TABLE sales_delta TO VERSION AS OF 2;
PySpark szintaxis:
# Verzió lekérdezése
df_v3 = spark.read.format("delta").option("versionAsOf", 3).load("/mnt/lakehouse/sales_delta")
# Időbélyeg lekérdezése
df_ts = spark.read.format("delta").option("timestampAsOf", "2026-08-20 10:00:00").load("/mnt/lakehouse/sales_delta")
# Visszaállítás
delta_table.restoreToVersion(2)
Karbantartási Parancsok (OPTIMIZE és VACUUM)
SQL szintaxis:
-- Kis fájlok tömörítése
OPTIMIZE sales_delta;
-- 7 napnál régebbi inaktív fájlok fizikai törlése
VACUUM sales_delta RETAIN 168 HOURS;
PySpark szintaxis:
# Optimize
delta_table.optimize().executeCompaction()
# Vacuum
delta_table.vacuum(168)
2. Z-Ordering és Data Skipping Belső Működési Mechanizmusa
A Probléma: Hagyományos Rendezés vs. Többdimenziós Lekérdezések
Lineáris rendezés (Lexicographical sort): Ha az adatokat több oszlop szerint rendezzük (pl.
ORDER BY col_a, col_b), a rendezés dominánsan az első oszlopra (col_a) érvényesül.Skálázódási korlát: Ha a lekérdezés csak a második oszlopra szűr (
WHERE col_b = 'X'), a hagyományos rendezés nem tudja kihasználni az oszlopstatisztikákat, és szinte a teljes adathalmazt be kell olvasni.
A Megoldás: Z-Order Görbe (Morton-kódok)
Térkitöltő görbe (Space-filling curve): A Z-ordering többdimenziós adatokat képez le egydimenziós térbe úgy, hogy a több dimenzióban egymáshoz közeli pontok az egydimenziós sorrendben (és így a fájlokban) is közel maradjanak egymáshoz.
Bit-interleaving (Bitek összefésülése):
A motor az indexelt oszlopok értékeit bináris formára alakítja.
Az oszlopok bináris bitjeit felváltva fésüli össze egyetlen kulccsá (pl. $A$ oszlop 1. bitje, $B$ oszlop 1. bitje, $A$ 2. bitje, $B$ 2. bitje...).
Az így kapott érték alapján rendezi a sorokat a Parquet fájlokba írás előtt.
Eredmény: Bármelyik megadott oszlopra (vagy azok kombinációjára) történik a szűrés, a keresési tartomány hatékonyan szűkíthető.
Data Skipping a Gyakorlatban
Fájlszintű metaadatok: Minden Parquet fájlhoz a
_delta_logtárolja az oszlopok minimum és maximum értékeit (stats: {"minValues": {...}, "maxValues": {...}, "numRecords": ...}).Fájlok kihagyása (File Pruning):
Ha a lekérdezés feltétele
WHERE customer_id = 'C123', a lekérdező motor összeveti ezt az értéket a JSON/Checkpoint naplóban lévő min/max értékekkel.Ha a keresett érték kívül esik a fájl intervallumán, a Spark az adott Parquet fájlt egyáltalán nem tölti le és nem olvassa be a tárhelyről.
Szintaxis és Gyakorlati Alkalmazás
SQL szintaxis:
OPTIMIZE sales_delta
ZORDER BY (order_date, customer_id);
PySpark szintaxis:
delta_table.optimize().executeZOrderBy("order_date", "customer_id")
Tervezési Szabályok és Korlátok
Oszlopok száma: Legfeljebb 2–4 nagy kardinalitású oszlopra érdemes alkalmazni; túl sok dimenzió esetén a térkitöltő görbe hatékonysága drasztikusan csökken.
Particionálás vs. Z-Ordering:
Alacsony kardinalitás (pl. év, ország): Particionálás (
PARTITIONED BY).Magas kardinalitás (pl. ügyfél-azonosító, dátum, tranzakció-azonosító): Z-Ordering.
Statisztika gyűjtési limit: A Delta Lake alapértelmezetten az első 32 oszlopra gyűjt statisztikát (
spark.databricks.delta.properties.defaults.dataSkippingNumIndexedCols). Ha a Z-Order oszlop a 32. oszlop után van, a Data Skipping nem működik, kivéve ha ezt a limitet manuálisan megemeljük.
3. Change Data Feed (CDF) Belső Működése és Konfigurációja
A CDF Célja és Működési Alapelve
Cél: Sor-szintű változáskövetés (Change Data Capture – CDC) biztosítása anélkül, hogy a teljes táblát újra kellene olvasni vagy verziókat kellene manuálisan összehasonlítani (
diff).Működési mechanizmus:
Bekapcsolásakor a Delta Lake az adatfájlok mellett egy rejtett
_change_datamappába speciális CDC fájlokat generál azUPDATEésDELETEműveleteknél.Az
INSERTműveletek nem generálnak extra CDC fájlokat; a motor közvetlenül az újonnan hozzáadott Parquet fájlokból olvassa ki az új sorokat.
CDF Bekapcsolása
Új tábla létrehozásakor (SQL):
CREATE TABLE sales_delta (
order_id BIGINT,
customer_id STRING,
amount DOUBLE,
country STRING,
order_date DATE
)
USING DELTA
TBLPROPERTIES (delta.enableChangeDataFeed = true);
Meglévő táblán (SQL / Python):
ALTER TABLE sales_delta
SET TBLPROPERTIES (delta.enableChangeDataFeed = true);
CDF Metaadat-oszlopok Struktúrája
A CDF lekérdezésekor a tábla eredeti oszlopai mellett 4 kiegészítő metaadat-oszlop jelenik meg:
| Oszlop neve | Típus | Leírás / Lehetséges értékek |
| _change_type | STRING | A változás típusa: insert, update_preimage, update_postimage, delete |
| _commit_version | BIGINT | A tranzakció naplósorszáma, amelyben a változás történt |
| _commit_timestamp | TIMESTAMP | A tranzakció pontos végrehajtási időbélyege |
Egy UPDATE művelet mindig két sort generál: egy update_preimage (módosítás előtti állapot) és egy update_postimage (módosítás utáni állapot) bejegyzést.
Változások Lekérdezése (Batch és Streaming)
Batch lekérdezés adott verziók között (SQL):
SELECT *
FROM table_changes('sales_delta', 2, 5)
WHERE _change_type IN ('insert', 'update_postimage');
Batch lekérdezés időbélyeg tartományra (PySpark):
changes_df = spark.read.format("delta") \
.option("readChangeFeed", "true") \
.option("startingTimestamp", "2026-08-20 08:00:00") \
.option("endingTimestamp", "2026-08-20 18:00:00") \
.load("/mnt/lakehouse/sales_delta")
Inkrementális streaming feldolgozás (Structured Streaming):
streaming_changes = spark.readStream.format("delta") \
.option("readChangeFeed", "true") \
.option("startingVersion", 1) \
.load("/mnt/lakehouse/sales_delta")
Gyakorlati Szabályok és Tervezési Szempontok
Tárhely- és Írási Többlet (Overhead): Az
UPDATEésDELETEműveletek során a CDF miatt extra adatkiírás történik a_change_datakönyvtárba, ami némi írási késleltetést és tárhelynövekedést okoz.VACUUM viselkedés: A
VACUUMparancs a_change_datakönyvtárban található inaktív fájlokat is törli a megőrzési idő lejárta után; ezen túlmenően a korábbi változások már nem olvashatók.
4. Párhuzamosság-vezérlés és Tranzakciós Ütközések (Concurrency Control)
Optimista Párhuzamosság-vezérlés (OCC - Optimistic Concurrency Control)
Alapelv: A Delta Lake feltételezi, hogy a párhuzamos műveletek többsége nem ütközik egymással, így az író folyamatok nem zárolják le előre a teljes táblát.
Működési lépések:
1. Snapshot készítés: A tranzakció indulásakor rögzíti a tábla aktuális verzióját (pl. Version 10).
2. Művelet végrehajtása: Elvégzi a számításokat és kiírja az új Parquet fájlokat a háttértárra.
3. Ellenőrzés és Commit kísérlet: Megkísérli a
00000000000000000011.jsonfájl létrehozását az objektumtár atomi put-if-absent műveletével.4. Ütközéskezelés (Conflict Resolution): Ha egy másik tranzakció közben már létrehozta a 11-es verziót, a motor automatikusan ellenőrzi, hogy a két tranzakció azonos fájlokat vagy partíciókat érintett-e. Ha nem, automatikus újrapróbálkozást (retry) hajt végre; ha igen,
ConcurrentAppendExceptionvagyConcurrentModificationExceptionhibával leáll.
Műveleti Kompatibilitási Mátrix
| Művelet 1 | Művelet 2 | Eredmény / Ütközési szabály |
| APPEND | APPEND | Mindig összefér (nincs konfliktus, automatikus újra-commit). |
| APPEND | UPDATE / DELETE / MERGE | Összefér, kivéve ha az UPDATE/DELETE szűrési feltétele az újonnan hozzáadott partíciót/fájlt érinti. |
| UPDATE / DELETE | UPDATE / DELETE | Konfliktus lép fel, ha mindkét folyamat ugyanazt a Parquet fájlt próbálja módosítani vagy törölni. |
| OPTIMIZE | UPDATE / DELETE / MERGE | Általában kompatibilis, az OPTIMIZE automatikusan újratervezi a tömörítést. |
5. Séma-menedzsment: Kikényszerítés és Evolúció
Séma Kikényszerítés (Schema Enforcement / Validation)
Cél: Megakadályozza a hibás, nem kompatibilis típusú vagy váratlan oszlopokat tartalmazó adatok beírását.
Alapértelmezett viselkedés: Ha az írni kívánt DataFrame sémája nem egyezik meg pontosan a Delta tábla sémájával, a Spark azonnal kivételt (
AnalysisException) dob, és a tranzakció megszakad.
Séma Evolúció (Schema Evolution)
Cél: Új oszlopok automatikus hozzáadása a tábla metaadataihoz anélkül, hogy a meglévő Parquet fájlokat módosítani kellene.
Működés: A korábbi Parquet fájlokban az új oszlop automatikusan
NULLértéket kap olvasáskor.Megvalósítás PySparkban:
df_new_columns.write \
.format("delta") \
.mode("append") \
.option("mergeSchema", "true") \
.save("/mnt/lakehouse/sales_delta")
Megvalósítás SQL-ben (
MERGEműveletnél):
-- Automatikus sémabővítés bekapcsolása session szinten
SET spark.databricks.delta.schema.autoMerge.enabled = true;
MERGE INTO sales_delta AS target
USING staging_sales AS source
ON target.order_id = source.order_id
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *;
Séma Módosítási Kivételek és Korlátok
Megengedett automatikus módosítások: Új oszlopok hozzáadása,
NULLmező típusának kompatibilis szigorítása (pl.ByteType$\rightarrow$IntegerType).Tiltott módosítások (Manuális beavatkozást igénylő műveletek):
Oszlop törlése vagy átnevezése (külön
ALTER TABLE ... RENAME COLUMNvagy Column Mapping támogatás szükséges).Inkompatibilis adattípus-váltás (pl.
StringType$\rightarrow$IntegerTypevagyDouble$\rightarrow$Long).
6. Column Mapping és Speciális Karakterek Támogatása
A Column Mapping Célja és Működési Alapelve
Cél: Lehetővé teszi az oszlopok átnevezését, törlését és nem szabványos karakterek (pl. szóközök, vesszők, pontosvesszők) használatát az oszlopnevekben anélkül, hogy a mögöttes Parquet adatfájlokat újra kellene írni.
Működési mechanizmus:
Szétválasztja a logikai oszlopnevet (amit a felhasználó lát) a fizikai oszlopazonosítótól (amit a Parquet fájl belső sémája tárol).
A napló metaadatai egyedi belső azonosítót (
id) és fizikai nevet rendelnek minden mezőhöz.
Column Mapping Bekapcsolása és Használata
Bekapcsolás meglévő táblán (SQL):
ALTER TABLE sales_delta
SET TBLPROPERTIES (
'delta.columnMapping.mode' = 'name',
'delta.minReaderVersion' = '2',
'delta.minWriterVersion' = '5'
);
Oszlop átnevezése és törlése (Metaadat-szintű művelet):
-- Oszlop átnevezése (nem írja újra a Parquet fájlokat)
ALTER TABLE sales_delta RENAME COLUMN customer_id TO client_id;
-- Oszlop logikai törlése
ALTER TABLE sales_delta DROP COLUMN country;
7. Liquid Clustering (Új Generációs Adatrendezés)
Miért váltja le a Hagyományos Particionálást és a Z-Orderinget?
A korábbi korlátok:
Hive-stílusú particionálás: Túlparticionálás esetén fellép a Small File Problem, és a partíciós oszlopok utólagos módosítása a tábla teljes újraírását igényli.
Z-Ordering: Nem inkrementális (minden futtatáskor újrarendezi a teljes érintett adathalmazt), nem támogatja a rendezési kulcsok egyszerű, menet közbeni módosítását.
A Liquid Clustering előnyei:
Rugalmasság: A klaszterezési kulcsok a tábla újraírása nélkül bármikor módosíthatók.
Inkrementális optimalizálás: Csak az újonnan beérkezett és a még nem klaszterezett adatokat rendezi át, jelentősen csökkentve az
OPTIMIZEszámítási költségét.Automata ferdeségkezelés (Data Skew): Dinamikusan igazítja a fájlok méretét anélkül, hogy merev könyvtárstruktúrákat (
/country=US/) hozna létre a felhőtárhelyen.
Liquid Clustering Szintaxis és Gyakorlati Műveletek
Tábla létrehozása klaszterezéssel (SQL):
CREATE TABLE sales_cluster
(
order_id BIGINT,
client_id STRING,
amount DOUBLE,
order_date DATE
)
USING DELTA
CLUSTER BY (order_date, client_id);
Klaszterezési kulcsok módosítása meglévő táblán:
ALTER TABLE sales_cluster
CLUSTER BY (order_date, amount);
Inkrementális karbantartás futtatása:
OPTIMIZE sales_cluster;
8. Összegző Architektúra-Tervezési Szabályok
Particionálás vs. Liquid Clustering: Új projekteknél particionálás helyett a Liquid Clustering az ajánlott megközelítés a dinamikus és skálázható adatkezeléshez.
Írási stratégia: Gyakori apró írások (streaming) esetén célszerű az Auto-Compaction és Optimized Writes beállításokat bekapcsolni a
_delta_logés a fájlstruktúra túlterhelésének megelőzésére.Biztonsági karbantartás: A
VACUUMparancsot mindig a downstream folyamatok olvasási ablakához és az elvárt Time Travel igényekhez igazított megőrzési idővel (RETAIN X HOURS) kell ütemezni.
Megjegyzések
Megjegyzés küldése