Adatok szelektív felülírása a Delta Lake használatával

A Delta Lake a következő különböző lehetőségeket kínálja a szelektív felülíráshoz:

Option Felhasználási eset Támogatott számítási típusok Minimális verzió
REPLACE WHERE Atomilag felülírja a predikátumnak megfelelő sorokat. Olyan cserékhez használható, amelyekhez rögzített egyező feltétel tartozik, például colA = 5 vagy int_col IN (1, 2, 3). Minden számítási típus. SQL a Databricks Runtime 12.2 LTS-ben és újabb verziókban. Python és Scala a Databricks Runtime 9.1 LTS-ben és újabb verziókban.
REPLACE USING Dinamikus adatok felülírása. A megadott oszlopoknak megfelelő összes sort lecseréli az oszlopértékek egyenlőségi összehasonlítása alapján a megadott adatkészletben. Minden számítási típus. SQL a Databricks Runtime 16.3-ban és újabb verziókban. Python és Scala a Databricks Runtime 18.2-ben és újabb verziókban.
REPLACE ON Dinamikus adatok felülírása logikai kifejezéssel. Használja helyettesítésekhez összetett vagy NULL-biztonságos illesztési feltétel esetén, például a(z) s.colA <=> t.colA AND s.colB <=> t.colB használatával. Minden számítási típus. SQL a Databricks Runtime 17.1-ben és újabb verziókban. Python és Scala a Databricks Runtime 18.2-ben és újabb verziókban.
partitionOverwriteMode Örökölt dinamikus partíciófelülírás, amely minden olyan partícióban felülírja az összes meglévő adatot, amelybe az írási művelet új adatokat fog véglegesíteni. Új számítási feladatokhoz nem ajánlott. Az SQL csak a klasszikus számításokat támogatja. Python és a Scala minden számítási típust támogat. SQL, Python és Scala a Databricks Runtime 11.3 LTS-ben és újabb verziókban.

A legtöbb felhasználási esetben a Databricks a REPLACE USING vagy a REPLACE WHERE használatát javasolja. Csak akkor használja REPLACE ON , ha a használati eset összetett vagy NULL-biztonságos egyeztetési feltételeket igényel.

Az egyes beállítások helyettesítési viselkedéséről további információt a következő témakörben talál INSERT: . A Delta Lake-lehetőségek teljes listáját DataFrameWriter a Delta Lake és az Apache Iceberg című témakörben találja.

A Scala és a Python nyelvben a replaceOn és a replaceUsing nem használható együtt a replaceWhere, a partitionOverwriteMode vagy a overwriteSchema elemmel.

Üres forráslekérdezések esetén sem a REPLACE USING, sem a REPLACE ON nem töröl adatokat, a REPLACE WHERE azonban adatokat törölhet.

Important

Ha véletlenül felülírták az adatokat, a visszaállítással visszavonhatja a módosítást.

REPLACE WHERE

Csak azokat az adatokat írhatja felül szelektíven, amelyek egy tetszőleges kifejezéssel egyezőek REPLACE WHERE.

Important

Ha szeretné kihasználni az inkrementális frissítés előnyeit a(z) REPLACE WHERE futtatásakor, használjon REPLACE WHERE-folyamatokat a Lakeflow-folyamataiban. Lásd: Kötegelt feldolgozás REPLACE WHERE folyamatokkal.

REPLACE WHERE nem kell a táblát az oszlopokra osztani a predikátumban. A predikátum a tábla oszlopaira vonatkozó tetszőleges kifejezés.

Az alábbi példák egy events táblát használnak, amely 2017 januárján belüli és kívüli sorokkal van feltöltve, valamint egy replace_data forrást, amely új, 2017 januári sorokat tartalmaz. A létrehozásához a következőket futtatjuk:

CREATE OR REPLACE TABLE main.default.events (
  event_id INT, start_date DATE, end_date DATE, payload STRING
);
INSERT INTO main.default.events VALUES
  (1, DATE '2017-01-05', DATE '2017-01-10', 'jan-original-a'),
  (2, DATE '2017-01-20', DATE '2017-01-25', 'jan-original-b'),
  (3, DATE '2017-02-05', DATE '2017-02-10', 'feb-unchanged');

CREATE OR REPLACE TABLE main.default.replace_data (
  event_id INT, start_date DATE, end_date DATE, payload STRING
);
INSERT INTO main.default.replace_data VALUES
  (10, DATE '2017-01-08', DATE '2017-01-12', 'jan-replacement-a'),
  (11, DATE '2017-01-15', DATE '2017-01-18', 'jan-replacement-b');

