Oktatóanyag: Térinformatikai folyamat létrehozása natív térbeli típusok használatával

Létrehoz és üzembe helyez egy olyan pipeline-t, amely GPS-adatokat fogad, a koordinátákat natív térbeli típusokká alakítja, majd a raktári geokerítésekkel összekapcsolva követi az érkezéseket a Lakeflow pipelines adatvezénylési funkcióival és az Auto Loaderrel. Ez az oktatóanyag a Databricks natív térbeli típusait (GEOMETRY, ) és olyan beépített térbeli függvényeket használ, GEOGRAPHYmint ST_Pointa , ST_GeomFromWKTés ST_Contains, hogy külső kódtárak nélkül is nagy léptékben futtathassa a térinformatikai munkafolyamatokat.

Ebben az oktatóanyagban a következőket meg fogja tanulni:

  • Hozzon létre egy csővezetéket, és generáljon minta GPS- és geokerítés-adatokat egy Unity Catalog-kötetben.
  • A nyers GPS-pingek fokozatos betöltése az Auto Loaderrel egy bronz streamelési táblába történik.
  • Készítsen egy ezüst adatfolyam táblát, amely a szélességet és a hosszúságot natív GEOMETRY ponttá alakítja.
  • A raktár geofenceseinek materializált nézetének létrehozása WKT-sokszögekből.
  • Térbeli illesztés futtatása, hogy létrehozzon egy táblázatot a raktár érkezéseiről (melyik eszköz lépett be melyik geokerítésbe).

Az eredmény egy medál stílusú folyamat: bronz (nyers GPS), ezüst (geometriai pontok) és arany (geofences és érkezési események). További információért lásd: Mi a medallion lakehouse architektúra?

Követelmények

Az oktatóanyag elvégzéséhez meg kell felelnie a következő követelményeknek:

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

Hozzon létre egy új ETL-folyamatot, és állítsa be a táblák alapértelmezett katalógusát és sémáját.

  1. A munkaterületen kattintson a Plusz ikonra.Új az oldalsávon, majd válassza az ETL-folyamat lehetőséget. Ezzel megnyitja a folyamatszerkesztőt egy olyan alapértelmezett folyamatnévvel, mint a New Pipeline <date> <time>.

  2. Jelölje ki a nevet, és adjon meg egy leíró nevet, például Spatial pipeline tutorial.

  3. A név jobb oldalán kattintson a katalógusra és a sémára az írási engedélyekkel rendelkező alapértelmezett beállítások kiválasztásához.

    Ez a katalógus és séma alapértelmezés szerint akkor használatos, ha nem ad meg katalógust vagy sémát a kódban. Cserélje le a <catalog> és a <schema> az alábbi lépésekben az itt kiválasztott értékekre.

  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 Lakeflow Pipelines szerkesztő a folyamatában található mintafájlokkal nyílik meg. Ezután hozza létre a minta GPS- és földrajziadat-adatokat.

2. lépés: A minta GPS- és geokerítésadatok létrehozása

