Adatfelvétel API-ból a pipeline-okban

Az API-ból történő befogadás azt jelenti, hogy HTTP-n keresztül húzzák ki az adatokat egy webszolgáltatásból, általában oldalas JSON-ként, nem pedig fájlból vagy adatbázisból olvasni. A fájlokkal vagy az üzenetbusszal ellentétben nincs beépített általános API forrás, így a hitelesítést, az oldalozást és a sebességkorlátokat magad kezeled. A Lakeflow-folyamatok három mintát támogatnak egy tetszőleges API-ból történő adatbeolvasáshoz. Melyik illeszkedik, az a hangerő és frissítési igényed függvénye.

Important

Mielőtt bármilyen egyedi API-felvételi kódot írnál, ellenőrizd, hogy létezik-e már menedzselt csatlakozó a forrásodhoz. A Lakeflow Connect számos elterjedt szoftverként szolgáltatás (SaaS) API-hoz beépített csatlakozókat szállít, mint például a Salesforce, Workday, ServiceNow és Google Analytics, és egyre több partner csatlakozó is megjelenik. Ha egy csatlakozó lefedi az adatforrásodat, akkor kezeli helyetted a hitelesítést, a lapozást és az inkrementális adatkinyerést, és ez szinte mindig kevesebb munkát jelent, mint egy egyedi fejlesztésű adatbeolvasás. Lásd: Felügyelt csatlakozók a Lakeflow Connectben. Az alábbi mintákat csak akkor használd, ha egyik csatlakozó sem megfelelő.

Előfeltételek

Minta kiválasztása

Nincs natív általános REST-API forrás a pipeline-ekben, ezért amikor egy tetszőleges API-ból húzol, válassz három mintát az adatmennyiség és a betöltés gyakorisága alapján:

Minta Használja amikor
Periodikus húzások mint materializált nézet Az adatmennyiségek kicsik vagy közepes méretűek, és a folyamat minden egyes futása során egyszer kerülnek lekérésre, például referenciaadatok, napi devizaárfolyamok vagy egy lapozott, de korlátozott terjedelmű API esetében.
Python adatforrás API Egy nagy forgalmú vagy streamelő API-t fokozatosan kell lekérdezned, ellenőrzőpontokkal rögzített előrehaladással, hogy újraindítás után ne kelljen mindent elölről újra beolvasni.
Leválasztott adatbetöltés Auto Loaderrel El szeretnéd különíteni az API-specifikus sajátosságokat a transzformációs logikádtól, és külön ráfordítás nélkül kapsz pontosan egyszeri fájlkövetést.

1. minta: Periodikus húzások mint materializált nézet

A kis-közepes méretű hasznos terhekhez, amelyeket egyszer húznak le csővezetékenként, írj egy Python függvényt, amely az API-t hívja és Spark DataFrame-et ad vissza. Mivel az adathalmaz egy materializált nézet, a folyamat minden alkalommal, amikor frissül, teljes egészében és idempotens módon újrafuttatja a függvényt.

Az alábbi lépések bemutatják, hogyan hozhat létre materializált nézetet időszakos lekérésekkel:

  1. Tárold az API-tokent egy titokban, majd képezd le egy Spark-konfigurációs tulajdonságra a folyamat beállításaiban, hogy a folyamat kódja olvasni tudja. Hozzáadjuk a tulajdonságot a spark_conf csővezeték klaszterkonfigurációjának blokkjához:

    {
      "clusters": [
        {
          "spark_conf": {
            "api.token": "{{secrets/<scope-name>/<secret-name>}}"
          }
        }
      ]
    }
    

    A következő lépésben a kód ezt az értéket a spark.conf.get("api.token")-vel olvassa. További információkért a titkok konfigurálásáról a csővezeték-beállításokban lásd: Biztonságos hozzáférés tárolási adatokhoz titkokkal a pipeline-ben.

  2. Definiáljunk egy materializált nézetet, amely az API-t hívja, és a választ DataFrame formájában adja vissza:

    import requests
    from pyspark import pipelines as dp
    from pyspark.sql import Row
    
    @dp.materialized_view(
        name="exchange_rates_bronze",
        comment="Daily FX rates pulled from a public REST API",
    )
    def exchange_rates_bronze():
        resp = requests.get(
            "https://api.example.com/v1/rates",
            params={"base": "USD"},
            headers={"Authorization": f"Bearer {spark.conf.get('api.token')}"},
            timeout=30,
        )
        resp.raise_for_status()
        rates = resp.json()["rates"]
        rows = [Row(currency=k, rate=float(v), as_of_date=resp.json()["date"]) for k, v in rates.items()]
        return spark.createDataFrame(rows)
    
  3. A függvényen belül az oldalozást úgy kezeljük, hogy az oldalak fölött ível, és összefűzed az eredményeket, mielőtt visszaadnád a DataFrame-et:

    import requests
    from pyspark import pipelines as dp
    from pyspark.sql import Row
    
    @dp.materialized_view(
        name="customers_bronze",
        comment="Customers pulled from a paginated REST API",
    )
    def customers_bronze():
        token = spark.conf.get("api.token")
        rows = []
        url = "https://api.example.com/v1/customers"
        while url:  # follow the API's next-page cursor until exhausted
            resp = requests.get(
                url,
                headers={"Authorization": f"Bearer {token}"},
                timeout=30,
            )
            resp.raise_for_status()
            payload = resp.json()
            rows.extend(Row(**record) for record in payload["data"])
            url = payload.get("next")  # None on the last page
        return spark.createDataFrame(rows)
    

    Adj újrapróbálkozási és visszalépési logikát a kérés köré a nagyobb hibatűrés érdekében.