A events januári sorainak atomikus lecseréléséhez a replace_data adataival futtassa az alábbit. A februári sor nem esik a predikátum hatókörébe, ezért változatlan marad:

Python

replace_data = spark.read.table("main.default.replace_data")

(replace_data.write
  .mode("overwrite")
  .option("replaceWhere", "start_date >= '2017-01-01' AND end_date <= '2017-01-31'")
  .saveAsTable("main.default.events")
)

Scala

val replace_data = spark.read.table("main.default.replace_data")

replace_data.write
  .mode("overwrite")
  .option("replaceWhere", "start_date >= '2017-01-01' AND end_date <= '2017-01-31'")
  .saveAsTable("main.default.events")

SQL

INSERT INTO TABLE main.default.events REPLACE WHERE start_date >= '2017-01-01' AND end_date <= '2017-01-31' SELECT * FROM main.default.replace_data

Ez a mintakód kiírja az adatokat replace_data, ellenőrzi, hogy az összes sor megfelel-e a predikátumnak, és overwrite szemantika segítségével hajt végre atomi pótlást. Ha a művelet bármely értéke kívül esik a predikátumon, ez a művelet alapértelmezés szerint hibával meghiúsul.

Klasszikus számítás esetén, ha ezt a viselkedést a predikátumtartományon belüli értékekre overwrite és insert a megadott tartományon kívüli rekordokra szeretné módosítani, távolítsa el a kényszerellenőrzést a következő beállítással spark.databricks.delta.replaceWhere.constraintCheck.enabledfalse:

Python

spark.conf.set("spark.databricks.delta.replaceWhere.constraintCheck.enabled", False)

Scala

spark.conf.set("spark.databricks.delta.replaceWhere.constraintCheck.enabled", false)

SQL

SET spark.databricks.delta.replaceWhere.constraintCheck.enabled=false

Note

REPLACE WHERE boolean_expression bizonyos korlátozásokat elfogad. Lásd INSERT az SQL nyelvi referenciájában.

Üres forráslekérdezések REPLACE WHERE esetén törölheti a táblázat sorait.

Örökölt viselkedés

replaceWhere Az örökölt függvény csak klasszikus számítási környezetben érhető el. Tekintse meg a klasszikus számítás áttekintését.

Ha a régi viselkedést replaceWherehasználja, a lekérdezések felülírják a predikátumnak megfelelő adatokat csak partícióoszlopok felett. A következő parancs atomilag lecseréli a januári hónapot a céltáblában, amelyet a date particionált, a df adatainak felhasználásával.

Python
(df.write
  .mode("overwrite")
  .option("replaceWhere", "birthDate >= '2017-01-01' AND birthDate <= '2017-01-31'")
  .saveAsTable("people10m")
)
Scala
df.write
  .mode("overwrite")
  .option("replaceWhere", "birthDate >= '2017-01-01' AND birthDate <= '2017-01-31'")
  .saveAsTable("people10m")

A korábbi működés használatához állítsa a(z) spark.databricks.delta.replaceWhere.dataColumns.enabled értékét erre: false:

Python
spark.conf.set("spark.databricks.delta.replaceWhere.dataColumns.enabled", False)
Scala
spark.conf.set("spark.databricks.delta.replaceWhere.dataColumns.enabled", false)
SQL
SET spark.databricks.delta.replaceWhere.dataColumns.enabled=false

Dinamikus adatok felülírása

A dinamikus adatok felülírják a megadott kulcsoszlopokkal vagy logikai kifejezéssel egyező adatokat, így az összes többi adat változatlan marad. A particionált táblák, a nem particionált táblák és a folyékony fürtözésű táblák mind támogatottak.

A dinamikus partíció felülírása a dinamikus adatok felülírásának egy részhalmaza. A dinamikus partíció felülírja az összes meglévő adatot minden olyan partíción, amelynek írása új adatokat véglegesít, és az összes többi partíciót változatlanul hagyja. Csak a particionált táblák támogatottak.

REPLACE USING

A Databricks Runtime 16.3-ban és újabb verziókban támogatott SQL. Python és Scala támogatott a Databricks Runtime 18.2-ben és újabb verziókban. A Databricks Runtime 16.3 és 17.1 közötti viselkedésbeli különbségekért tekintse meg az örökölt viselkedést.