Ez a lépés mintaadatokat hoz létre egy kötetben: nyers GPS-pingeket (JSON) és raktári geofeneket (JSON WKT-sokszögekkel). A GPS-pontok egy határolókeretben jönnek létre, amely átfedésben van a két raktári sokszöggel, így a térbeli illesztés egy későbbi lépésben visszaadja az érkezési sorokat. Ezt a lépést kihagyhatja, ha már rendelkezik saját adatokkal egy kötetben vagy táblában.

  1. A Lakeflow Pipelines-szerkesztő eszközböngészőjében kattintson a Plusz ikonra,Hozzáadás, majd Feltárás.

  2. Állítsa be a Név beállítást Setup spatial data, válassza a Pythont, és hagyja meg az alapértelmezett célmappát.

  3. Kattintson a Létrehozás gombra.

  4. Illessze be az alábbi kódot az új jegyzetfüzetbe. Cserélje le a <catalog> és <schema> bejegyzéseket az 1. lépésben beállított alapértelmezett katalógust és sémát tartalmazó értékekre.

    A jegyzetfüzetben az alábbi kóddal generálhat GPS- és geofizikai adatokat.

    from pyspark.sql import functions as F
    
    catalog = "<catalog>"   # for example, "main"
    schema = "<schema>"    # for example, "default"
    
    spark.sql(f"USE CATALOG `{catalog}`")
    spark.sql(f"USE SCHEMA `{schema}`")
    spark.sql(f"CREATE VOLUME IF NOT EXISTS `{catalog}`.`{schema}`.`raw_data`")
    volume_base = f"/Volumes/{catalog}/{schema}/raw_data"
    
    # GPS: 5000 rows in a box that overlaps both warehouse geofences (LA area)
    gps_path = f"{volume_base}/gps"
    df_gps = (
        spark.range(0, 5000)
        .repartition(10)
        .select(
            F.format_string("device_%d", F.col("id").cast("long")).alias("device_id"),
            F.current_timestamp().alias("timestamp"),
            (-118.3 + F.rand() * 0.2).alias("longitude"),   # -118.3 to -118.1
            (34.0 + F.rand() * 0.2).alias("latitude"),     # 34.0 to 34.2
        )
    )
    df_gps.write.format("json").mode("overwrite").save(gps_path)
    print(f"Wrote 5000 GPS rows to {gps_path}")
    
    # Geofences: two warehouse polygons (WKT) in the same region
    geofences_path = f"{volume_base}/geofences"
    geofences_data = [
        ("Warehouse_A", "POLYGON ((-118.35 34.02, -118.25 34.02, -118.25 34.08, -118.35 34.08, -118.35 34.02))"),
        ("Warehouse_B", "POLYGON ((-118.20 34.05, -118.12 34.05, -118.12 34.12, -118.20 34.12, -118.20 34.05))"),
    ]
    df_geo = spark.createDataFrame(geofences_data, ["warehouse_name", "boundary_wkt"])
    df_geo.write.format("json").mode("overwrite").save(geofences_path)
    print(f"Wrote {len(geofences_data)} geofences to {geofences_path}")
    
  5. Futtassa a jegyzetfüzetcellát (Shift + Enter).

A futtatás befejezése után a kötet tartalmazza gps a (nyers pingeket) és geofences a (WKT-ben lévő sokszögeket). A következő lépésben a GPS-adatokat egy bronz táblázatba tölti.

3. lépés: GPS-adatok betöltése egy bronz adatfolyam táblába

Töltse be a nyers GPS JSON-t a kötetből növekményesen az Auto Loader használatával, és írja be egy bronz szintű adatfolyam táblába.

  1. Az eszközböngészőben kattintson a Plusz ikonra.Hozzáadás, majd átalakítás.

  2. Állítsa be a Név beállítást gps_bronze, válassza az SQL vagy a Python lehetőséget, majd kattintson a Létrehozás gombra.

  3. Cserélje le a fájl tartalmát a következőre (használja a nyelvnek megfelelő lapot). Cserélje le a <catalog> és <schema> az alapértelmezett katalógusra és sémára.

    SQL

    CREATE OR REFRESH STREAMING TABLE gps_bronze
    COMMENT "Raw GPS pings ingested from volume using Auto Loader";
    
    CREATE FLOW gps_bronze_ingest_flow AS
    INSERT INTO gps_bronze BY NAME
    SELECT *
    FROM STREAM read_files(
      "/Volumes/<catalog>/<schema>/raw_data/gps",
      format => "json",
      inferColumnTypes => "true"
    )
    

    Python

    from pyspark import pipelines as dp
    
    path = "/Volumes/<catalog>/<schema>/raw_data/gps"
    
    dp.create_streaming_table(
      name="gps_bronze",
      comment="Raw GPS pings ingested from volume using Auto Loader",
    )
    
    @dp.append_flow(target="gps_bronze", name="gps_bronze_ingest_flow")
    def gps_bronze_ingest_flow():
        return (
            spark.readStream.format("cloudFiles")
            .option("cloudFiles.format", "json")
            .option("cloudFiles.inferColumnTypes", "true")
            .load(path)
        )
    
  4. Kattintson a Lejátszás ikonra.Futtassa a fájlt vagy a futtatási folyamatot egy frissítés futtatásához.

