このページでは

For AI agents: a documentation index is available at /docs/llms.txt. Append .md to any page URL for markdown, or send Accept: text/markdown.

Databricks

AmplitudeのDatabricksインポートソースを使用すると、DatabricksからAmplitudeアカウントにデータをインポートできます。 Databricksのインポートは、Databricksの変更データフィード(CDF)機能を使用して、Databricksワークスペースからライブデータに安全にアクセスして抽出します。

ガイド付きのセットアップ手順については、Loomのビデオをご覧ください。

機能

  • イベント、ユーザープロパティ、プロファイル、およびグループプロパティをインポートします。
  • デルタ同期をサポートすることで、Amplitudeは新しいデータや変更されたデータのみをインポートできます。

制限事項

  • ユーザー検索ページには、取り込まれた最新のイベントが 100 件表示されません。

  • 以下のDatabricks機能はタイムトラベルをサポートしておらず、この連携ではサポートされていません。

  • ミラー同期のデータ フィード タイプを変更するための SQL 入力制限事項:

    • ソースの Delta Table は 1 つだけです(「メイン テーブル」と呼ばれます)。
    • 単一のSELECTステートメント。
    • 共通表式 (CTE) (WITH句など) はサポートされていません。
    • UNION、INTERSECT 、MINUS 、EXCEPT などの設定操作はサポートされていません。
    • 句を含むステートメントは、メインテーブルのミューテーションメタデータを使用し、JOIN結合されたテーブルのミューテーション履歴を無視します。Amplitudeはデータ同期中に結合されたテーブル内のデータの最新バージョンを使用します。
    • 明示的なSQL検証はすべてのエッジケースをカバーするとは限りません。 たとえば、複数のソーステーブルを指定した場合、ソース作成時には検証は成功しますが、インポート実行時には失敗することがあります。
  • ユーザープライバシーAPI:ユーザープライバシーAPIは、以前に取り込まれたデータを削除し、Amplitudeがユーザーに関する新しい情報を処理することを妨げることはありません。 CDFを使用する場合、ユーザーに関するデータの送信を停止してからユーザープライバシーAPIを使用してユーザーを削除する必要があります。これにより、Amplitudeは次の同期時にユーザーを再作成しません。

    Amplitudeのシステムからエンドユーザーに関連付けられたすべてのデータを削除するには、データウェアハウスからユーザーを削除するだけでは不十分です。 このプロセスでは、Amplitudeがユーザーのデータをシステムから確実に削除できるようにするため、ユーザープライバシーAPIリクエストが必要です。

CDF とイベントボリューム:

CDF を使用することにより、Databricks は同期頻度に基づいて、統合された行INSERT、UPDATE挿入、およびDELETE削除操作を Amplitude に送信します。同期ウィンドウ中のイベントに対する複数の操作は、既存のイベントボリュームに対して1つのイベントとしてカウントされます。ただし、同期ウィンドウ外のイベントに対する操作は、既存のイベントボリュームに対する追加イベントとしてカウントされます。 このため、既存のイベントボリュームを使用する割合に影響を与える可能性があります。 必要に応じて、追加のイベントボリュームを購入するには、営業担当者に連絡してください。

  • ミラー同期イベントとミューテーションは不明なユーザーをサポートしていません。 行にはユーザーIDが含まれている必要があります。そうしないと、Amplitudeはイベントをドロップします。 大量の匿名イベントが発生している場合、Amplitudeはこのモードの使用を推奨しません。

Databricksの設定

AmplitudeでDatabricksソースの設定を開始する前に、Databricksで以下のタスクを完了してください。

汎用コンピューティング クラスタの検索または作成

Amplitudeは、同期ジョブを開始するために、ユーザーに代わってこのクラスターにワークフローを作成します。 完了したら、後のステップで使用するためにサーバのホスト名とHTTPパス値をコピーします。「構成」→「JDBC/ODBC」タブで両方の値を確認してください。 クラスタタイプの詳細については、「コンピューティング」を参照してください。

クラスタのポリシーに以下の設定を含まないことにより、新しいクラスタがジョブを実行できるようにしてください。詳細については、Databricksの記事「ポリシー定義」を参照してください。

json
"workload_type.clients.jobs": {
    "type": "fixed",
    "value": false
}

クラスターのPythonバージョンが3.9以上であることを確認してください。そうしないと、ワークフロー・ジョブで次のエラーが表示されることがあります。