REPLACE USING lehetővé teszi a databricks SQL-tárolókon, kiszolgáló nélküli számításon és klasszikus számításon működő, számításfüggetlen, atomi felülírásos viselkedést. REPLACE USING nincs szükség a Spark-munkamenet konfigurálásának beállítására.

REPLACE USING akkor cseréli le a sorokat, ha a megadott oszlopok egyenlőség esetén egyenlőek. Az összes többi adat változatlan marad.

Az alábbi példák visszaállítják a events táblát, és egy source_data forrást használnak. A létrehozásához a következőket futtatjuk:

CREATE OR REPLACE TABLE main.default.events (
  event_id INT, start_date DATE, end_date DATE, payload STRING
);
INSERT INTO main.default.events VALUES
  (1, DATE '2017-01-05', DATE '2017-01-10', 'original-a'),
  (2, DATE '2017-01-20', DATE '2017-01-25', 'original-b'),
  (3, DATE '2017-02-05', DATE '2017-02-10', 'original-c');

CREATE OR REPLACE TABLE main.default.source_data (
  event_id INT, start_date DATE, end_date DATE, payload STRING
);
INSERT INTO main.default.source_data VALUES
  (1, DATE '2017-01-05', DATE '2017-01-11', 'replacement-a'),
  (4, DATE '2017-03-01', DATE '2017-03-02', 'new-d');

REPLACE USING (event_id, start_date) kicseréli azt a events sort, amelynek event_id és start_date egyezik source_data (1. sor), új kulcsot helyez be a forrássorra (4. sor), és a többi sort változatlan marad. Dinamikus adat felülírása a következővel REPLACE USING:

Python

sourceDataDF = spark.read.table("main.default.source_data")

(sourceDataDF.write
  .mode("overwrite")
  .option("replaceUsing", "event_id, start_date")
  .saveAsTable("main.default.events")
)

Scala

val sourceDataDF = spark.read.table("main.default.source_data")

sourceDataDF.write
  .mode("overwrite")
  .option("replaceUsing", "event_id, start_date")
  .saveAsTable("main.default.events")

SQL

INSERT INTO TABLE main.default.events
  REPLACE USING (event_id, start_date)
  SELECT * FROM main.default.source_data

Üres forrás lekérdezések REPLACE USING esetén nem töröl táblázatsorokat.

Összetett vagy NULL-biztonságos egyeztetési feltételek esetén használja REPLACE ON inkább. Lásd a(z) REPLACE ON.

Lásd INSERT az SQL nyelvi referenciájában.

Örökölt viselkedés

A Databricks Runtime 16.3-17.1-ben örökölt viselkedést használ, REPLACE USING és csak a dinamikus partíciót írja felül, míg a Databricks Runtime 17.2 és újabb verziók lehetővé teszik a dinamikus adat felülírást.

Tartsa szem előtt az alábbi korlátozásokat és viselkedéseket az örökölt viselkedésre vonatkozóan REPLACE USING :

  • A USING záradékban meg kell adnia a tábla partícióoszlopainak teljes készletét.
  • Mindig ellenőrizze, hogy az írott adatok csak a várt partíciókat érintik-e. A rossz partíció egyetlen sora véletlenül felülírhatja a teljes partíciót.

REPLACE ON

A Databricks Runtime 17.1-ben és újabb verziókban támogatott SQL. Python és Scala támogatott a Databricks Runtime 18.2-ben és újabb verziókban.

REPLACE ON lecseréli azokat a sorokat, amelyek megfelelnek egy felhasználó által definiált feltételnek, ellentétben a REPLACE USING-val, amely akkor cseréli le a sorokat, ha a megadott oszlopok az egyenlőségi összehasonlítás szerint megegyeznek. Akkor használja REPLACE ON , ha olyan egyező logikára van szüksége, amely REPLACE USING nem támogatja, például az értékek egyenlőként való kezelését NULL .

Ha szeretné, a targetAlias beállítással megadhat egy aliast a céltábla és az .as().alias() API-k számára a forrásadatok aliasának megadásához.

Az SQL szintaxisát lásd: INSERT.

A következő példák visszaállítják a events táblát és a source_data forrást, és mindkettőbe beszúrnak egy-egy sort, amely NULLstart_date értéket tartalmaz, hogy bemutassák a NULL-biztos egyezést. A létrehozásához a következőket futtatjuk:

CREATE OR REPLACE TABLE main.default.events (
  event_id INT, start_date DATE, end_date DATE, payload STRING
);
INSERT INTO main.default.events VALUES
  (1, DATE '2017-01-05', DATE '2017-01-10', 'original-a'),
  (2, NULL, DATE '2017-01-25', 'original-null-key'),
  (3, DATE '2017-02-05', DATE '2017-02-10', 'original-c');

CREATE OR REPLACE TABLE main.default.source_data (
  event_id INT, start_date DATE, end_date DATE, payload STRING
);
INSERT INTO main.default.source_data VALUES
  (1, DATE '2017-01-05', DATE '2017-01-11', 'replacement-a'),
  (2, NULL, DATE '2017-01-26', 'replacement-null-key');

Mivel a <=> két NULL értéket egyenlőnek tekint, a NULLstart_date értéket tartalmazó sor (2. sor) egyezőnek minősül, és az 1. sorral együtt lecserélődik. A 3. sorban nincs megfelelő forrássor, és változatlan marad:

Python

sourceDataDF = spark.read.table("main.default.source_data")

(sourceDataDF.alias("s")
  .write
  .mode("overwrite")
  .option("targetAlias", "t")
  .option("replaceOn", "s.event_id <=> t.event_id AND s.start_date <=> t.start_date")
  .saveAsTable("main.default.events")
)

Scala

val sourceDataDF = spark.read.table("main.default.source_data")

sourceDataDF.as("s")
  .write
  .mode("overwrite")
  .option("targetAlias", "t")
  .option("replaceOn", "s.event_id <=> t.event_id AND s.start_date <=> t.start_date")
  .saveAsTable("main.default.events")

SQL

INSERT INTO TABLE main.default.events AS t
  REPLACE ON (s.event_id <=> t.event_id AND s.start_date <=> t.start_date)
  (SELECT * FROM main.default.source_data) AS s

Üres forrás lekérdezések REPLACE ON esetén nem töröl táblázatsorokat.

Dinamikus partíció felülírása partitionOverwriteMode használatával (örökölt)

Important

Ez a funkció nyilvános előzetes verzióban van.

A Databricks Runtime 11.3 LTS és újabb verziója támogatja a dinamikus partíció felülírását a particionált táblákhoz felülírási módban: SQL-ben INSERT OVERWRITE vagy DataFrame-írással df.write.mode("overwrite"). Ez a felülírástípus csak a klasszikus számításhoz érhető el, a Databricks SQL Warehouse-hoz és a kiszolgáló nélküli számításhoz nem.

Figyelmeztetés

Ha lehetséges, használja a INSERT REPLACE USING-t a partíció felülírása helyett, amely magában foglalja a INSERT OVERWRITE PARTITION és spark.sql.sources.partitionOverwriteMode=dynamic elemeket. A partíciók módosításakor történő felülírás során elavult adatok kerülhetnek felhasználásra.

A dinamikus partíció felülírási módjának használatához állítsa a Spark-munkamenet konfigurációját spark.sql.sources.partitionOverwriteMode a következőre dynamic: . Alternatívaként beállíthatja a DataFrameWriter lehetőséget partitionOverwriteModedynamic. Ha van ilyen, a lekérdezésspecifikus beállítás felülírja a munkamenet-konfigurációban definiált módot. A spark.sql.sources.partitionOverwriteMode alapértelmezett értéke static.

Az alábbi példa partitionOverwriteModehasznál:

SQL

SET spark.sql.sources.partitionOverwriteMode=dynamic;
INSERT OVERWRITE TABLE default.people10m SELECT * FROM morePeople;

Python

(df.write
  .mode("overwrite")
  .option("partitionOverwriteMode", "dynamic")
  .saveAsTable("default.people10m")
)

Scala

df.write
  .mode("overwrite")
  .option("partitionOverwriteMode", "dynamic")
  .saveAsTable("default.people10m")

Tartsa szem előtt a következő korlátozásokat és működési módokat partitionOverwriteMode:

  • A overwriteSchema nem állítható be true.
  • Nem adhatja meg mind a kettőt, partitionOverwriteMode és replaceWhere, egyazon DataFrameWriter műveletben.
  • Ha egy feltételt replaceWhere egy DataFrameWriter beállítással ad meg, a Delta Lake ezt a feltételt alkalmazza annak szabályozására, hogy mely adatokat írja felül a rendszer. Ez a beállítás elsőbbséget élvez a partitionOverwriteMode munkamenetszintű konfigurációval szemben.
  • Mindig ellenőrizze, hogy az írott adatok csak a várt partíciókat érintik-e. A rossz partíció egyetlen sora véletlenül felülírhatja a teljes partíciót.