Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
Egy folyamatot egy adatfeldolgozási munkafolyamat részeként futtathat a Lakeflow Jobs, az Apache Airflow vagy az Azure Data Factory használatával.
A folyamatok automatikusan feloldják az adathalmazok közötti függőségeket, így önállóan kezeli az egyszerű, folyamaton belüli vezénylést. Az olyan vezénylési igényekhez, amelyekre egy folyamat nincs felkészítve – például feltételes végrehajtás, a feladateredmények alapján történő elágazás, újrapróbálkozások vagy egy folyamat más típusú munkákkal való összehangolása esetén –, a logika folyamatba építése helyett inkább dedikált munkafolyamat-vezénylőt használjon.
Készítse elő a folyamatát az orkesztráláshoz
A vezénylés akkor működik a legjobban, ha minden folyamat egy különálló munkaegységet fed le, amelyet önállóan szeretne ütemezni, ellenőrizni vagy futtatni. Ezeket a határokat úgy tervezheti meg, hogy a munkafolyamatok különálló tevékenységekként koordinálhassák őket, beleértve a felső és az alsóbb rétegbeli tevékenységek közötti megfelelő vezérlési folyamatot.
Ha már rendelkezik egy nagy folyamatcsatornával, amely olyan munkafolyamatokat fog össze, amelyeket külön szeretne összehangolni, a táblák egy új folyamatcsatornába való áthelyezésével bontsa kisebb folyamatcsatornákra. Lásd Táblák áthelyezése csővezetékek között.
Lakeflow feladatok
A Lakeflow Jobsban több feladatot is koordinálhat egy adatfeldolgozási munkafolyamat megvalósításához. Ha egy folyamatot szeretne belefoglalni egy feladatba, használja a Folyamat feladatot egy feladat létrehozásakor. Tekintse meg a pipeline feladatot a munkákhoz.
Apache Airflow
Az Apache Airflow egy nyílt forráskódú megoldás az adat-munkafolyamatok kezelésére és ütemezésére. Az Airflow a munkafolyamatokat a műveletek irányított aciklikus gráfjaiként (DAG-k) jelöli. Definiálhat egy munkafolyamatot egy Python-fájlban, és az Airflow kezeli az ütemezést és a végrehajtást. Az Airflow és az Azure Databricks telepítésével és használatával kapcsolatos információkért lásd: Orchestrate Lakeflow Jobs with Apache Airflow.
Ahhoz, hogy egy pipelinet egy Airflow-munkafolyamat részeként futtasson, használja a DatabricksSubmitRunOperatort.
Requirements
A Lakeflow-folyamatok Airflow-támogatásának használatához a következők szükségesek:
- Airflow 2.1.0-s vagy újabb verzió.
- A Databricks szolgáltatói csomag 2.1.0-s vagy újabb verziója.
Example
Az alábbi példa létrehoz egy Airflow DAG-t, amely elindítja a folyamat frissítését az azonosítóval 8279d543-063c-4d63-9926-dae38e35ce8b:
from airflow import DAG
from airflow.providers.databricks.operators.databricks import DatabricksSubmitRunOperator
from airflow.utils.dates import days_ago
default_args = {
'owner': 'airflow'
}
with DAG('ldp',
start_date=days_ago(2),
schedule_interval="@once",
default_args=default_args
) as dag:
opr_run_now=DatabricksSubmitRunOperator(
task_id='run_now',
databricks_conn_id='CONNECTION_ID',
pipeline_task={"pipeline_id": "8279d543-063c-4d63-9926-dae38e35ce8b"}
)
Cserélje le a CONNECTION_ID-t a munkaterülethez tartozó Airflow-kapcsolat azonosítójára.
Mentse ezt a példát a airflow/dags könyvtárban, és használja az Airflow felhasználói felületét a DAG megtekintéséhez és aktiválásához . A folyamatfrissítés részleteinek megtekintéséhez használja a folyamat felhasználói felületét.
Azure Data Factory
Megjegyzés:
A Lakeflow-folyamatok és Azure Data Factory mindegyik tartalmazza az újrapróbálkozások számának konfigurálását meghibásodás esetén. Ha az újrapróbálkoztatási értékek a folyamaton és a folyamatot meghívó Azure Data Factory tevékenységen vannak konfigurálva, akkor az újrapróbálkozások száma az Azure Data Factory újrapróbálkozásainak értékét a folyamat újrapróbálkozásainak értékével megszorozva adódik.
Ha például egy folyamat frissítése sikertelen, a folyamat alapértelmezés szerint legfeljebb ötször újrapróbálkozza a frissítést. Ha az Azure Data Factory újrapróbálkozása háromra van állítva, és a folyamat öt újrapróbálkozás alapértelmezett értékét használja, a sikertelen folyamat akár tizenötször is újrapróbálkozhat. A folyamatfrissítések sikertelensége esetén a túlzott újrapróbálkozási kísérletek elkerülése érdekében a Databricks javasolja az újrapróbálkozások számának korlátozását a folyamat konfigurálásakor vagy a folyamatot meghívó Azure Data Factory-tevékenység konfigurálásakor.
A folyamat újrapróbálkozási konfigurációjának módosításához használja a pipelines.numUpdateRetryAttempts beállítást a folyamat konfigurálásakor.
Az Azure Data Factory egy felhőalapú ETL-szolgáltatás, amely lehetővé teszi az adatintegrációs és átalakítási munkafolyamatok vezénylét. Az Azure Data Factory közvetlenül támogatja az Azure Databricks-feladatok munkafolyamatokban való futtatását, beleértve a jegyzetfüzeteket, a JAR-feladatokat és a Python-szkripteket. Munkafolyamatba is felvehet egy pipeline-t, ha meghívja a pipeline REST API-ját egy Azure Data Factory webtevékenységből. Például folyamatfrissítés aktiválása az Azure Data Factoryből:
Hozzon létre egy adat-előállítót , vagy nyisson meg egy meglévő adat-előállítót.
Amikor a létrehozás befejeződött, nyissa meg a data factory lapját, és kattintson az Azure Data Factory Studio megnyitása csempére. Megjelenik az Azure Data Factory felhasználói felülete.
Hozzon létre egy új Azure Data Factory-folyamatot az Azure Data Factory Studio felhasználói felületén az Új legördülő menü Folyamat elemének kiválasztásával.
A Tevékenységek eszközkészletben bontsa ki az Általános elemet, és húzza a webes tevékenységet a folyamatvászonra. Kattintson a Beállítások fülre , és adja meg a következő értékeket:
Megjegyzés:
Ajánlott biztonsági eljárásként, ha automatizált eszközökkel, rendszerekkel, szkriptekkel és alkalmazásokkal hitelesít, a Databricks azt javasolja, hogy munkaterület-felhasználók helyett a szolgáltatásnevekhez tartozó személyes hozzáférési jogkivonatokat használja. Szolgáltatási főszereplők jogkivonatainak létrehozásához lásd a Szolgáltatási főszereplők jogkivonatainak kezelése című részt.
URL-cím:
https://<databricks-instance>/api/2.0/pipelines/<pipeline-id>/updates.Cserélje le
<get-workspace-instance>.Cserélje le
<pipeline-id>a folyamatazonosítóra.Metódus: Válassza a POST lehetőséget a legördülő menüből.
Fejlécek: Kattintson az + Új gombra. A Név szövegmezőbe írja be a következőt
Authorization: Az Érték szövegmezőbe írja be a következőtBearer <personal-access-token>:Cserélje le
<personal-access-token>egy Azure Databricks személyes hozzáférési jogkivonatra.Törzs: További kérelemparaméterek megadásához adjon meg egy JSON-dokumentumot, amely tartalmazza a paramétereket. Például egy frissítés indításához és a folyamat összes adatának újrafeldolgozásához:
{"full_refresh": "true"}. Ha nincsenek további kérelemparaméterek, írjon be üres zárójeleket ({}).
A webes tevékenység teszteléséhez kattintson a Hibakeresés gombra a Data Factory felhasználói felületén található folyamat eszköztárán. A futtatás kimenete és állapota, beleértve a hibákat is, az Azure Data Factory-folyamat Kimenet lapján jelenik meg. A folyamatfrissítés részleteinek megtekintéséhez használja a folyamatok felhasználói felületét.
Jótanács
Gyakori munkafolyamat-követelmény, hogy egy feladatot egy korábbi tevékenység befejezése után kezdjen el. Mivel a folyamatkérés updates aszinkron, és a frissítés elindítása után, de még a befejezés előtt visszatér, a Azure Data Factory folyamat azon tevékenységeinek, amelyeknek függősége van a folyamatfrissítéstől, várniuk kell a frissítés befejezésére. A frissítés befejezésére való várakozás egyik lehetősége, ha a folyamatfrissítést aktiváló webes tevékenységet követő Until tevékenységet ad hozzá. A Until tevékenységben:
- Adjon hozzá egy várakozási tevékenységet , hogy a frissítés befejezéséhez konfigurált számú másodpercet várjon.
- Adjon hozzá egy webes tevékenységet a várakozási tevékenység után, amely a folyamatfrissítés részleteire vonatkozó kérést használja a frissítés állapotának lekéréséhez. A
stateválasz mezője a frissítés aktuális állapotát adja vissza, beleértve azt is, hogy befejeződött-e. - Használja a mező értékét a
stateUntil tevékenység befejezési feltételének beállításához. Használhat egy Változó beállítása tevékenységet is, hogy hozzáadjon egy folyamatváltozót azstateérték alapján, és ezt a változót a megszüntetési feltételhez használja.