json
TypeError: 'type' object is not subscriptable

クラスタポリシーとアクセスモード

Amplitudeはすべてのポリシーとアクセスモードをサポートしています。 ただし、クラスターに以下のポリシーとアクセスモードがある場合は、個人用アクセストークンがAmplitudeで認証されるワークスペースユーザーまたはサービスプリンシパルにデータリーダー権限(USE CATALOG、USE SCHEMA、EXECUTE、READ VOLUME、SELECT)を付与してください。そうしないと、インポートソースの Unity カタログ内のテーブルにアクセスできません。

AWS Databricks:

GCP Databricks:

認証

AmplitudeのDatabricksインポートは3つの認証方法をサポートしています。 お客様の組織のセキュリティ要件に合ったものを選択してください。

  • 個人用アクセストークン(PAT):Databricksワークスペースのユーザーまたはサービスプリンシパルとして認証します。最速でセットアップするにはワークスペースのユーザー認証を選択し、よりきめ細かな制御を行うにはサービスプリンシパルを選択してください。
  • OAuth 2.0 (サービスプリンシパル) (ベータ): クライアントIDとクライアントシークレットを使用して、Databricksが管理するOAuthサービスプリンシパルで認証を行います (OAuthマシンツーマシン、またはM2M)。組織でパーソナル・アクセストークンを許可していない場合は、この方法を選択してください。
  • Microsoft Entra IDを使用したOAuth 2.0 (ベータ): Azure Databricksの場合、クライアントID、クライアントシークレット、およびテナントIDを使用してMicrosoft Entra IDサービスプリンシパルで認証を行います。

詳細については、「Databricks Automation のための認証」を参照してください。

どちらの方法を選択しても、サービス プリンシパルまたはユーザーには以下に示す権限が必要です。

ワークスペースユーザーのパーソナル・アクセストークン(PAT)を作成する

AmplitudeのDatabricksインポートは、認証にパーソナル・アクセストークンを使用します。 最も迅速なセットアップを行うには、Databricksでワークスペースユーザー用にPATを作成します。 詳細については、Databricksの記事「ワークスペースユーザーのためのパーソナル・アクセス・トークン」を参照してください。

サービスプリンシパルのパーソナルアクセストークン(PAT)を作成する

Amplitudeは、Databricksでサービス・プリンシパルを作成して、より細かいアクセス制御を可能にすることをお勧めします。

  1. Databricksの指示に従ってサービス・プリンシパルを作成してください。 後の手順で使用するためにUUIDをコピーします。
  2. このサービス・プリンシパルに対してPATを生成します。
ベータ

Databricksのインポート用のOAuth 2.0認証方法はベータです。一般提供前に動作が変更される可能性があります。

OAuth 2.0(サービス・プリンシパル)認証の設定

組織がパーソナル アクセストークンを許可していない場合は、OAuth マシンツーマシン(M2M)認証を使用してください。この方法は、PATではなくDatabricksが管理するOAuthシークレットを使用して認証を行います。

  1. Databricksでサービス プリンシパルを作成します。
  2. サービス プリンシパルのOAuthシークレットを生成します。Databricksの記事「OAuth(OAuth M2M)を使用したサービス プリンシパルによるDatabricksへのアクセス認証」に従ってください。クライアントID(サービス プリンシパルのアプリケーションID)とクライアント シークレットをコピーして、後の手順で使用します。

Microsoft Entra ID 認証を使用した OAuth 2.0 の設定

Azure Databricks の場合、Microsoft Entra ID サービス プリンシパルを使用して認証を行います。

  1. Microsoft Entra ID サービス プリンシパルを作成し、それを Azure Databricks ワークスペースに追加します。マイクロソフトの記事「Microsoft Entra サービス プリンシパルによる認証」に従ってください。
  2. クライアント ID (Entra アプリケーションのクライアント ID)、クライアント シークレット、および Microsoft Entra ID テナント ID をコピーして、後の手順で使用します。

権限

