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.
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
- Egy csővezeték. A létrehozáshoz lásd a Lakeflow pipelines oktatóanyagokat.
- API hitelesítő adatok, például token vagy kulcs, Azure Databricks titkos értékként tárolva. Sose kódolj hitelesítést a pipeline forráskódban. Lásd: Titkos kódok kezelése.
- Hálózati hozzáférés a pipeline számítástól az API végpontig.
- Ismerős a streaming táblák és a materializált nézetek iránt, az adathalmaztípusok, amelyeket ezek a minták hoznak. Lásd a Streaming táblázatokat és a Materializált nézeteket.
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:
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_confcső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.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)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:
Valósíts meg egy
DataSourceés egyDataSourceStreamReaderelemet, 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.Regisztráld az adatforrást, hogy a csővezeték formátum alapján hivatkozhasson rá:
spark.dataSource.register(MyApiDataSource)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:
Í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)Ütemezd a jegyzetfüzetet vagy szkriptet, hogy önállóan futjon a Lakeflow Jobs-szal. Lásd Lakeflow Jobs.
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.