從 API 匯入是指從網路服務透過 HTTP 拉取資料,通常是分頁 JSON,而非從檔案或資料庫讀取資料。 不像檔案或訊息匯流排,這裡沒有內建通用的 API 原始碼,所以認證、分頁和速率限制都是自己處理的。 Lakeflow 管線支援三種從任意 API 匯入的模式。 哪一個適合,取決於你的音量和刷新需求。
這很重要
在撰寫任何自訂 API 導入程式碼前,先確認你的原始碼是否已經有受管理連接器。 Lakeflow Connect 內建多種常見軟體即服務(SaaS)API 連接器,如 Salesforce、Workday、ServiceNow 和 Google Analytics,且合作夥伴連接器數量也日益增加。 如果連接器能覆蓋你的來源,它會幫你處理認證、分頁和增量擷取,而且幾乎總是比手動捲製的讀取工作量少。 詳見 Lakeflow Connect 連接器概念。 僅在沒有合適的接頭時,才使用以下樣式。
先決條件
- 一條管線。 要建立一個,請參考 Lakeflow pipelines 的教學。
- API 憑證,例如令牌或金鑰,則儲存為 Azure Databricks 秘密。 絕不要在管線原始碼中硬編碼憑證。 請參閱機密管理。
- 從你的管線運算到API端點的網路存取。
- 熟悉串流表格與具體化視圖,這些模式所產生的資料集類型。 請參見 串流表 與 實體化視圖。
選擇一個圖案
管線中沒有原生通用 REST-API 來源,所以當你從任意 API 拉取資料時,根據資料量和你匯入的頻率,選擇三種模式之一:
| 樣式 | 何時使用 |
|---|---|
| 週期性拉取作為具體視角 | 資料負載屬於小到中等規模,且每次管線執行時只會擷取一次,例如參考資料、每日外匯匯率,或雖為分頁但可限定範圍的 API。 |
| Python 資料來源 API | 你需要以增量方式輪詢高吞吐量或串流 API,並以檢查點記錄進度,避免重啟時從頭重讀所有資料。 |
| 使用 AutoLoader 的分離擷取 | 你要把 API 特定的怪異現象從轉換邏輯中隔離出來,免費獲得 exact-once 檔案追蹤。 |
模式一:定期擷取作為實體化檢視
對於每次管線執行時只擷取一次的中小型資料量,請撰寫一個 Python 函式,用來呼叫 API 並回傳 Spark DataFrame。 由於資料集是實體化檢視表,管線每次更新時都會以完整且冪等的方式重新執行該函式。
下列步驟說明如何建立具體化檢視,並定期擷取資料:
將 API 權杖存入秘密,然後映射到 Spark 設定中的屬性,讓管線程式碼能讀取。 將該屬性加入管線叢集配置的
spark_conf區塊中:{ "clusters": [ { "spark_conf": { "api.token": "{{secrets/<scope-name>/<secret-name>}}" } } ] }下一步的程式碼會使用
spark.conf.get("api.token")讀取此值。 欲了解更多關於在管線設定中設定秘密的資訊,請參閱 「安全存取帶有管線內秘密的儲存憑證」。定義一個具體化檢視,呼叫 API 並以 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)在函式內處理分頁,方法是循環頁面並將結果串接起來,然後再回傳 DataFrame:
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)為了提升韌性,請在該請求周圍加入重試與退避機制。
此模式會在每次管線更新時重新讀取完整的 API 回應,因此只有在有效載荷有限制時才使用。 對於增量閱讀,請使用模式二。
模式二:使用Python資料來源API進行高容量或串流API。
對於需要增量輪詢並使用偏移追蹤的 API,請使用 Spark 的 Python Data Source API 實作自訂資料來源。 這可提供正確的串流語義,包括檢查點進度資訊和增量讀取,因此重新啟動時會從上次的偏移量繼續,而不是再次從整個 API 拉取資料。
以下步驟示範如何從自訂資料來源擷取:
實作一個
DataSource和DataSourceStreamReader,以呼叫 API 並追蹤讀取偏移量。 關於撰寫自訂資料來源的詳細資訊,請參見 PySpark 自訂資料來源。註冊資料來源,讓管線能以格式名稱參考:
spark.dataSource.register(MyApiDataSource)從串流資料表中的已註冊來源讀取:
from pyspark import pipelines as dp @dp.table(name="events_bronze") def events_bronze(): return spark.readStream.format("my_api_source").load()
模式三:將資料擷取與排程工作及自動載入器分離
一個常見的生產模式是將 API 呼叫與管線分離。 排程工作會將原始 API 回應以檔案形式落在 Unity 目錄卷中,管線會用 Auto Loader 接收這些回應。 這樣可將 API 特有的細節(例如分頁和速率限制)與您的宣告式轉換邏輯隔離開來,並讓您直接享有 Auto Loader 的僅一次檔案追蹤功能。
以下步驟教你如何將攝取與排定工作分離:
撰寫一個筆記本或腳本,呼叫 API,並將原始 JSON 回應寫入到 Unity Catalog 磁碟區。 從機密中讀取 API 憑證。 請參閱機密管理。
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)透過 Lakeflow Jobs 排程筆記本或指令碼,使其自動執行。 請參閱 Lakeflow 職位。
在你的管線中,定義一個使用 Auto Loader 讀取已落地檔案的串流資料表:
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") )
想了解更多關於使用 Auto Loader 可靠檔案擷取的資訊,請參見 「從雲端物件儲存載入檔案 」和 「什麼是 Auto Loader?」。
API 資料擷取的最佳實務
- 將秘密排除在原始碼之外。 將 API 權杖和金鑰儲存在 Azure Databricks 的秘密範圍中,並在執行時讀取它們。 請參閱機密管理。
- 及早驗證回應。 在擷取的資料列上加入 期望條件,以便在格式錯誤的 API 回應流向下游之前將其攔截。
- 處理分頁和速率限制。 逐頁處理,並加入退避重試機制,以免暫時性故障導致整個更新失敗。