Oktatóanyag: Hozd létre az első pipeline-t a Lakeflow Pipelines Editor használatával

Az Automatikus betöltővel új Lakeflow-folyamatot hozhat létre az adatok vezényléséhez, majd kibővítheti a mintafolyamatot az adatok megtisztításával és egy lekérdezés létrehozásával, hogy megtalálja a 100 legjobb felhasználót.

Ebben az oktatóanyagban megtudhatja, hogyan használhatja a Lakeflow Pipelines-szerkesztőt a következőkre:

  • Hozzon létre egy új folyamatot az alapértelmezett mappastruktúrával, és kezdje a mintafájlok készletével.
  • Az adatminőségre vonatkozó korlátozások meghatározása elvárások alapján.
  • A szerkesztőfunkciókkal új átalakítással bővítheti a folyamatot az adatok elemzéséhez.

Requirements

Az oktatóanyag megkezdése előtt a következőket kell tennie:

  • Be kell jelentkeznie egy Azure Databricks munkaterületre.
  • Engedélyezze a Unity-katalógust a munkaterületen.
  • Rendelkezik engedéllyel számítási erőforrás létrehozására vagy egy számítási erőforráshoz való hozzáférésre.
  • Új séma katalógusban való létrehozására vonatkozó engedélyekkel rendelkezik. A szükséges engedélyek a következők: ALL PRIVILEGES vagyUSE CATALOG.CREATE SCHEMA
  • A folyamatok és kimeneteik létrehozásához, futtatásához, frissítéséhez és megtekintéséhez szükséges jogosultságok teljes készletét a folyamatok identitásainak, engedélyeinek és jogosultságainak kezelése című témakörben találhatja meg.

1. lépés: Folyamat létrehozása

Ebben a lépésben létrehoz egy folyamatot az alapértelmezett mappastruktúra és kódminták használatával. A kódminták a users mintaadatforrás táblára wanderbricks hivatkoznak.

  1. Az Azure Databricks munkaterületen kattintson a Plusz ikonra.Új, majd a Folyamat ikonra.ETL-folyamat. Ezzel megnyitja a folyamatszerkesztőt egy olyan alapértelmezett folyamatnévvel, mint a New Pipeline <date> <time>.

  2. (Nem kötelező) Válassza ki a nevet, és adjon meg egy leíró nevet a folyamatnak.

  3. (Nem kötelező) A név jobb oldalán kattintson a katalógusra és a sémára a különböző alapértelmezett értékek beállításához.

  4. (Nem kötelező) Az Ön számára létrehozott my_transformation forrásfájlban válassza Python vagy SQL a nyelvi legördülő listából a fájl nyelvének beállításához.

  5. Kattintson a Kód ikonra.Használjon mintakódot.

    A kiválasztott nyelv mintakódja megjelenik a my_transformation mappában lévő forrásfájlban transformations . A kimeneti adathalmazok még nem lettek létrehozva, és a folyamatdiagram a képernyő jobb oldalán üres.

  6. A folyamatkód (a transformations mappában lévő kód) futtatásához kattintson a folyamat futtatása elemre a képernyő jobb felső részén.

    A futtatás befejezése után a munkaterület alsó része megjeleníti a létrehozott két új táblát, sample_users_<date_time> valamint sample_aggregation_<date_time>a . A munkaterület jobb oldalán lévő Folyamatábra most a két táblát is megjeleníti, beleértve azt is, hogy a(z) sample_users a(z) sample_aggregation forrása. Jegyezze fel a teljes sample_users_<date_time> táblanevet. A következő lépésben hivatkozhat rá.

2. lépés: Adatminőség-ellenőrzések alkalmazása

