パイプライン内のAPIからデータを取り込む

APIからの取り込みとは、ファイルやデータベースから読み取るのではなく、通常はページ付きJSONとしてウェブサービスからHTTP経由でデータを引き出すことを意味します。 ファイルやメッセージバスとは異なり、汎用APIソースは内蔵されていないため、認証、ページ化、レート制限は自分で管理できます。 Lakeflowパイプラインは任意のAPIからの取り込みに3つのパターンをサポートしています。 どのモデルが合うかは音量やリフレッシュの必要度によります。

Important

カスタムAPIインジェスションコードを書く前に、ソース用のマネージドコネクターがすでに存在しているか確認してください。 Lakeflow Connectは、Salesforce、Workday、ServiceNow、Google Analyticsなど、多くの一般的なSaaS(ソフトウェア・アズ・ア・サービス)API向けの組み込みコネクタを出荷しており、パートナーコネクタも増加しています。 コネクターがソースをカバーすれば、認証、ページ化、インクリメンタル抽出を代行し、手作業での取り込みよりもほとんどの場合手間が少なくなります。 Lakeflow Connect のマネージド コネクタを参照してください。 コネクターが合わない時のみ以下のパターンを使用してください。

前提条件

  • パイプラインだ。 作成方法については、 Lakeflowパイプラインのチュートリアルをご覧ください。
  • トークンやキーなどのAPI認証情報は、Azure Databricksの秘密として保存されます。 パイプラインのソースコードに認証情報をハードコーディングしないでください。 「シークレットの管理」を参照してください。
  • パイプライン コンピューティング環境から API エンドポイントへのネットワーク アクセス。
  • ストリーミングテーブルやマテリアル化されたビュー、これらのパターンが生み出すデータセットタイプへの精通。 ストリーミング テーブルマテリアライズドビューを参照してください。

パターンを選択する

パイプラインにはネイティブの汎用 REST-API ソースは存在しないため、任意のAPIから取得する際は、データ量と取り込む頻度に基づいて3つのパターンのいずれかを選びます。

Pattern 次の場合に使用します。
マテリアライズド ビューとしての定期的なプル ペイロードは小規模から中規模で、参照データ、日次FXレート、またはページ付きだがバウンド可能なAPIなど、パイプライン実行ごとに一度だけプルされます。
Python データソース API 大量データまたはストリーミング API を増分的にポーリングし、再起動してもすべてを読み直さないように、進捗をチェックポイントとして保存する必要があります。
Auto Loaderによる分離インジェスト API固有の癖を変換ロジックから切り離し、正確に一度だけのファイル追跡を無料で得たいです。

パターン 1: マテリアライズド ビューとしての定期的なプル

パイプラインごとに1回だけ引き取る小規模から中規模のペイロードについては、APIを呼び出してSpark DataFrameを返すPython関数を書きます。 データセットがマテリアライズされたビューであるため、パイプラインは更新するたびに関数を完全にかつ冪等的に再実行します。

以下の手順で、定期的なプルを使った物質化されたビューの構築方法を示します。

  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レスポンス全体を再読み込むため、ペイロードが境界化されている場合にのみ使用してください。 インクリメンタルリードにはパターン2を使いましょう。

パターン 2: Python Data Source API を使用した高ボリュームまたはストリーミング API

APIの場合はオフセットトラッキングでインクリメントポーリングが必要で、SparkのPython Data Source APIを使ってカスタムデータソースを実装してください。 これにより、チェックポイントの進捗やインクリメンタルリードを含む適切なストリーミングセマンティクスが得られ、再起動はAPI全体を再度引き出すのではなく、最後のオフセットから再開されます。

以下の手順で、カスタムデータソースから取り込む方法を示します:

  1. APIを呼び出して読み取りオフセットを追跡する DataSourceDataSourceStreamReader を実装します。 カスタムデータソースの作成方法については、 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()
    

パターン 3: スケジュールされたジョブと Auto Loader を使用して、データ取り込みを分離する

一般的な本番パターンとして、API呼び出しをパイプラインから分離することがあります。 スケジュールジョブは生のAPI応答をUnityカタログのボリュームにファイルとして格納し、パイプラインがAuto Loaderでそれらを拾います。 これにより、ページ番号やレート制限などのAPI固有の癖を宣言型変換ロジックから切り離し、Auto Loaderの正確に一度だけのファイル追跡を無料で提供できます。

以下の手順で、スケジュールされた作業と摂取を切り離す方法を示します:

  1. APIを呼び出し、Unityカタログのボリュームに生のJSON応答を書くノートブックやスクリプトを書きます。 シークレットから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応答が下流する前に検出します。
  • ページ設定やレート制限を管理しましょう。 ページをループオーバーし、バックオフ付きのリトライを追加して、一時的な失敗がアップデート全体に失敗しないようにしましょう。

その他のリソース