Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
Ingestování z API znamená stahování dat přes HTTP z webové služby, obvykle jako stránkovaný JSON, místo čtení ze souboru nebo databáze. Na rozdíl od souborů nebo sběrnice zpráv zde není vestavěný generický API zdroj, takže autentizaci, stránkování a omezení rychlosti řešíte sami. Kanály Lakeflow podporují tři způsoby načítání z libovolného API. To, který z nich je vhodný, závisí na vašich požadavcích na objem a frekvenci aktualizace.
Důležité
Než napíšete jakýkoli vlastní kód pro API ingester, zkontrolujte, zda už pro váš zdrojový kód neexistuje spravovaný konektor. Lakeflow Connect nabízí vestavěné konektory pro mnoho běžných API softwaru jako služby (SaaS), jako jsou Salesforce, Workday, ServiceNow a Google Analytics, a také roste počet partnerských konektorů. Pokud konektor pokrývá váš zdroj, postará se o autentizaci, stránkování a inkrementální extrakci za vás, a téměř vždy je to méně práce než ručně svíjený ingester. Viz koncepty konektorů Lakeflow Connect. Použijte níže uvedené vzory jen tehdy, když žádný konektor nepasuje.
Předpoklady
- Potrubí. Pro vytvoření jednoho si přečtěte návody na potrubí Lakeflow.
- API přihlašovací údaje, jako je token nebo klíč, uložené jako Azure Databricks secret. Nikdy nenatvrdy nekódujte přihlašovací údaje do zdrojového kódu pipeline. Viz správa tajemství.
- Přístup k síti z vašeho pipeline výpočetního systému do API endpointu.
- Znalost streamovacích tabulek a materializovaných pohledů, typy datových sad, které tyto vzory vytvářejí. Viz Streamovací tabulky a Materializované zobrazení.
Zvolte vzor
V pipelinech není nativní generický REST-API zdroj, takže když vybíráte z libovolného API, vyberte jeden ze tří vzorů podle objemu dat a frekvence ingestu:
| Vzor | Použít, když |
|---|---|
| Periodické tahy jako materializovaný pohled | Datové objemy jsou malé až střední a načítají se jednou během každého spuštění pipeline, například referenční data, denní devizové kurzy nebo API se stránkováním, jehož rozsah lze omezit. |
| Python Data Source API | API s vysokým objemem dat nebo streamovací API musíte dotazovat inkrementálně a ukládat kontrolní body postupu, aby se po restartu nemuselo vše načítat znovu. |
| Oddělená ingesce dat pomocí Auto Loader | Chcete izolovat specifické specifika API od transformační logiky a získat přesně jednorázové sledování souborů zdarma. |
Vzor 1: Periodické tahy jako materializovaný pohled
Pro malé až střední payloady načtené jednou za běh pipeline napište Python funkci, která volá API a vrací Spark DataFrame. Protože datový soubor je materializovaný pohled, pipeline při každé aktualizaci znovu spustí funkci v plném rozsahu a idempotentně.
Následující kroky vám ukážou, jak vytvořit materializovaný pohled pomocí periodických tahů:
Uložte API token do tajného kódu a pak ho namapujte na konfigurační vlastnost Spark v nastavení pipeline, aby ho pipeline kód mohl číst. Přidejte vlastnost do bloku
spark_confkonfigurace clusteru kanálu:{ "clusters": [ { "spark_conf": { "api.token": "{{secrets/<scope-name>/<secret-name>}}" } } ] }Kód v dalším kroku přečte tuto hodnotu s
spark.conf.get("api.token"). Další informace o konfiguraci tajných klíčů v nastavení kanálu najdete v tématu Jak bezpečně přistupovat k přihlašovacím údajům k úložišti pomocí tajných klíčů v kanálu.Definujte materializovaný pohled, který volá API a vrací odpověď jako DataFrame:
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)Zpracujte stránkování uvnitř funkce tak, že budete v cyklu procházet jednotlivé stránky a před vrácením objektu DataFrame zřetězíte výsledky:
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)Přidejte logiku opakování a ustupování kolem požadavku na odolnost.
Tento vzor znovu čte celou odpověď API při každé aktualizaci pipeline, proto jej používejte pouze tehdy, když je payload omezen. Pro inkrementální čtení použijte vzor 2.
Vzor 2: API pro vysoké objemy nebo streamování s Python Data Source API
Pro API, u nichž je třeba provádět inkrementální dotazování se sledováním offsetu, implementujte vlastní zdroj dat pomocí rozhraní Spark Python Data Source API. To vám poskytne správnou sémantiku streamování, včetně postupu ukládaného pomocí kontrolních bodů a inkrementálních čtení, takže se po restartu pokračuje od posledního offsetu místo nového načítání celého API.
Následující kroky vám ukážou, jak importovat z vlastního datového zdroje:
Implementujte komponentu
DataSourceaDataSourceStreamReader, které volají API a sledují offset čtení. Podrobnosti o tvorbě vlastního datového zdroje najdete v PySpark custom data sources.Zaregistrujte zdroj dat, aby ho pipeline mohl odkazovat podle názvu formátu:
spark.dataSource.register(MyApiDataSource)Přečtěte si z registrovaného zdroje v tabulkě streamování:
from pyspark import pipelines as dp @dp.table(name="events_bronze") def events_bronze(): return spark.readStream.format("my_api_source").load()
Vzor 3: Oddělit příjem od plánované práce a automatického nakladače
Běžným výrobním vzorem je oddělení volání API od pipeline. Plánovaná úloha přiloží surové API odpovědi jako soubory do svazku Unity Catalog a pipeline je vyzvedne pomocí Auto Loaderu. To izoluje specifické API zvláštnosti, jako je stránkování a omezení rychlosti, od deklarativní transformační logiky a poskytuje vám bezplatné sledování souborů přesně jednou v AutoLoaderu.
Následující kroky vám ukážou, jak oddělit požití od plánované práce:
Napiš zápisník nebo skript, který volá API a zapisuje surové JSON odpovědi do svazku Unity Catalog. Přečti API přihlašovací údaje z tajného systému. Viz správa tajemství.
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)Naplánujte zápisník nebo skript tak, aby běžel samostatně s Lakeflow Jobs. Podívejte se na Úlohy Lakeflow.
Ve svém pipeline definujte streamovací tabulku, která čte přistávané soubory pomocí Auto Loaderu:
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") )
Pro více informací o spolehlivém zasílání souborů pomocí Auto Loaderu viz Načítání souborů z cloudového objektového úložiště a Co je Auto Loader?.
Nejlepší postupy pro vstup API
- Držte tajemství mimo zdrojový kód. Ukládejte API tokeny a klíče v Azure Databricks secret scopes a čtěte je za běhu. Viz správa tajemství.
- Ověřujte odpovědi včas. Přidejte očekávání na ingestované řádky, abyste zachytili deformované odpovědi API dříve, než budou postupovat dál.
- Zvládejte stránkování a limity rychlosti. Projít stránky a přidat opakování s backoffem, aby kvůli přechodnému selhání neselhala celá aktualizace.