認証に使用するサービス・プリンシパルまたはユーザーには、Databricksで次の権限が必要です。

  • ワークスペース:
    • 理由: Databricksワークスペースへのアクセスを許可します。
    • Databricks内の場所:
      • ワークスペース → → 権限 → 権限の追加。
      • ユーザー権限で作成したサービスプリンシパルを追加し、[保存] をクリックします。
  • テーブル:
    • 理由: リスト テーブルへのアクセスを許可し、データを読み取ります。
    • Databricks内の場所:
      • 「カタログ」→「カタログ」を選択→「権限」→「付与」。
      • Data Reader権限(USE CATALOG、USE SCHEMA、EXECUTE、READ VOLUME、SELECT)を選択します。
  • クラスタ:
    • 理由: クラスタに接続し、ユーザーに代わってワークフローを実行するためのアクセスを許可します。
    • Databricks内の場所:
      • コンピューティング → 汎用コンピューティング → 編集権限。
      • サービス プリンシパルにCan Restart権限を追加します。
  • エクスポート:
    • 理由: サービスプリンシパルが spark を通じてデータをアンロードし、S3 にエクスポートできるようにします。
    • Databricksでの場所:任意のノートブックで以下のSQLコマンドを実行します。 GRANT MODIFY ON ANY FILE TO ``; GRANT SELECT ON ANY FILE TO ``;

テーブルで CDF を有効にする

AmplitudeはDatabricksのChange Data Feedを使用してデータを継続的にインポートしています。 DatabricksテーブルでCDFを有効にするには、「Databricks | 変更データフィードを有効にする」を参照してください。

Amplitude Databricksソースを設定する

DatabricksをAmplitudeのソースとして追加するには、以下の手順を実行してください。

Databricksに接続する

  1. Amplitudeデータで、[カタログ] -> [ソース] に移動します。
  2. Databricksを検索します。
  3. 「Connect Databricks」画面の「Credentials」タブで、Databricksの設定時に設定した認証情報を入力します。
    • サーバのホスト名。
    • HTTPパス。
    • 認証方法を選択し、対応する認証情報を入力してください:
      • パーソナルアクセストークン:ワークスペースユーザーまたはサービスプリンシパルのPAT。
      • OAuth 2.0(サービスプリンシパル)(ベータ):OAuth クライアント ID とクライアント シークレットです。
      • OAuth 2.0 — Microsoft Entra ID (Azure)(ベータ):OAuth クライアント ID、クライアントシークレット、および Microsoft Entra ID テナント ID。
  4. [Next] をクリックしてアクセスを確認します。

インポートするデータを選択

  1. インポートするデータタイプを選択します。 Databricksソースは3つのデータタイプをサポートしています。

  2. データタイプとして選択した場合は、インポート戦略を選択します。Event

  • Append Only Sync:ID解決、プロパティとアトリビューションの同期、位置情報の解決など、Amplitudeの標準的なエンリッチメントサービスを使用してデータウェアハウスデータを取り込みます。
  • ミラー同期:挿入、更新、削除操作を使用してDatabricks内のデータを直接ミラーリングします。これにより、Amplitudeのエンリッチメントサービスが無効になり、真実のソースと同期を保ちます。
  1. Amplitudeがデータをインポートする前に、Databricksでデータを変換するSQLコマンドを設定します。

    • AmplitudeはSQL実行出力内の各レコードをインポートするイベントとして扱います。 インポートする各レコードが準拠していることを確認するには、バッチイベントアップロード API ドキュメントの Example ボディを参照してください。
    • Amplitudeは上記の手順1で指定したテーブルからのみ変換/インポートできます。
      • たとえば、テーブル A へのアクセス権がありながら手順 1 Bでのみ選択した場合、C からAのみデータをインポートできますA。
    • SQL コマンドで参照するテーブル名は、ステップ 1 で選択したテーブルの名前と一致している必要があります。たとえば、catalog.schema.table1 を選択した場合、SQL でその正確な値を使用します。
    sql
    select
        unix_millis(current_timestamp())                                       as time,
        id                                                                     as user_id,
        "demo"                                                                 as event_type,
        named_struct('name', name, 'home', home, 'age', age, 'income', income) as user_properties,
        named_struct('group_type1', ARRAY("group_A", "group_B"))               as groups,
        named_struct('group_property', "group_property_value")                 as group_properties
    from catalog.schema.table1;
    

VARIANTデータ型サポート

Amplitudeは、プロパティ列をマッピングするためのVARIANTデータ型をサポートしています。SQLクエリーでは、ユーザープロパティ、イベントプロパティ、およびグループプロパティのVARIANT列を使用できます。これにより、半構造化またはネストされた JSON データを直接 Databricks からインポートできます。

