Aller au contenu principal

Didacticiel : Construire un pipeline géospatial avec des types spatiaux natifs

Vous créez et déployez un pipeline qui ingère des données GPS, convertit les coordonnées en types spatiaux natifs et se joint aux géorepérages de warehouse pour suivre les arrivées à l'aide des LakeFlow Pipelines pour l'orchestration des données et Auto Loader. Ce tutoriel utilise les types spatiaux natifs de Databricks (GEOMETRY, GEOGRAPHY) et des fonctions spatiales intégrées telles que ST_Point, ST_GeomFromWKT et ST_Contains, afin que vous puissiez exécuter des workflows géospatiaux à l'échelle sans bibliothèques externes.

Dans ce tutoriel, vous allez :

  • Créer un pipeline et générer des exemples de données GPS et de géorepérage dans un volume Unity Catalog.
  • Ingérez les pings GPS bruts de manière incrémentielle avec Auto Loader dans une table de streaming bronze.
  • Construire une table de streaming silver qui convertit la latitude et la longitude en un point GEOMETRY natif.
  • Créer une vue matérialisée des géofences de warehouse à partir de polygones WKT.
  • Exécuter une jointure spatiale pour produire une table des arrivées de warehouse (quel appareil est entré dans quelle zone de géorepérage).

Le résultat est un pipeline de style médaillon : bronze (GPS brut), argent (points en tant que géométrie) et Gold (géorepérages et événements d'arrivée). Voir Qu'est-ce que l'architecture lakehouse en médaillon ? pour plus d'informations.

Exigences

Pour suivre ce didacticiel, vous devez satisfaire aux exigences suivantes :

Étape 1 : Créez une pipeline

Créez un nouveau pipeline ETL et définissez le catalogue et le schéma default pour vos tables.

  1. Dans votre Workspace, cliquez sur Icône Plus. Nouveau dans la barre latérale, puis sélectionnez Pipeline ETL . Ceci ouvre l'éditeur de pipeline avec un nom de pipeline default comme New Pipeline <date> <time>.

  2. Sélectionnez le nom et saisissez un nom descriptif, tel que Spatial pipeline tutorial.

  3. À droite du nom, cliquez sur le catalogue et le schéma pour choisir les valeurs par default pour lesquelles vous disposez de permissions d’écriture.

    Ce catalogue et ce schéma sont utilisés par default si vous ne spécifiez pas de catalogue ou de schéma dans votre code. Remplacez <catalog> et <schema> dans les étapes suivantes par les valeurs que vous choisissez ici.

  4. (Facultatif) Dans le fichier source my_transformation créé pour vous, sélectionnez Python ou SQL dans la liste déroulante des langues pour définir la langue du fichier.

  5. Cliquez Icône de code. sur **Utiliser l'exemple de code**.

L'éditeur de Lakeflow Pipelines s'ouvre avec des exemples de fichiers dans votre pipeline. Ensuite, créez les exemples de données GPS et de géofence.

Étape 2 : créez l’échantillon de données GPS et de géorepérage

Cette étape génère des échantillons de données en volume : des pings GPS bruts (JSON) et des géo-clôtures de warehouse (JSON avec des polygones WKT). Les points GPS sont générés dans une boîte englobante qui chevauche les deux polygones de warehouse, de sorte que la jointure spatiale à une étape ultérieure renverra des lignes d'arrivée. Vous pouvez ignorer cette étape si vous avez déjà vos propres données dans un volume ou une table.

  1. Dans l'éditeur Lakeflow Pipelines, dans le navigateur d'assets, cliquez sur Icône Plus. Ajouter , puis sur Exploration .

  2. Définissez Nom sur Setup spatial data, choisissez Python , et laissez le dossier de destination par default.

  3. Cliquez sur Créer .

  4. Dans le nouveau notebook, collez le code suivant. Remplacez <catalog> et <schema> par le catalogue et le schéma default que vous avez définis à l'Étape 1.

    Utilisez le code suivant dans le Notebook pour générer des données GPS et de géorepérage.

    Python
    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. Exécutez la cellule du Notebook (Shift + Enter).

Une fois l'exécution terminée, le volume contient gps (pings bruts) et geofences (polygones en WKT). À l'étape suivante, vous ingérez les données GPS dans une table bronze.

Étape 3 : Ingestion des données GPS dans une table de streaming bronze

Ingérez le JSON GPS brut du volume de manière incrémentielle à l'aide d'Auto Loader et écrivez dans une table de streaming bronze.

  1. Dans l'explorateur d'actifs, cliquez sur Icône Plus. Ajouter , puis sur Transformation .

  2. Définissez le Nom sur gps_bronze, choisissez SQL ou Python , et cliquez sur Créer .

  3. Remplacez le contenu du fichier par ce qui suit (utilisez le tab qui correspond à votre langue). Remplacez <catalog> et <schema> par votre catalogue et schéma default.

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"
)
  1. Cliquez sur Icône de lecture. Exécuter le fichier ou Exécuter le pipeline pour exécuter une mise à jour.