Ez a minta minden pipeline frissítéskor újraolvassa az API teljes válaszát, tehát csak akkor használd, ha a hasznos terhelés korlátozott. Növekményes olvasásokhoz használja a 2. mintát.

2. minta : Nagy volumenű vagy streaming API-k a Python Data Source API-val

API-knál lépésesen kell lekérdezést igényelni offset követéssel, egyedi adatforrást kell valósítani a Spark Python Data Source API-jával. Ez megfelelő streamelési szemantikát biztosít, beleértve az ellenőrzőpontokkal mentett előrehaladást és az inkrementális olvasásokat, így újraindításkor az utolsó offsettől folytatódik a feldolgozás, ahelyett hogy a teljes API-t újra lekérné.

Az alábbi lépések bemutatják, hogyan tölthet be adatokat egy egyedi adatforrásból:

  1. Valósíts meg egy DataSource és egy DataSourceStreamReader elemet, amely meghívja az API-t, és nyomon követi az olvasási eltolást. Egyedi adatforrás létrehozásáról részletekért lásd: PySpark egyedi adatforrások.

  2. Regisztráld az adatforrást, hogy a csővezeték formátum alapján hivatkozhasson rá:

    spark.dataSource.register(MyApiDataSource)
    
  3. Olvasás a regisztrált forrásból egy adatfolyam-táblában:

    from pyspark import pipelines as dp
    
    @dp.table(name="events_bronze")
    def events_bronze():
        return spark.readStream.format("my_api_source").load()
    

3. minta: Az adatbetöltés leválasztása ütemezett feladattal és az Auto Loaderrel

Egy gyakori gyártási minta az, hogy az API hívást elválasztjuk a csővezetéktől. Egy ütemezett feladat fájlként adja le a nyers API válaszokat egy Unity Catalog kötetben, és a pipeline az Auto Loaderrel veszi fel azokat. Ez elkülöníti az API-specifikus sajátosságokat, például a lapozást és a sebességkorlátozást a deklaratív átalakítási logikától, és ráadásul az Auto Loader pontosan egyszeri fájlkövetését is biztosítja.

Az alábbi lépések megmutatják, hogyan lehet leválasztani a bevitést egy ütemezett munkától:

  1. Írj egy jegyzetfüzetet vagy szkriptet, amely az API-t hívja, és a nyers JSON válaszokat írja le egy Unity Catalog kötetre. Olvasd ki az API-hitelesítő adatokat egy titokból. Lásd: Titkos kódok kezelése.

    import requests, json, time
    
    token = dbutils.secrets.get(scope="<scope-name>", key="<secret-name>")
    volume_path = "/Volumes/main/raw/landing/api_events"
    
    resp = requests.get(
        "https://api.example.com/v1/events",
        headers={"Authorization": f"Bearer {token}"},
        timeout=30,
    )
    resp.raise_for_status()
    # One file per run; the pipeline's Auto Loader tracks which files it has ingested.
    with open(f"{volume_path}/events_{int(time.time())}.json", "w") as f:
        json.dump(resp.json()["data"], f)
    
  2. Ütemezd a jegyzetfüzetet vagy szkriptet, hogy önállóan futjon a Lakeflow Jobs-szal. Lásd Lakeflow Jobs.

  3. A pipeline-ban definiálj egy streaming táblát, amely az Auto Loaderrel olvassa a leérkezett fájlokat:

    from pyspark import pipelines as dp
    
    @dp.table(name="api_events_bronze")
    def api_events_bronze():
        return (
            spark.readStream.format("cloudFiles")
                .option("cloudFiles.format", "json")
                .load("/Volumes/main/raw/landing/api_events")
        )
    

További információkért a megbízható fájlfelvételről az Auto Loaderrel lásd: Fájlok betöltése felhőobjektum-tárolóból és Mi az Auto Loader?.

Legjobb gyakorlatok az API felvételéhez

  • Tartsd távol a titkokat a forráskódtól. Tárold az API tokeneket és kulcsokat Azure Databricks titkos scope-jaiban, és olvasd el őket futásidőben. Lásd: Titkos kódok kezelése.
  • A válaszokat korán ellenőrizd. Hozzáadj elvárásokat a felvett sorokon, hogy elkapd a hibás API-válaszokat, mielőtt azok lefelé áramlanának.
  • Kezeld az oldalszámolást és a sebességhatárokat. Az oldalak bejárása során adj hozzá újrapróbálkozást növekvő késleltetéssel, hogy egy átmeneti hiba ne hiúsítsa meg a teljes frissítést.

További források