Ebben a lépésben adatminőség-ellenőrzést ad hozzá a sample_users táblához. A pipeline elvárásokat használva korlátozza az adatokat. Ebben az esetben töröl minden olyan felhasználói rekordot, amely nem rendelkezik érvényes e-mail-címmel, és a megtisztított táblát users_cleaneda következőképpen adja ki.

  1. A bal oldali folyamat-objektumböngészőben kattintson a Plusz ikonra, és válassza az Átalakítás lehetőséget.

  2. Az Új átalakítási fájl létrehozása párbeszédpanelen végezze el a következő beállításokat:

    • Válassza a Python vagy SQLLanguage lehetőséget. Ennek nem kell megegyeznie az előző választásoddal.
    • Adjon nevet a fájlnak. Ebben az esetben válassza a users_cleanedlehetőséget.
    • A Cél elérési útja beállításnál hagyja meg az alapértelmezett értéket.
    • Adathalmaztípus esetén hagyja meg a nincs kijelölve, vagy válassza a Materialized nézetet. Ha a Materialized nézetet választja, az létrehoz egy mintakódot.
  3. Kattintson a Létrehozás gombra az átalakítási kódfájl létrehozásához.

  4. Az új kódfájlban szerkessze a kódot az alábbiak szerint (használja az SQL-t vagy a Pythont az előző képernyős választása alapján). Cserélje le a(z) sample_users_<date_time> elemet az előző szakaszban szereplő sample_users tábla teljes nevére.

    SQL

    -- Drop all rows that do not have an email address
    
    CREATE MATERIALIZED VIEW users_cleaned
    (
      CONSTRAINT non_null_email EXPECT (email IS NOT NULL) ON VIOLATION DROP ROW
    ) AS
    SELECT *
    FROM sample_users_<date_time>;
    

    Python

    from pyspark import pipelines as dp
    
    # Drop all rows that do not have an email address
    
    @dp.materialized_view
    @dp.expect_or_drop("no null emails", "email IS NOT NULL")
    def users_cleaned():
        return (
            spark.read.table("sample_users_<date_time>")
        )
    
  5. Kattintson a Folyamat futtatása elemre a folyamat frissítéséhez. Most már három táblával kell rendelkeznie.

3. lépés: A legnépszerűbb felhasználók elemzése

Ezután szerezze be a 100 legjobb felhasználót a létrehozott foglalások száma alapján. Csatlakoztassa a wanderbricks.bookings táblát a users_cleaned materializált nézethez.

  1. A bal oldali folyamat-objektumböngészőben kattintson a Plusz ikonra, és válassza az Átalakítás lehetőséget.

  2. Az Új átalakítási fájl létrehozása párbeszédpanelen végezze el a következő beállításokat:

    • Válassza a Python vagy SQLLanguage lehetőséget. Ennek nem kell megegyeznie a korábbi kijelölésekkel.
    • Adjon nevet a fájlnak. Ebben az esetben válassza a users_and_bookingslehetőséget.
    • A Cél elérési útja beállításnál hagyja meg az alapértelmezett értéket.
    • Adathalmaztípus esetén hagyja a Nincs kijelölve értéket.
  3. Kattintson a Létrehozás gombra az átalakítási kódfájl létrehozásához.

  4. Az új kódfájlban szerkessze a kódot az alábbiak szerint (használja az SQL-t vagy a Pythont az előző képernyős választása alapján).

    SQL

    -- Get the top 100 users by number of bookings
    
    CREATE OR REFRESH MATERIALIZED VIEW users_and_bookings AS
    SELECT u.name AS name, COUNT(b.booking_id) AS booking_count
    FROM users_cleaned u
    JOIN samples.wanderbricks.bookings b ON u.user_id = b.user_id
    GROUP BY u.name
    ORDER BY booking_count DESC
    LIMIT 100;
    

    Python

    from pyspark import pipelines as dp
    from pyspark.sql.functions import col, count, desc
    
    # Get the top 100 users by number of bookings
    
    @dp.materialized_view
    def users_and_bookings():
        return (
            spark.read.table("users_cleaned")
            .join(spark.read.table("samples.wanderbricks.bookings"), "user_id")
            .groupBy(col("name"))
            .agg(count("booking_id").alias("booking_count"))
            .orderBy(desc("booking_count"))
            .limit(100)
        )
    
  5. Kattintson a Folyamat futtatása az adathalmazok frissítéséhez. Ha a futtatás befejeződött, a Folyamatdiagramon láthatja, hogy négy tábla van, köztük az új users_and_bookings tábla.

    Pipeline grafikon, amely négy táblát mutat a pipeline-ban

További erőforrások

Most, hogy megtanulta, hogyan használhatja a Lakeflow-folyamatok szerkesztőjének néhány funkcióját, és létrehozott egy folyamatot, az alábbiakban további információkat talál: