Streamelés kiszolgáló nélküli számításon

Ez a lap bemutatja, hogyan választhatja ki a kiszolgáló nélküli streamelési számítási feladatok megfelelő konfigurációját Azure Databricks, beleértve a folyamatos folyamatokat, a növekményes betöltési folyamatokat és a felügyelt összekötőket. A megfelelő konfiguráció kiválasztása a stream forrás-, alakzat- és késési igényeitől függ.

Mi számít streamelési számítási feladatnak?

Az adatfolyam-feldolgozási feladat egy forrásból (például felhőalapú objektumtárolóból, üzenetbuszból vagy változási hírcsatornából) olvas korlátlan mennyiségű adatot, és növekményesen ír egy célba. Azure Databricks a streamelési számítási feladatok két mintáját támogatja:

  • Folyamatos: Egy folyamat, amely leállítás nélkül fut, és új adatokat dolgoz fel a beérkezéskor. A késés mérése másodpercben történik.
  • Növekményes (más néven aktivált): Ütemezésen vagy eseményindítón futó folyamat, amely az utolsó futtatás óta érkezett összes adatot feldolgozza, és leáll. A késés mérése percekben történik.

Egyes munkaterhelések adatfolyam-feldolgozási folyamatoknak tűnnek, de technikai értelemben nem folyamatok. Ilyenek például az események figyelésére nyitott websocketet tartalmazó szolgáltatás, egy felhasználónként állandó kapcsolatot fenntartó csevegőalkalmazás vagy a bejövő HTTP-kéréseket kezelő webhook-fogadó. Ezek alkalmazások, nem streamelési folyamatok. A megfelelő szerver nélküli lehetőséggel kapcsolatban lásd a Nem streamelési folyamatnak minősülő számítási feladatok című részt.

Válassza ki a megfelelő streamkonfigurációt

Ez a táblázat megfelelteti az eseteket azoknak a kiszolgáló nélküli konfigurációknak, amelyek a legjobban illeszkednek hozzájuk. Az ezen az oldalon található szakaszok részletesebben ismertetik ezeket a javaslatokat.

Felhasználási eset Ajánlott konfiguráció Miért
Folyamatos, alacsony késleltetésű streaming ETL vagy transzformációk Lakeflow-folyamatok folyamatos módban A folyamatos mód az állandóan aktív adatfolyamokhoz készült. A folyamatos csővezetékes feldolgozás lehetővé teszi a mikrokötegek egyidejű futtatását, javítva az áteresztőképességet és csökkentve a késleltetést. A felügyelt állapot automatikusan tartja a helyreállítást.
Növekményes betöltés a felhőbeli tárolóból Használja az automatikus betöltőt a Lakeflow-folyamatokban (alacsony késés esetén) vagy kiszolgáló nélküli feladatbanTrigger.AvailableNow() (ha az alacsonyabb késés elfogadható). Az Automatikus betöltő hatékonyan követi nyomon az új fájlokat. Trigger.AvailableNow() feldolgozza a hátralékot, majd kilép, amely megfelel egy ütemezett vagy igény szerinti ütemezésnek.
Felügyelt betöltés SaaS-forrásokból vagy adatbázis CDC-ből Standard összekötők a Lakeflow Connectben Teljes mértékben felügyelt összekötők kiszolgáló nélküli betöltési folyamatokkal. A támogatott forrásokhoz nincs szükség kódra.
SQL streamelése Delta-táblákon keresztül adatfolyam táblák SQL-natív inkrementális feldolgozás hozzáfűzésorientált forrásokhoz, felügyelt folyamatokkal és frissítéssel.
Időszakos mikroköteg-feldolgozás jegyzetfüzetben vagy feladatban Szerver nélküli feladat ezzel: Trigger.AvailableNow() Költséghatékony, ha a perc szintű frissesség elegendő. A szerver nélküli számítás gyorsan elindul, és leáll, amikor a kötegelt feladat befejeződik.

Folyamatos streamelés

A kiszolgáló nélküli számításon futó folyamatos streameléshez használja a Lakeflow-folyamatokat folyamatos módban. A feldolgozási folyamat folyamatosan fut, feldolgozza a beérkező rekordokat, és automatikusan helyreáll a hibák után.

Folyamatos stream konfigurálása:

Tip

A streamcsövezés alapértelmezés szerint engedélyezve van a kiszolgáló nélküli Lakeflow-folyamatokban. A mikrokötegek egymással párhuzamosan futnak, nem egymás után, ami javítja az adatbetöltés-igényes adatfolyamok áteresztőképességét.

Az időalapú strukturált streamelési eseményindítók, például Trigger.ProcessingTime(interval) és Trigger.Continuous(interval), nem érhetők el kiszolgáló nélküli jegyzetfüzetekben vagy feladatokban. A Lakeflow-folyamatokat folyamatos módban használja a mindig aktív minta esetén. Lásd a streamelési korlátozásokat. Trigger.Once() támogatott, de elavult — migrálja a meglévő lekérdezéseket ide: Trigger.AvailableNow()

Növekményes és aktivált streamelés

Növekményes streameléshez futtassa a strukturált streamelést kiszolgáló nélküli feladatban Trigger.AvailableNow() . Minden futtatás feldolgozza az utolsó ellenőrzőpont óta érkezett összes adatot, majd kilép.

Kiszolgáló nélküli feladat konfigurálása növekményes streameléssel:

Az alábbi példa beolvassa az új fájlokat a felhőbeli tárolóból (source_path) az Automatikus betöltővel, feldolgozza a futtatáskor rendelkezésre álló összes adatot, és egy Delta-táblába ír:

