Samanaikaisuuden hallinta Delta-taulukoille

Kun useat Fabric-muistikirjat, putket tai Spark-työt kirjoittavat samaan Delta-taulukkoon samanaikaisesti, Delta Lake käyttää optimistista samanaikaisuuden hallintaa (OCC) pitääkseen taulukon johdonmukaisena. Jokainen transaktio lukee snapshotin, kirjoittaa uusia tiedostoja ja varmistaa, ettei välissä ole tapahtunut ristiriitaisia committeja. Jos ristiriita havaitaan, transaktio epäonnistuu poikkeuksella sen sijaan, että tiedot korruptoituisi.

Tämä artikkeli käsittelee käytännön malleja samanaikaisten kirjoitusten hallintaan Fabric-ohjelmassa. OCC-protokollan täydellisen määrittelyn löydät katsosta Delta Lake conparallel control (avoimen lähdekoodin dokumentaatio).

Eristystasot

Kaikki Delta-taulukot käyttävät Serializable Isolation -tasoa. Sarjoitettava on tiukin taso ja ainoa tuettu. Se varmistaa, että samanaikaisten tapahtumien tulos on identtinen jonkin peräkkäisen suoritusmääräyksen kanssa.

Delta Lake käyttää myös sisäistä SnapshotIsolation-tasoa operaatioissa, jotka eivät muuta loogista dataa (kuten OPTIMIZE). SnapshotIsolation ohittaa samanaikaisen lisäyksen tarkistuksen ja mahdollistaa tiivistämisen ilman ristiriitaa samanaikaisten lisäysten kanssa. Et voi määrittää SnapshotIsolationia suoraan—Delta Lake soveltaa sitä automaattisesti tarvittaessa.

Eristyksessä Serializable samanaikainen sokea liite (INSERT INTO) voi olla ristiriidassa saman MERGE osion lukeman tai UPDATE -ryhmän kanssa.

Mitkä operaatiot ovat ristiriidassa

Kaikki samanaikaiset kirjoitukset eivät ole ristiriidassa. Keskeinen tekijä on, kosketetaanko kaksi operaatiota samoja taustatiedostoja.

Samanaikainen pari Ristiriita? Miksi
Kaksi INSERT (liite) operaatiota Ei Jokainen lisää uusia tiedostoja lukematta olemassa olevia (blind append).
INSERT + OPTIMIZE Ei OPTIMIZE commit, koska SnapshotIsolation se ei muuta loogista dataa, joten se ohittaa samanaikaisen liitteen tarkistuksen kokonaan. Liitteet lisäävät uusia tiedostoja, jotka eivät mene päällekkäin tiivistettyjen tiedostojen kanssa.
Kaksi UPDATE, DELETE, eli MERGE operaatiota Kyllä, jos he lukevat tai muokkaavat päällekkäisiä tiedostoja Jokainen kirjoittaa tiedostoja uudelleen, joten toisen kirjoittajan tilannekuva on vanhentunut.
OPTIMIZE + UPDATE/DELETE/MERGE Kyllä, jos he koskettavat samoja tiedostoja OPTIMIZEPoistaa ja lisää tiedostoja uudelleen ().dataChange=false Jos datan muokkausoperaatio lukee samoja tiedostoja, a ConcurrentDeleteReadException nousee a.
Kaksi OPTIMIZE juoksua Kyllä, jos he valitsevat samat tiedostot Molemmat yrittävät poistaa ja kirjoittaa uudelleen saman tiedostojoukon, mikä laukaisee .ConcurrentDeleteDeleteException
INSERT + MERGE/UPDATE/DELETE Kyllä, jos datan muokkausoperaatio lukee saman osion Alaisuudessa Serializablesokea liite voi olla ristiriidassa samanaikaisten datamuutosten kanssa, jos operaatio lukee osion, johon liite kirjoitti.

Vinkki