Eventデータタイプと追加専用インジェストについては、オプションでユーザープロパティの_同期_または_グループプロパティ_の同期を選択して、イベント内の対応するプロパティを同期します。

  1. SQL を追加したら、「SQL のテスト」をクリックします。 AmplitudeはDatabricksインスタンスに対してテストを実行し、SQLが有効であることを確認します。[次へ] をクリックします。

  2. 初期インポート用のテーブルのバージョンを選択します。最初のインポートでは、選択したバージョンの時点でテーブルからすべての情報が取り込まれます。 [最初] または [最新] を選択します。

    • First は最初のバージョンを意味します。これは 0 です。
    • Latest 最新バージョンを意味します。
  3. 同期頻度を設定します。 ソースを設定するときに同期頻度を設定できます。 この頻度によって、AmplitudeがDatabricksからデータを取得する間隔が決まります。

    使用できる同期頻度オプションは、インポートするデータの種類(イベント、ユーザープロパティ、グループプロパティ、プロファイルなど)によって異なります。 例:

    • 毎日の同期: 毎日1回、指定された時間に実行されます。
    • 時間ごとの同期: 1時間ごとに実行されます。

    Amplitudeはベストエフォートベースで同期を実行します。 同期は通常、設定された頻度で実行されますが、それほど頻繁に実行されない場合もあります。

  4. このソースのインスタンスを説明できる名前を入力します。

  5. ソースは、ワークスペースの [ソース] リストに表示されます。

データインポートを確認する

Amplitudeがインポートするイベントは、SQLステートメントで割り当てた名前を前提とします。上記の例では、イベント名は demo です。

Amplitudeに送られるデータを検証するには:

  • トラッキングプランの「イベント」ページを表示します。
  • 指定したイベント名に基づいてフィルタリングを行うセグメンテーションチャートを作成します。
  • ソース内のIngestion Jobsタブに移動します。 必要に応じて、ERROR LOG を使用して、取り込みとデバッグのステータスを表示できます。

会社のネットワークポリシーによっては、AmplitudeのサーバーがDatabricksインスタンスにアクセスできるようにするには、以下のIPアドレスを許可リストに追加する必要がある場合があります。

  • Amplitude米国のIPアドレス:
    • 52.33.3.219
    • 35.162.216.242.
    • 52.27.10.221.
  • Amplitude EU の IP アドレス:
    • 3.124.22.25.
    • 18.157.59.125.
    • 18.192.47.195.

トラブルシューティング

text
shaded.databricks.org.apache.hadoop.fs.s3a.AWSClientIOException: getFileStatus on s3a://com-amplitude-falcon/databricks_import/unloaded_data/source_destination_158631/batch_712300169/meta: com.amazonaws.SdkClientException: Unable to execute HTTP request: Remote host terminated the handshake: Unable to execute HTTP request: Remote host terminated the handshake
---------------------------------------------------------------------------
 Py4JJavaError                             Traceback (most recent call last)
 File /databricks/spark/python/pyspark/errors/exceptions.py:228, in capture_sql_exception.<locals>.deco(*a, **kw)
 227 try:
 --> 228     return f(*a, **kw)
 229 except Py4JJavaError as e:
  • 根本原因 1: 共有アクセスモードで実行されているクラスターの制限が原因で、セキュリティ強化対策が追加されているため、このエラーが発生しています。
  • 解決策 1: シングルユーザーまたは分離なしの共有アクセスモードで構成されたクラスターを使用します。
  • 根本原因 2: Azure Databricks の場合、ファイアウォール設定により、Databricks クラスターが Amplitude s3 バケットにアクセスできないことがあります。
  • 解決策 2:DatabricksクラスターがAmplitude s3バケットにアクセスすることをブロックしているネットワークルールがないか確認してください。
  1. plaintext
    [Databricks][JDBCDriver](500593) Communication link failure. Failed to connect to server. Reason: HTTP Response code: 403, Error message: PERMISSION_DENIED: You do not have permission to autostart 0108-111840-mc2khhh6.. isCausedByCustomer=true,isAutomaticallyRecoverable=false,errorType=databricks-jdbc-connection-error.
    
