從管線中的 API 擷取資料

從 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。 由於資料集是實體化檢視表,管線每次更新時都會以完整且冪等的方式重新執行該函式。

下列步驟說明如何建立具體化檢視,並定期擷取資料:

  1. 將 API 權杖存入秘密,然後映射到 Spark 設定中的屬性,讓管線程式碼能讀取。 將該屬性加入管線叢集配置的 spark_conf 區塊中:

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

    下一步的程式碼會使用 spark.conf.get("api.token") 讀取此值。 欲了解更多關於在管線設定中設定秘密的資訊,請參閱 「安全存取帶有管線內秘密的儲存憑證」。

  2. 定義一個具體化檢視,呼叫 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)
    
  3. 在函式內處理分頁,方法是循環頁面並將結果串接起來,然後再回傳 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 拉取資料。

以下步驟示範如何從自訂資料來源擷取:

  1. 實作一個 DataSource 和 DataSourceStreamReader,以呼叫 API 並追蹤讀取偏移量。 關於撰寫自訂資料來源的詳細資訊,請參見 PySpark 自訂資料來源。

  2. 註冊資料來源,讓管線能以格式名稱參考:

    spark.dataSource.register(MyApiDataSource)
    
  3. 從串流資料表中的已註冊來源讀取:

    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 的僅一次檔案追蹤功能。

以下步驟教你如何將攝取與排定工作分離:

  1. 撰寫一個筆記本或腳本,呼叫 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)
    
  2. 透過 Lakeflow Jobs 排程筆記本或指令碼,使其自動執行。 請參閱 Lakeflow 職位。

  3. 在你的管線中,定義一個使用 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 回應流向下游之前將其攔截。
  • 處理分頁和速率限制。 逐頁處理,並加入退避重試機制,以免暫時性故障導致整個更新失敗。

其他資源