Vain liitännäiset putket (INSERT INTO, df.write.mode("append")) ovat yksinkertaisin tapa välttää konfliktit kokonaan. Jos työkuormasi voidaan lisätä ensin ja sovittaa vasta myöhemmin, poistat kirjoitus-kirjoitus-kiistan.

Eristäkää kirjoittajat osioinnilla

Yleisin tapa ajaa samanaikaista DML:ää samaa taulukkoa vastaan ilman ristiriitoja on jakaa taulukko sarakkeen mukaan, joka erottaa kirjoittajasi, ja sisällyttää kyseinen sarake jokaiseen operaatioehtoon. Kun kumpikin kirjoittaja kohdistaa eri osion, operaatiot koskettavat erillisiä tiedostojoukkoja eivätkä ole ristiriidassa.

Tyypillinen tilanne: useat putket käsittelevät dataa eri liiketoimintayksikölle tai vuokralaiselle. Jaa tuolle ulottuvuudelle ja kiinnitä jokainen putki MERGE omaan osioonsa.

-- Each pipeline targets its own partition, so concurrent runs don't conflict
MERGE INTO events AS target
USING staged AS source
ON target.event_id = source.event_id
    AND target.business_unit = 'EMEA'
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *

Important

Osion sarakkeen täytyy esiintyä yhdistämisehdossa – ei vain lähdedatassa. Ilman sitä Delta Lake ei voi validointihetkellä todeta, että molemmat operaatiot koskettivat erillisiä tiedostokokonaisuuksia, ja konfliktitarkistelija käsittelee operaation täyden taulukon lukemisena.

Lisätietoja osiointistrategioista löytyy kohdasta Partitioning for Delta Tables.

Sisäänrakennettu commit-uudelleenyritys

Delta Lake yrittää automaattisesti uudelleen commitin, kun se havaitsee toisen tapahtuman sitoutuneen ensin. Jokaisella uusintayrityksellä se lukee voittavan commitin, suorittaa konfliktin tarkistuksen ja—jos loogista ristiriitaa ei ole—yrittää sitoutumisen seuraavassa saatavilla olevassa versiossa. Tämä prosessi toistuu läpinäkyvästi ilman mitään toimenpiteitä koodistasi.

Looginen ristiriita (esimerkiksi kaksi operaatiota, jotka kirjoittavat samaa tiedostoa uudelleen) ei voi ratkaista automaattisesti. Uudelleenyritys nostaa esiin yhden poikkeuksista, jotka on lueteltu Common conflict exceptions -artikkelissa. Monet tilapäiset versiotörmäykset – kuten kaksi sokeaa lisäosaa kilpailemassa samasta versiopaikasta – ratkeavat automaattisesti eivätkä koskaan ilmesty sovellukseesi.

Yleiset konfliktipoikkeukset

Kun konflikti havaitaan, Delta Lake nostaa erityisen poikkeuksen. Ymmärtämällä, minkä poikkeuksen näet, auttaa tunnistamaan juurisyyn.

Poikkeus Mitä tapahtui
ConcurrentAppendException Toinen kirjoittaja lisäsi tiedostoja osioon (tai tiedostojoukkoon), jota operaatiosi luki. Yleistä, kun A MERGE toimii osiota vastaan, joka vastaanottaa myös inserttejä toisesta putkesta. Serializable Eristettynä jopa sokeat liitteet (tavalliset INSERT operaatiot) voivat laukaista tämän poikkeuksen.
ConcurrentDeleteReadException Toinen kirjoittaja poisti tai kirjoitti uudelleen tiedoston, jonka operaatiosi oli lukenut. Tyypillisesti, kun OPTIMIZE tiedostot tiivistetään samanaikaisesti UPDATE tai MERGE myös lukivat, tai kun kaksi datan muokkausoperaatiota menee päällekkäin samoilla riveillä.
ConcurrentDeleteDeleteException Molemmat operaatiot yrittivät poistaa tai kirjoittaa saman tiedoston uudelleen. Usein syynä on päällekkäiset OPTIMIZE ajot tai kaksi putkea, jotka kirjoittavat saman osion uudelleen samanaikaisesti.
ConcurrentWriteException Yleinen ristiriita syntyi, kun toinen transaktio sitoutui samaan taulukkoversioon ennen konfliktinratkaisua – esimerkiksi tiedostojärjestelmän hallittujen committien päivityksen yhteydessä.
MetadataChangedException Taulukkoskeema tai ominaisuudet muuttuivat kesken transaktion—esimerkiksi samanaikainen ALTER TABLE tai skeeman evoluutiokirjoitus.
ConcurrentTransactionException Kaksi Structured Streaming -kyselyä, joilla on sama tarkistuspisteen sijainti, kirjoitettiin taulukkoon samanaikaisesti. Poista suoratoistotyösi tai käytä erillisiä tarkistuspistepolkuja.
ProtocolChangedException Samanaikainen transaktio päivitti tai alensi taulukkoprotokollaa, kun nykyinen transaktio yritti myös muuttaa protokollaa. Voi myös tapahtua, kun taulukon ominaisuus poistetaan samanaikaisesti.

