Állapotalapú lekérdezések aszinkron állapot-ellenőrzőpontja

Megjegyzés

A Databricks Runtime 10.4 LTS-ben és újabb verziókban érhető el.

Az aszinkron állapotellenőrzés biztosítja az épp egyszeri garantálást a streamelési lekérdezésekhez, de csökkentheti a teljes késleltetést egyes állapotfrissítéseknél szűk keresztmetszettel rendelkező strukturált állapotalapú streamelési számítási feladatok esetében. Ez a következő mikroköteg feldolgozásának megkezdésével érhető el, amint az előző mikroköteg számítása befejeződött anélkül, hogy az állapot-ellenőrzőpontozás befejezésére vártak volna. Az alábbi táblázat összehasonlítja a szinkron és az aszinkron ellenőrzőpontokra vonatkozó kompromisszumokat:

Jellemző Szinkron ellenőrzőpont-készítés Aszinkron ellenőrzőpont-készítés
Késleltetés Nagyobb késés az egyes mikrokötegekhez. Csökken a késés, mivel a mikrocsomagok áthatolhatnak egymásba.
Újraindítás Gyors helyreállítás, mert csak az utolsó köteget kell újra lefuttatni. Nagyobb újraindítási késleltetés, mivel több mikroköteget is újra kell futtatni.

Az alábbiakban találhatók a streamelési feladatok azon jellemzői, amelyek előnyösek lehetnek az aszinkron állapot-ellenőrzőpontozás során:

  • A feladat egy vagy több állapotalapú művelettel rendelkezik (például aggregáció, flatMapGroupsWithState, mapGroupsWithState, stream-stream csatlakozások)
  • Az állapotellenőrzési pontok késése a kötegelt végrehajtás általános késésének egyik fő közreműködője. Ezek az információk a StreamingQueryProgress eseményekben találhatók. Ezek az események a Spark-illesztőprogram log4j-naplóiban is megtalálhatók. Íme egy példa a streamelési lekérdezések előrehaladására, valamint arra, hogyan állapíthatja meg, hogy az állapot-ellenőrzőpont milyen hatással van a kötegelt végrehajtás teljes késésére.
    • {
         "id" : "2e3495a2-de2c-4a6a-9a8e-f6d4c4796f19",
         "runId" : "e36e9d7e-d2b1-4a43-b0b3-e875e767e1fe",
         "...",
         "batchId" : 0,
         "durationMs" : {
           "...",
           "triggerExecution" : 547730,
           "..."
         },
         "stateOperators" : [ {
           "...",
           "commitTimeMs" : 3186626,
           "numShufflePartitions" : 64,
           "..."
         }]
      }
      
    • Állapotellenőrzési pont késésének elemzése a fenti lekérdezési folyamat eseményéről

      • A sorozat időtartama (durationMs.triggerDuration) körülbelül 547 másodperc.
      • Az állapottároló véglegesítési késése (stateOperations[0].commitTimeMs) körülbelül 3186 másodperc. A véglegesítés késése összesítve van az állapottárolót tartalmazó tevékenységek között. Ebben az esetben 64 ilyen tevékenység van (stateOperators[0].numShufflePartitions).
      • Az állapotkezelőt tartalmazó feladatok ellenőrzőpontig átlagosan 50 másodpercet (3,186/64) igényeltek. Ez egy további késés, amely hozzájárul a tétel időtartamához. Feltételezve, hogy mind a 64 feladat egyidejűleg fut, az ellenőrzőpont művelet a köteg időtartamának körülbelül 9%-át (50 másodperc / 547 másodperc) tette ki. Ez a százalék még magasabb lesz, ha az egyidejű tevékenységek maximális értéke kevesebb, mint 64.

Aszinkron állapotellenőrzés engedélyezése

Az aszinkron állapotellenőrzéshez a RocksDB-alapú állapottárolót kell használnia. Állítsa be a következő konfigurációkat:


spark.conf.set(
  "spark.databricks.streaming.statefulOperator.asyncCheckpoint.enabled",
  "true"
)

spark.conf.set(
  "spark.sql.streaming.stateStore.providerClass",
  "com.databricks.sql.streaming.state.RocksDBStateStoreProvider"
)

Az aszinkron ellenőrzőpontokra vonatkozó korlátozások és követelmények

Megjegyzés

A számítási kapacitás automatikus skálázása korlátozott a strukturált streamelési munkaterhelések fürtméretének csökkentésében. A Databricks a Lakeflowban elérhető, továbbfejlesztett automatikus skálázással rendelkező Spark deklaratív adatfeldolgozási folyamatait ajánlja streaming számítási feladatokhoz. Lásd A Lakeflow-adatcsatorna fürt kihasználásának optimalizálása automatikus skálázással.

  • Az aszinkron ellenőrzőpontokban egy vagy több tárolóban előforduló hibák meghiúsulnak a lekérdezésben. Szinkron ellenőrzőpont-módban a rendszer a feladat részeként hajtja végre az ellenőrzőpontot, és a Spark többször is újrapróbálkozza a feladatot, mielőtt a lekérdezés meghiúsul. Ez a mechanizmus nem rendelkezik aszinkron állapot-ellenőrzőpontokkal. A Databricks azt javasolja, hogy folyamatos feladatokat használva automatikus újrapróbálkozások történjenek feladathiba esetén. Tekintse meg a feladatok folyamatos futtatása című témakört.
  • Az aszinkron ellenőrzőpontok akkor működnek a legjobban, ha az állapottároló helye nem változik a mikroköteg-végrehajtások között. Előfordulhat, hogy a fürt átméretezése az aszinkron állapotellenőrzéssel együtt nem működik megfelelően azon okból kifolyólag, hogy az állapottárolók példányai újraelosztódhatnak, miközben csomópontokat adnak hozzá vagy törölnek a fürt átméretezési eseménye során.
  • Az aszinkron állapot-ellenőrzőpontozás csak a RocksDB állapottároló-szolgáltató implementációjában támogatott. A memóriabeli állapottár alapértelmezett implementációja nem támogatja azt.