Une fois la mise à jour terminée, le Graphe du pipeline affiche la table gps_bronze. Ensuite, ajoutez une table argent qui convertit les coordonnées en un point de géométrie natif.

Étape 4 : Ajouter une table de streaming silver avec des points de géométrie

Créez une table de streaming qui lit à partir de la table bronze et ajoute une colonne GEOMETRY à l'aide de ST_Point(longitude, latitude).

  1. Dans l'explorateur d'actifs, cliquez sur Icône Plus. Ajouter , puis sur Transformation .

  2. Définissez le Nom sur raw_gps_silver, choisissez SQL ou Python , et cliquez sur Créer .

  3. Collez le code suivant dans le nouveau fichier.

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)
  1. Cliquez sur Icône de lecture. Exécuter le fichier ou Exécuter le pipeline .

Le graphe du pipeline affiche maintenant gps_bronze et raw_gps_silver. Ensuite, ajoutez les géorepérages du warehouse en tant que vue matérialisée.

Étape 5 : Créer la table Gold des geofences du warehouse

Créez une vue matérialisée qui lit les géorepérages du volume et convertit la colonne WKT en une colonne GEOMETRY à l'aide de ST_GeomFromWKT.

  1. Dans l'explorateur d'actifs, cliquez sur Icône Plus. Ajouter , puis sur Transformation .

  2. Définissez le Nom sur warehouse_geofences_gold, choisissez SQL ou Python , et cliquez sur Créer .

  3. Collez le code suivant. Remplacez <catalog> et <schema> par votre catalogue et schéma default.

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"
)
  1. Cliquez sur Icône de lecture. Exécuter le fichier ou Exécuter le pipeline .

Le pipeline inclut désormais la table geofences. Ensuite, ajoutez la jointure spatiale au compute des arrivées de warehouse.

Étape 6 : Créer la table d'arrivées du warehouse avec une jointure spatiale

Ajoutez une vue matérialisée qui joint les points GPS silver aux géofences en utilisant ST_Contains(boundary_geom, point_geom) pour déterminer quand un appareil se trouve à l'intérieur d'un polygone de warehouse.

  1. Dans l'explorateur d'actifs, cliquez sur Icône Plus. Ajouter , puis sur Transformation .

  2. Définissez le Nom sur warehouse_arrivals, choisissez SQL ou Python , et cliquez sur Créer .

  3. Collez le code suivant.

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)
  1. Cliquez sur Icône de lecture. Exécuter le fichier ou Exécuter le pipeline .

Lorsque la mise à jour est terminée, le graphe du pipeline affiche les quatre datasets : gps_bronze, raw_gps_silver, warehouse_geofences_gold et warehouse_arrivals.

Vérifier la jointure spatiale

Confirmez que la jointure spatiale a produit des lignes : les points de la table silver qui se trouvent à l'intérieur d'un geofence apparaissent dans warehouse_arrivals. Exécutez l'une des opérations suivantes dans un Notebook ou un éditeur SQL (utilisez le même catalogue et le même schéma que votre cible de pipeline).

Compter les arrivées par warehouse (SQL) :

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

Vous devriez voir des nombres non nuls pour Warehouse_A et Warehouse_B (les données GPS échantillonnées chevauchent les deux polygones). Pour inspecter les lignes d'échantillon :

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

Mêmes vérifications en Python (notebook) :

Python
# 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))

Si vous voyez des lignes dans warehouse_arrivals, la jointure ST_Contains(boundary_geom, point_geom) fonctionne correctement.

Étape 7 : planifier le pipeline (facultatif)

Pour maintenir le pipeline à jour lorsque de nouvelles données GPS arrivent dans le volume, créez un Job pour exécuter le pipeline selon un calendrier.

  1. En haut de l'éditeur, choisissez le bouton Planifier .
  2. Si la boîte de dialogue Planifications apparaît, choisissez Ajouter une planification .
  3. Vous pouvez, si vous le souhaitez, donner un nom au Job.
  4. Par default, la programmation s’exécute une fois par jour. Vous pouvez accepter cela ou définir le vôtre. Le choix d’ Avancé vous permet de définir une heure spécifique ; Plus d’options vous permet d’ajouter des notifications d’exécution.
  5. Sélectionnez Créer pour appliquer la planification.

Consultez Monitor Lakeflow Jobs pour plus d'informations sur les exécutions de Jobs.

Ressources supplémentaires