Amikor a frissítés befejeződött, a folyamatdiagram megjeleníti a táblát gps_bronze . Ezután adjon hozzá egy ezüsttáblát, amely a koordinátákat natív geometriai ponttá alakítja.

4. lépés: Ezüst adatfolyam táblázat hozzáadása geometriai pontokkal

Hozzon létre egy streamelési táblát, amely a bronz táblából olvas, és a GEOMETRY használatával hozzáad egy ST_Point(longitude, latitude) oszlopot.

  1. Az eszközböngészőben kattintson a Plusz ikonra.Hozzáadás, majd átalakítás.

  2. Állítsa be a Név beállítást raw_gps_silver, válassza az SQL vagy a Python lehetőséget, majd kattintson a Létrehozás gombra.

  3. Illessze be a következő kódot az új fájlba.

    SQL

    CREATE OR REFRESH STREAMING TABLE raw_gps_silver
    COMMENT "GPS pings with native geometry point for spatial joins";
    
    CREATE FLOW raw_gps_silver_flow AS
    INSERT INTO raw_gps_silver BY NAME
    SELECT
      device_id,
      timestamp,
      longitude,
      latitude,
      ST_Point(longitude, latitude) AS point_geom
    FROM STREAM(gps_bronze)
    

    Python

    from pyspark import pipelines as dp
    from pyspark.sql import functions as F
    
    dp.create_streaming_table(
      name="raw_gps_silver",
      comment="GPS pings with native geometry point for spatial joins",
    )
    
    @dp.append_flow(target="raw_gps_silver", name="raw_gps_silver_flow")
    def raw_gps_silver_flow():
        return (
            spark.readStream.table("gps_bronze")
            .select(
                "device_id",
                "timestamp",
                "longitude",
                "latitude",
                F.expr("ST_Point(longitude, latitude)").alias("point_geom"),
            )
        )
    
  4. Kattintson a Lejátszás ikonra.Futtassa a fájlt vagy a futtatási folyamatot.

Az adatcsatorna-diagram most már megjeleníti a gps_bronze és raw_gps_silver elemeket. Ezután adja hozzá a raktár geokerítéseit materializált nézetként.

5. lépés: A raktár geofences aranytáblájának létrehozása

Hozzon létre egy materializált nézetet, amely beolvassa a kötet geofenceseit, és a WKT oszlopot oszlopmá GEOMETRY alakítja a segítségével ST_GeomFromWKT.

  1. Az eszközböngészőben kattintson a Plusz ikonra.Hozzáadás, majd átalakítás.

  2. Állítsa be a Név beállítást warehouse_geofences_gold, válassza az SQL vagy a Python lehetőséget, majd kattintson a Létrehozás gombra.

  3. Illessze be a következő kódot. Cserélje le a <catalog> és <schema> az alapértelmezett katalógusra és sémára.

    SQL

    CREATE OR REPLACE MATERIALIZED VIEW warehouse_geofences_gold AS
    SELECT
      warehouse_name,
      ST_GeomFromWKT(boundary_wkt) AS boundary_geom
    FROM read_files(
      "/Volumes/<catalog>/<schema>/raw_data/geofences",
      format => "json"
    )
    

    Python

    from pyspark import pipelines as dp
    from pyspark.sql import functions as F
    
    path = "/Volumes/<catalog>/<schema>/raw_data/geofences"
    
    @dp.table(name="warehouse_geofences_gold", comment="Warehouse geofence polygons as geometry")
    def warehouse_geofences_gold():
        return (
            spark.read.format("json").load(path).select(
                "warehouse_name",
                F.expr("ST_GeomFromWKT(boundary_wkt)").alias("boundary_geom"),
            )
        )
    
  4. Kattintson a Lejátszás ikonra.Futtassa a fájlt vagy a futtatási folyamatot.

A folyamat most már tartalmazza a geofences táblát. Ezután adja hozzá a térbeli illesztést a raktári érkezések számításához.

6. lépés: A raktár érkezési táblájának létrehozása térbeli illesztéssel

