Lakeflow Connect を使用して Azure Databricks に取り込むためのデータ ソースとして Microsoft Dynamics 365 を設定する方法について説明します。
注
このページでは、Azure Synapse Analytics workspaceを使わないCSVエクスポートワークフローについて説明しています。 Azure Synapse Analyticsワークスペースを使ってParquet形式のDeltaテーブルをエクスポートするには、「ParquetデータソースをMicrosoft Dynamics 365インジェスションのために設定する」を参照してください。 Databricksは、大規模または大ボリュームのインスタンスに対してParquetのワークフローを推奨しており、その理由は大規模でのパフォーマンスと安定性が向上するためです。
コネクタがソース データにアクセスする方法については、「コネクタ が D365 データにアクセスする方法」を参照してください。 サポートされている Dataverse アプリケーションの一覧については、「 サポートされている Dynamics 365 アプリケーション」を参照してください。
[前提条件]
Dynamics 365 データ ソースを構成する前に、次のものが必要です。
- リソースを作成するアクセス許可を持つアクティブな Azure サブスクリプション。
- 管理者アクセス権を持つ Microsoft Dynamics 365 環境。
- Dynamics 365 インスタンスに関連付けられている Dataverse 環境。
- Azure Databricks のワークスペース管理者またはメタストア管理者のアクセス許可。
- Dataverse 環境で Azure Synapse Link を作成および構成するためのアクセス許可。
- Azureのサブスクリプションで、ストレージアカウントが他のSynapse Linkプロファイルに紐づいていない場合。 Dataverseのテーブルを別のプロファイルに紐づけたストレージアカウントに追加することはできません。新しいSynapse Linkプロファイルを作成する必要があります。
- ADLS Gen2 ストレージ アカウント (またはアカウントを作成するためのアクセス許可)。
- Microsoft Entra ID アプリケーションを作成および構成するためのアクセス許可。
- Dataverse API v9.2 以降。
- Azure Storage REST API バージョン 2021-08-06。
- Azure Synapse Link for Dataverse バージョン 1.0 以降。
仮想エンティティまたは直接テーブルを構成する (省略可能)
仮想エンティティとダイレクト テーブルを使用すると、データをコピーすることなく、Dataverse 以外のソース (Dynamics 365 Finance & Operations など) のデータを Dataverse で利用できるようになります。 Dataverse 以外のソースの場合は、Azure Synapse Link を設定する前に、仮想エンティティまたは直接テーブルを構成する必要があります。
仮想エンティティを構成するには:
Power Apps で、[環境] ページに移動し、[Dynamics 365 アプリ] をクリックします。
Dataverse で F&O エンティティを仮想エンティティとしてリンクするには、 Finance and Operations Virtual Entity ソリューションをインストールします。
Dataverse と F&O アプリケーションの間でサービス間 (S2S) 承認を設定します。 これにより、Dataverse はアプリケーションと通信できます。 詳細については、Microsoft のドキュメント 「Dataverse 仮想エンティティの構成」を参照してください。
取り込む仮想エンティティごとに、[詳細プロパティ] で [変更履歴の追跡] を有効にします。
既定では、F&O 仮想エンティティ ソリューションは、Dataverse テーブルの一覧で既定で一部の仮想エンティティを公開します。 ただし、追加のエンティティは手動で公開できます。
- Dataverse 環境の [詳細設定] ページに移動します。
- 右上にあるフィルター アイコンをクリックして、高度な検索にアクセスします。
- ドロップダウン メニューから [利用可能な財務エンティティ] と [操作エンティティ ] を選択し、[ 結果] をクリックします。
- 公開する仮想エンティティを選択します。
- [ エンティティ管理者 ] ページで、[ 表示] を [True] に切り替え、[ 保存して閉じる] をクリックします。
これで、dataverse テーブルの一覧に、 mserp_で始まる名前のエンティティを表示できるようになりました。
Important
仮想エンティティやダイレクトテーブルは、Dataverseが同期を終えた後にのみAzure Synapse Linkに現れます。 通常は最大15分かかりますが、最大で30分かかることもあります。 30分経ってもテーブルが欠落している場合は、 スキーマ発見に現れない仮想エンティティを参照してください。
Azure Synapse Link を構成する
この手順では、 Azure Data Lake への Synapse Link for Dataverse を 使用して、取り込むテーブルを選択します。 このサービスは、かつて「データエクスポート」として知られていたサービスを置き換えAzure Data Lake Storage Gen2。 名前に反して、Azure Synapse Analyticsを使用しておらず、依存もしていません。 これはDataverseからADLS Gen2への継続的なエクスポートサービスです。
Power Apps ポータルで、[分析] をクリックし、[Azure Synapse にリンク] をクリックします。
[ 新しいリンク] をクリックします。 Dataverseは同じテナントからのアクティブなサブスクリプションを自動的に入力します。 ドロップダウンから適切なサブスクリプションを選択します。
[ Connect to your Azure Synapse Analytics Workspace]\(Azure Synapse Analytics ワークスペースに接続 する\) チェック ボックスをオンにしないでください。 データは直接CSVとしてADLS Gen2ストレージアカウントに格納され、このワークフローにはAzure Synapse Analyticsのワークスペースは必要ありません。
[ Synapse Link Creation]\(Synapse リンクの作成\) ページで、[ 詳細設定] をクリックします。 次に、[ 詳細な構成設定の表示] を切り替えます。
[ 増分更新フォルダー構造を有効にする] を切り替え、目的の Synapse Link 更新間隔を設定します。 最小値は 5 分です。 この間隔は、この Synapse Link に含まれるすべてのテーブルに適用されます。 (別の手順で Databricks パイプラインのスケジュールを設定します)。
同期したいテーブルを選択し、 Appendのみ と パーティション 設定をデフォルトに残してください。
- Dataverse ネイティブ アプリから取り込む場合は、Dataverse セクションから関連する Dataverse テーブルを直接選択します。
- F&O から取り込む場合は、 D365 Finance & Operations セクションから直接テーブルを選択するか、 Dataverse セクション (プレフィックス
mserp_) から仮想エンティティを選択できます。 仮想エンティティの詳細については、 手順 1 を参照してください。
[保存] をクリックします。 最初のSynapse Link同期が始まります。
F&O ユーザーの場合、この初期同期は、数百ギガバイトの大きなテーブルに対して数時間かかる場合があります。
注
F&Oエンティティの初期同期に時間がかかる場合は、F&Oアプリのテーブルにインデックスを作成して高速化できます:
- F&O 環境でインデックスを作成するテーブルに移動します。
- テーブルの拡張機能を作成します。
- テーブル拡張機能内で、新しいインデックスを定義します。
- インデックスに含めたいフィールドを追加すれば、それらのフィールドでのデータベース検索が速くなります。
- 変更を保存して F&O 環境にデプロイします。
インジェスト用の Entra ID アプリケーションを作成する
この手順では、Azure Databricksへのインジェストをサポートする Unity カタログ接続を作成するために必要なEntra ID情報を収集します。
右側のパネルに表示されている Entra ID テナントの テナント ID (
portal.azure.com>>Microsoft Entra ID>>Overview タブ>>テナント ID) を収集します。Azure Synapse Link を作成すると、選択したテーブルを同期するための ADLS コンテナーが Azure Synapse によって作成されます。 Synapse Linkの管理ページにアクセスして、ADLS コンテナー名を見つけます。
ADLS コンテナーのアクセス資格情報を収集します。
- Microsoft Entra ID アプリをまだ作成していない場合は作成します。
- クライアント シークレットを収集します。
-
アプリ ID (
portal.azure.com>>Microsoft Entra ID>>管理>>アプリ登録) を収集します。
ADLS コンテナーへのアクセス権を Entra ID アプリに付与します (まだアクセス権がない場合)。
注
Entra ID アプリケーションが、各 Synapse Link プロファイルに関連付けられている ADLS コンテナーにアクセスできることを確認します。 複数の環境またはアプリケーションからデータを取り込む場合は、アプリケーションに関連するすべてのコンテナーにロールの割り当てがあることを確認します。
- Azure ストレージ アカウントに移動し、コンテナーまたはストレージ アカウントを選択します。 (Azure Databricksは、最小限の特権を維持するためにコンテナー レベルを推奨します)。
- [ アクセス制御 (IAM)]、[ ロールの割り当ての追加] の順にクリックします。
- ストレージ BLOB データ共同作成者→読み取り/書き込み/削除アクセス ロールを選択します。 組織でこれを許可していない場合は、Azure Databricks アカウント チームにお問い合わせください。
- [ 次へ]、[メンバーの選択] の順に クリックします。
- [ ユーザー、グループ、またはサービス プリンシパル] を選択し、[ アプリの登録を検索] を選択します。 (検索結果にアプリが存在しない場合は、検索バーにオブジェクト ID を明示的に入力し、 Enter キーを押します)。
- [ Review + Assign をクリックします。
- アクセス許可が正しく構成されていることを確認するには、コンテナーのAccess Controlを確認します。
Dynamics 365 パイプラインを作成する
パイプラインはUIやAPIを通じて作成できます。 UIウィザードは接続とパイプラインを同時に処理し、APIパスはそれらを2つの別々のステップとして作成します。
UI を使用する
ウィザードは、前のステップで集めたEntra IDのアプリケーションの認証情報とストレージの詳細を入力するよう促し、その後、接続とパイプラインを一緒に作成します。
- 左側のメニューで、[ 新規] をクリックし、[ データの追加またはアップロード] をクリックします。
- [ データの追加] ページで、[ Dynamics 365 ] タイルをクリックします。
- そこからウィザードの指示に従います。
API の使用
まず接続を作成し、それを使うパイプラインを作りましょう。 最初のステップで取得した接続名を、2番目のステップでパイプラインを定義する必要があります。
ステップ1:Dynamics 365接続を作成する
この手順では、Unity カタログ接続を作成して、Dynamics 365資格情報を安全に格納し、Azure Databricksへのインジェストを開始します。
- ワークスペースで、
[カタログ] をクリックします。
- [
[接続] をクリックし、[ 接続] をクリックします。
- [ 接続の作成 ] ボタンをクリックします。
- 一意の接続名を指定し、接続の種類として Dynamics 365 を選択します。
- 前の手順で作成した Entra ID アプリの クライアント シークレット と クライアント ID を入力します。 スコープを変更しないでください。 [次へ] をクリックします。
- Azure ストレージ アカウント名、テナント ID、ADLS コンテナー名を入力し、[接続の作成] をクリックします。
- 接続名を書き留めます。
ステップ2:インジェスションパイプラインの作成
この手順では、インジェスト パイプラインを設定します。 取り込まれた各テーブルは、宛先内で同じ名前の対応するストリーミング テーブルを取得します。 ノートブックかDatabricksのCLIを使うことができます。 どちらの方法でも、パイプラインを作成する Databricks サービスへの API 呼び出しが行われます。
ノートブックを使用する
このページの最後にあるテンプレートは、パイプラインの作成と管理のためのヘルパー関数を定義しています。 最初のセルはそれらの関数を設定し、2つ目は自分でパイプラインを定義する場所です。
- ノートブック テンプレートをコピーします。
- ノートブックの最初のセルを変更せずに実行します。
- パイプラインの詳細 (取り込むテーブル、データを格納する場所など) を使用して、ノートブックの 2 番目のセルを変更します。
- テンプレート ノートブックの 2 番目のセルを実行します。これにより、
create_pipelineが実行されます。 -
list_pipelineを実行して、パイプライン ID とその詳細を表示できます。 -
edit_pipelineを実行してパイプライン定義を編集できます。 -
delete_pipelineを実行してパイプラインを削除できます。
Databricks CLI を使用する
パイプラインを作成するために
databricks pipelines create --json "<pipeline_definition OR json file path>"
パイプラインを編集するには:
databricks pipelines update --json "<<pipeline_definition OR json file path>"
パイプライン定義を取得するには:
databricks pipelines get "<your_pipeline_id>"
パイプラインを削除するには:
databricks pipelines delete "<your_pipeline_id>"
詳細については、いつでも次のコマンドを実行できます。
databricks pipelines --help
databricks pipelines <create|update|get|delete|...> --help
追加機能の構成 (省略可能)
コネクタには、履歴の追跡、列レベルの選択、選択解除のための SCD タイプ 2 などの追加機能が用意されています。 マネージド インジェスト パイプラインの一般的なパターンを参照してください。
ノートブック テンプレート
両方のセルをワークスペース内のノートブックにコピーしてください。 セル1はパイプラインAPIを呼び出すヘルパー関数を定義し、セル2は作成したいパイプラインを定義する場所です。
セル1:API設定
このセルをコピーして変更せずに as-is 実行してください。 セル2が呼び出す create_pipeline、 list_pipeline、 edit_pipeline、 delete_pipeline、その他のヘルパーを定義します。
# DO NOT MODIFY
# This sets up the API utils for creating managed ingestion pipelines in Databricks.
import requests
import json
notebook_context = dbutils.notebook.entry_point.getDbutils().notebook().getContext()
api_token = notebook_context.apiToken().get()
workspace_url = notebook_context.apiUrl().get()
api_url = f"{workspace_url}/api/2.0/pipelines"
headers = {
'Authorization': 'Bearer {}'.format(api_token),
'Content-Type': 'application/json'
}
def check_response(response):
if response.status_code == 200:
print("Response from API:\n{}".format(json.dumps(response.json(), indent=2, sort_keys=False)))
else:
print(f"Failed to retrieve data: error_code={response.status_code}, error_message={response.json().get('message', response.text)}")
def create_pipeline(pipeline_definition: str):
response = requests.post(url=api_url, headers=headers, data=pipeline_definition)
check_response(response)
def edit_pipeline(id: str, pipeline_definition: str):
response = requests.put(url=f"{api_url}/{id}", headers=headers, data=pipeline_definition)
check_response(response)
def delete_pipeline(id: str):
response = requests.delete(url=f"{api_url}/{id}", headers=headers)
check_response(response)
def list_pipeline(filter: str):
body = "" if len(filter) == 0 else f"""{{"filter": "{filter}"}}"""
response = requests.get(url=api_url, headers=headers, data=body)
check_response(response)
def get_pipeline(id: str):
response = requests.get(url=f"{api_url}/{id}", headers=headers)
check_response(response)
def start_pipeline(id: str, full_refresh: bool=False):
body = f"""
{{
"full_refresh": {str(full_refresh).lower()},
"validate_only": false,
"cause": "API_CALL"
}}
"""
response = requests.post(url=f"{api_url}/{id}/updates", headers=headers, data=body)
check_response(response)
def stop_pipeline(id: str):
print("cannot stop pipeline")
セル2:パイプライン定義
Synapse Linkのデータを取り込みたい量に応じて、以下の2つの選択肢から選んでください。
- オプションA、スキーマレベルの仕様:Azure Synapse Linkで同期されたすべてのテーブルを取り込みます。 Azure Databricksはパイプラインあたり250以上のテーブルを推奨していないので、Synapse Linkがそれ以上同期する場合は、複数のパイプラインに分散してテーブルを分けてください。
-
オプションB、テーブルレベルの仕様:指定したテーブルのみを取り込みます。 各
source_table値は、Synapse Link管理ページの「Name」欄のテーブル名と一致しなければなりません。
プレースホルダーの値は自分のものに置き換え、 "channel": "PREVIEW" as-isはそのままにしてください。
# Option A: schema-level spec
pipeline_spec = """
{
"name": "<YOUR_PIPELINE_NAME>",
"ingestion_definition": {
"connection_name": "<YOUR_CONNECTION_NAME>",
"objects": [
{
"schema": {
"source_schema": "objects",
"destination_catalog": "<YOUR_DATABRICKS_CATALOG>",
"destination_schema": "<YOUR_DATABRICKS_SCHEMA>"
}
}
]
},
"channel": "PREVIEW"
}
"""
create_pipeline(pipeline_spec)
# Option B: table-level spec
pipeline_spec = """
{
"name": "<YOUR_PIPELINE_NAME>",
"ingestion_definition": {
"connection_name": "<YOUR_CONNECTION_NAME>",
"objects": [
{
"table": {
"source_schema": "objects",
"source_table": "<YOUR_F_AND_O_TABLE_NAME>",
"destination_catalog": "<YOUR_DATABRICKS_CATALOG>",
"destination_schema": "<YOUR_DATABRICKS_SCHEMA>"
}
}
]
},
"channel": "PREVIEW"
}
"""
create_pipeline(pipeline_spec)
例:SCDタイプ2のトラック履歴
既定では、API は SCD 型 1 を使用します。 ソースでデータが編集された場合、それに応じて宛先のデータが上書きされます。 履歴データを保持し、SCD タイプ 2 を使用する場合は、構成で指定します。 例えば次が挙げられます。
# Schema-level spec with SCD type 2
pipeline_spec = """
{
"name": "<YOUR_PIPELINE_NAME>",
"ingestion_definition": {
"connection_name": "<YOUR_CONNECTION_NAME>",
"objects": [
{
"schema": {
"source_schema": "objects",
"destination_catalog": "<YOUR_DATABRICKS_CATALOG>",
"destination_schema": "<YOUR_DATABRICKS_SCHEMA>",
"table_configuration": {
"scd_type": "SCD_TYPE_2"
}
}
}
]
},
"channel": "PREVIEW"
}
"""
create_pipeline(pipeline_spec)
# Table-level spec with SCD type 2
pipeline_spec = """
{
"name": "<YOUR_PIPELINE_NAME>",
"ingestion_definition": {
"connection_name": "<YOUR_CONNECTION_NAME>",
"objects": [
{
"table": {
"source_schema": "objects",
"source_table": "<YOUR_F_AND_O_TABLE_NAME>",
"destination_catalog": "<YOUR_DATABRICKS_CATALOG>",
"destination_schema": "<YOUR_DATABRICKS_SCHEMA>",
"table_configuration": {
"scd_type": "SCD_TYPE_2"
}
}
}
]
},
"channel": "PREVIEW"
}
"""
create_pipeline(pipeline_spec)
例:特定の列を含めるか除外するか
既定では、選択したテーブルのすべての列が API によって取り込まれます。 ただし、特定の列を含めるか除外するかを選択できます。 例えば次が挙げられます。
# Table spec with included and excluded columns.
pipeline_spec = """
{
"name": "<YOUR_PIPELINE_NAME>",
"ingestion_definition": {
"connection_name": "<YOUR_CONNECTON_NAME>",
"objects": [
{
"table": {
"source_schema": "objects",
"source_table": "<YOUR_F_AND_O_TABLE_NAME>",
"destination_catalog": "<YOUR_DATABRICKS_CATALOG>",
"destination_schema": "<YOUR_DATABRICKS_SCHEMA>",
"table_configuration": {
"include_columns": ["<COLUMN_A>", "<COLUMN_B>", "<COLUMN_C>"]
}
}
}
]
},
"channel": "PREVIEW"
}
"""
create_pipeline(pipeline_spec)