Yleisiä strategioita kirjoitusristiriitojen välttämiseksi

Ota automaattinen tiivistys käyttöön

Automaattinen tiivistys toimii synkronisesti osana kirjoitusoperaatioita. Synkroninen tiivistäminen estää erikseen aikataulutettujen tiivistymistöiden päällekkäisyyttä datan muokkausoperaatioiden kanssa ja siten samanaikaisten kirjoittajien poikkeusten syntymisen.

Aikatauluta ylläpito kirjoitusikkunan ulkopuolelle

OPTIMIZE ja VACUUM voi olla ristiriidassa samanaikaisten tietojen muokkaustoimintojen kanssa. Vuonna Fabric aikatauluta muistikirjan työt tai putkilinjan toiminnot taulukon tiivistämiselle ja VACUUM matalan aktiivisuuden ikkunoiden aikana—esimerkiksi sen jälkeen, kun syöttö on päättynyt eikä sen aikana.

Käytä append + merge -kuvioita

Korkean samanaikaisuuden vastaanottoon siirrä raakadata, jossa on vain liitännäisiä kirjoituksia, staging tableen (ei ristiriitoja mahdollisia), ja suorita sitten yksi MERGE työ, joka sovittaa sen kohdetauluun. Kuvio sarjallistaa konfliktialttiin toiminnon samalla kun niiden vastaanotto pysyy täysin rinnakkaisena.

Lisää uudelleenyrityslogiikka loogisille ristiriidoille

Sisäänrakennettu commit-uudelleenyritys käsittelee transienttiversioiden törmäykset automaattisesti, mutta loogiset ristiriidat – joissa kaksi operaatiota todella menevät päällekkäin – aiheuttavat poikkeuksen. Koska Delta Lake ei koskaan tuota osittaisia kirjoituksia, epäonnistunut transaktio on turvallista yrittää uudelleen sovellustasolla. Putkistoissa, joissa odotetaan satunnaisia loogisia ristiriitoja, kääri kirjoitus retry-logiikkaan:

from delta.exceptions import ConcurrentAppendException
import time

# Retry with backoff on transient concurrent write conflicts
max_retries = 3
for attempt in range(max_retries):
    try:
        spark.sql("MERGE INTO target USING source ON ...")
        break
    except ConcurrentAppendException:
        if attempt < max_retries - 1:
            time.sleep(2 ** attempt)
        else:
            raise

Valitse oikea asettelustrategia

Nesteen klusterointi ja osiointi ratkaisevat erilaisia ongelmia. Liquid clustering optimoi tiedostoasettelun lukusuorituskyvyn kannalta. Jakaminen luo fyysisiä rajoja, jotka estävät samanaikaiset kirjoittajakonfliktit. Jos työkuormasi tarvitsee molempia, jaa se writer-isolation-sarakkeen mukaan ja käytä Z-järjestystä jokaisen osion sisällä lukusuorituskyvyn parantamiseksi.