(spark.readStream
   .format("cloudFiles")
   .option("cloudFiles.format", "json")
   .option("cloudFiles.maxFilesPerTrigger", 1000)
   .load(source_path)
   .writeStream
   .trigger(availableNow=True)
   .option("checkpointLocation", checkpoint_path)
   .toTable("catalog.schema.target_table"))

Az ütemezett Trigger.AvailableNow() feladat a kiszolgáló nélküli számítás legköltséghatékonyabb streamelési mintája, ha a percszintű késés elfogadható. A számítási erőforrás másodpercek alatt elindul, lefuttatja a kötegelt feladatot, majd leáll.

Felügyelt betöltés

Ha a forrás SaaS-alkalmazás vagy operatív adatbázis, strukturált streamelési kód írása helyett használja a Lakeflow Connectet. A Lakeflow Connect kiszolgáló nélküli betöltési folyamatokat futtat olyan összekötőkhöz, mint a Salesforce, a Workday, a SQL Server CDC és a PostgreSQL CDC. Lásd: Felügyelt csatlakozók a Lakeflow Connectben.

Ez az elérési út a megfelelő válasz, ha:

  • A forráshoz létezik összekötő.
  • Egyéni kód helyett felügyelt folyamatot szeretne.
  • Szüksége van a sémaváltozások kezelésére, az adatok származáskövetésére és a monitorozásra, alapból.

SQL által felügyelt növekményes adatfeldolgozás

Az SQL-t előnyben részesítő csapatok számára használjon streaming táblákat natív SQL-es streaming számítási feladatokhoz. A streamelési táblákat a Lakeflow-folyamatokban vagy önálló streamelő táblákként is definiálhatja.

Az SQL-utasítással CREATE OR REFRESH STREAMING TABLE létrehozott különálló streamtáblák esetében a kezdeti adatfrissítés és -sokaság azonnal megkezdődik. A rendszer automatikusan létrehoz és felügyel egy dedikált kiszolgáló nélküli folyamatot az egyes streamelési táblákhoz.

Ha kötegelt szemantikai lekérdezési eredményekre van szüksége felügyelt frissítéssel, használjon materializált nézeteket. Lásd Materializált nézetek.

Nem streamingadat-folyamatnak minősülő munkaterhelések

Az a számítási feladat, amely hosszú élettartamú kapcsolatot igényel, porton figyel, vagy a bejövő HTTP-kérésekre válaszol, nem streamelési folyamat; ez egy alkalmazás. Ne futtassa ezeket a számítási feladatokat kiszolgáló nélküli feladaton. A Databricks megfelelő beállításai a következők:

  • Hosszú ideig futó szolgáltatások, amelyekhez állandó kapcsolatra vagy HTTP-végpontra van szükség: A szolgáltatás létrehozása a Databricks Apps használatával. A Databricks Apps kiszolgáló nélküli platform az egyéni alkalmazások Azure Databricks való üzemeltetéséhez, beleértve a FastAPI- és Flask-, Streamlit-, Dash-, Gradio-, Node.js- és Shiny-alkalmazásokat. Lásd : Databricks Apps.
  • Bejövő webhookok vagy eseményfigyelők: HTTP-végpont közzététele a Databricks Appsben, vagy leállítja a webhookot egy külső szolgáltatásban, és eseményeket ír a felhőbeli tárolóba vagy egy üzenetbuszba, majd egy kiszolgáló nélküli streamelési folyamattal felveszi őket.
  • Egyéni token- vagy hitelesítőadat-csere: Használjon szolgáltatásfiókokat az OAuth-tal, vagy hívja meg a Databricks REST API-kat egy alkalmazásból. Az adatfolyam-feldolgozási folyamatok nem tárolnak felhasználói munkameneteket vagy egyéni tokenállapotot.

Ha azt értékeli, hogy a munkaterhelése illeszkedik-e egy streamfeldolgozási folyamathoz, tegye fel a következő kérdést:

  • Olvas a munkaterhelés egy nem korlátos adatforrásból, és ír egy nyelőbe? Ha igen, akkor ez egy adatfolyam-feldolgozási folyamat.
  • Szükséges-e, hogy a munkaterhelés fenn tudjon tartani egy kapcsolatot egy ügyféllel? Ha igen, ez egy alkalmazás; a Databricks Apps használata.

Limitations

A kiszolgáló nélküli számítás a következő streamelési korlátozásokat szabja meg. Egyik sem akadályozza meg a fenti számítási feladatokat a megfelelő termékkel párosítva.

  • Az időalapú strukturált streamelési eseményindítók (Trigger.ProcessingTime(interval) és ) kiszolgáló nélküli jegyzetfüzetekben és Trigger.Continuous(interval)feladatokban nem támogatottak. Használja a Lakeflow-folyamatokat folyamatos módban a folyamatos adatfolyamokhoz, vagy Trigger.AvailableNow() az eseményindítású futtatásokhoz. Lásd a streamelési korlátozásokat.
  • Az explicit trigger nélküli streaming lekérdezések a következő hibával meghiúsulnak: INFINITE_STREAMING_TRIGGER_NOT_SUPPORTED Az Apache Spark alapértelmezés szerint Trigger.ProcessingTime("0 seconds")a kiszolgáló nélküli számításban nem támogatott. Mindig állítsa be Trigger.AvailableNow() az összes streamelési lekérdezést, vagy használja a Lakeflow-folyamatokat folyamatos módban.
  • A standard hozzáférési módban történő streamelésre vonatkozó összes korlátozás a kiszolgáló nélküli számításra is vonatkozik. Lásd a streamelési korlátozásokat.

Következő lépések