plaintext

  - **根本原因**: このエラーは、Amplitudeが汎用クラスターとJDBC接続を確立しようとしたときに発生します。クラスタが停止している場合、Databricks は自動的にクラスタを再起動しようとします。このため、PAT にクラスタを起動するための十分な権限がない場合、権限エラーが発生します。
  - **解決策**: 「Can Restart」権限を持つ[サービス プリンシパル アクセストークン](#create-a-service-principal-personal-access-token-pat)を作成し、クラスター上のユーザーに割り当ててアクセスを許可します。

3. ```
  java.sql.SQLException: [Databricks][JDBCDriver](500593) Communication link failure. Failed to connect to server. Reason: HTTP Response code: 502, Error message: Unknown.
  • 根本原因: ドライバまたはエグゼクタのいずれかがメモリ使用量が多いために応答していない可能性があります。 メモリがいっぱいになると、ドライバ/エグゼクタはメモリを解放するためのガベージ・コレクション中に数秒間応答しなくなります。 このため、上記のように通信リンク障害が発生することがあります。
  • 解決策: 1) ワークロードと基盤となるテーブルを最適化してメモリ使用量を減らすか、2) ドライバのメモリを増やすか、3) ワーカーの最大数を増やして、リソース競合が激しいときにクラスタがより多くのワーカーを追加できるようにします。
  1. plaintext
    Caused by: java.sql.SQLException: [Databricks][JDBCDriver](500593) Communication link failure. Failed to connect to server. Reason: HTTP Response code: 554, Error message: Service is under maintenance..
    
plaintext

  - **根本原因**: Databricks プロキシ API の問題です。
  - **解決策**: Databricks サポートにエスカレーションしてください。

5. ```
  Py4JJavaError: An error occurred while calling o450.isEmpty.
  : com.databricks.sql.transaction.tahoe.DeltaFileNotFoundException: [DELTA_EMPTY_DIRECTORY] No file found in the directory: s3
  • 根本原因:デルタログは切り詰められ、期限切れになっているため、Amplitudeサービスはインポートするデルタログファイルを見つけられませんでした。
  • 解決策:Amplitude サポートに連絡して、期限切れのデータをスキップし、インポートを続行してください。デフォルトでは、リテンションは30日です。データを少なくとも7日間保持するようにしてください。
  1. plaintext
    Fail worker job since databricks job run finished with state MAXIMUM_CONCURRENT_RUNS_REACHED.
    
plaintext

  - **根本原因**:Databricksジョブでは一度に1つのみ実行が許可されています。前の実行がまだ進行中の状態で新しいジョブをトリガーした場合、Databricksはそれをスキップします。
  - **解決策**:Databricksジョブ構成を更新して、[同時実行の最大数を増や](https://docs.databricks.com/aws/en/jobs/configure-job#configure-maximum-concurrent-runs)します。

7. ```
  [Databricks][JDBCDriver](700100) Connection timeout expired. Details: None.
  • 根本原因:これは、AmplitudeがDatabricksエンドポイントへのJDBC接続を確立できなかったことを意味します。 これは次のような理由でよく発生します:
    • Databricksクラスタが停止しており、再起動に時間がかかりすぎているため、接続試行がタイムアウトになりました。
    • Databricksワークスペースがリソース制限に達している可能性があります(たとえば、同時実行SQLエンドポイントの最大数、クラスタのクォータ数など)。
  • ソリューション:
    • Databricksのワークスペースとクラスタのステータスをチェックして、接続試行中にクラスタが停止したか、または再起動しているかを確認してください。
    • クラスタの自動起動と自動終了の設定を確認して、必要に応じてクラスタが迅速に再起動できるようにしてください。
    • Databricks のリソースをモニターし(同時接続の上限やクラスタ容量の問題など)、必要に応じてクォータを調整します。
  1. plaintext
    [DELTA_MISSING_CHANGE_DATA] Error getting change data for range [2 , 3] as change data was not recorded for version [2]
    
plaintext

  - **根本原因**:これは、Amplitudeが特定のバージョン範囲のデータをテーブルから取得できなかったことを意味します。 これが発生する理由は次のとおりです。
    - 特定のテーブルバージョンの後に変更データフィード (CDF) を有効にしたため、その範囲には変更データが存在しません。
    - 特定のテーブルバージョンはバキュームされ、Amplitudeは対応するデータファイルを削除しました。
  - **ソリューション**:
    - 新しいソースを作成し、最新のテーブルバージョンからのインポートを開始します。
    - 同じソースを再利用し、テーブルバージョンをスキップしたい場合は、[Amplitude サポート](https://gethelp.amplitude.com)までお問い合わせください。

これは役に立ちましたか?