Adjon hozzá egy materializált nézetet, amely összekapcsolja az ezüst GPS-pontokat a geofencesekkel ST_Contains(boundary_geom, point_geom) annak meghatározásához, hogy egy eszköz mikor található egy raktári sokszögben.

  1. Az eszközböngészőben kattintson a Plusz ikonra.Hozzáadás, majd átalakítás.

  2. Állítsa be a Név beállítást warehouse_arrivals, válassza az SQL vagy a Python lehetőséget, majd kattintson a Létrehozás gombra.

  3. Illessze be a következő kódot.

    SQL

    CREATE OR REPLACE MATERIALIZED VIEW warehouse_arrivals AS
    SELECT
      g.device_id,
      g.timestamp,
      w.warehouse_name
    FROM raw_gps_silver g
    JOIN warehouse_geofences_gold w
      ON ST_Contains(w.boundary_geom, g.point_geom)
    

    Python

    from pyspark import pipelines as dp
    from pyspark.sql import functions as F
    
    @dp.table(name="warehouse_arrivals", comment="Devices that have entered a warehouse geofence")
    def warehouse_arrivals():
        g = spark.read.table("raw_gps_silver")
        w = spark.read.table("warehouse_geofences_gold")
        return (
            g.alias("g")
            .join(w.alias("w"), F.expr("ST_Contains(w.boundary_geom, g.point_geom)"))
            .select(
                F.col("g.device_id").alias("device_id"),
                F.col("g.timestamp").alias("timestamp"),
                F.col("w.warehouse_name").alias("warehouse_name"),
            )
        )
    
  4. Kattintson a Lejátszás ikonra.Futtassa a fájlt vagy a futtatási folyamatot.

A frissítés befejeződésekor a folyamatdiagram mind a négy adathalmazt jeleníti meg: gps_bronze, raw_gps_silver, warehouse_geofences_goldés warehouse_arrivals.

A térbeli illesztés ellenőrzése

Győződjön meg arról, hogy a térbeli illesztés sorokat hozott létre, amelyekben az ezüsttáblából származó pontok, amelyek egy geofence belsejébe esnek, megjelennek itt: warehouse_arrivals. Futtassa az alábbiak egyikét egy jegyzetfüzetben vagy SQL-szerkesztőben (használja ugyanazt a katalógust és sémát, mint a folyamatcél).

Érkezések száma raktár szerint (SQL):

SELECT warehouse_name, COUNT(*) AS arrival_count
FROM warehouse_arrivals
GROUP BY warehouse_name
ORDER BY warehouse_name;

Nem nulla értékeket kell látnia a Warehouse_A és Warehouse_B esetében (a minta GPS-adatok átfedésben vannak mindkét sokszöggel). Mintasorok vizsgálata:

SELECT device_id, timestamp, warehouse_name
FROM warehouse_arrivals
ORDER BY timestamp DESC
LIMIT 10;

Ugyanezek az ellenőrzések a Pythonban (jegyzetfüzetben):

# Count by warehouse
display(spark.table("warehouse_arrivals").groupBy("warehouse_name").count().orderBy("warehouse_name"))

# Sample rows
display(spark.table("warehouse_arrivals").orderBy("timestamp", ascending=False).limit(10))

Ha sorok jelennek meg warehouse_arrivals, akkor a ST_Contains(boundary_geom, point_geom) illesztés megfelelően működik.

7. lépés: A folyamat ütemezése (nem kötelező)

Ha naprakészen szeretné tartani a folyamatot, amikor az új GPS-adatok megjelennek a kötetben, hozzon létre egy feladatot a folyamat ütemezés szerinti futtatásához.

  1. A szerkesztő felületének tetején válassza ki az Ütemezés gombot.
  2. Ha megjelenik az Ütemezések párbeszédpanel, válassza az Ütemezés hozzáadása lehetőséget.
  3. Ha szeretné, adjon nevet a feladatnak.
  4. Alapértelmezés szerint az ütemezés naponta egyszer fut. Elfogadhatja ezt, vagy beállíthatja a sajátját. A Speciális lehetőség kiválasztásával beállíthat egy adott időpontot; A további lehetőségek lehetővé teszi a futtatási értesítések hozzáadását.
  5. Válassza a Létrehozás lehetőséget az ütemezés alkalmazásához.

A feladatfuttatásokról további információt a Lakeflow-feladatok monitorozása című témakörben talál.

További erőforrások