SQL Server からのデータの取り込み

Lakeflow Connect を使用して SQL Server から Azure Databricks にデータを取り込む方法について説明します。

SQL Server コネクタは、Azure SQL Database、Azure SQL Managed Instance、および Amazon RDS SQL データベースをサポートしています。 これには、Azure 仮想マシン (VM) と Amazon EC2 で実行されている SQL Server が含まれます。 このコネクタでは、Azure ExpressRoute と AWS Direct Connect ネットワークを使用したオンプレミスの SQL Server もサポートされています。

必要条件

  • インジェスト ゲートウェイとインジェスト パイプラインを作成するには、まず次の要件を満たす必要があります。

    • ワークスペースが Unity Catalog に対して有効になっている。

    • ワークスペースに対してサーバーレス コンピューティングが有効になっています。 サーバーレス コンピューティング要件を参照してください。

    • 接続を作成する場合: メタストアに対する CREATE CONNECTION 特権があります。 「Unity Catalog の特権の管理」を参照してください。

      コネクタで UI ベースのパイプライン作成がサポートされている場合は、このページの手順を完了することで、接続とパイプラインを同時に作成できます。 ただし、API ベースのパイプライン作成を使用する場合は、このページの手順を完了する前に、カタログ エクスプローラーで接続を作成する必要があります。 「マネージド インジェスト ソースへの接続」を参照してください。

    • 既存の接続を使用する場合: 接続に USE CONNECTION 特権または ALL PRIVILEGES があります。

    • ターゲット カタログに対する USE CATALOG 権限があります。

    • 既存のスキーマに対する USE SCHEMACREATE TABLE、および CREATE VOLUME 特権、またはターゲット カタログに対する CREATE SCHEMA 特権があります。

    • プライマリ SQL Server インスタンスにアクセスできます。 変更の追跡と変更データ キャプチャ機能は、読み取りレプリカまたはセカンダリ インスタンスではサポートされていません。

    • クラスターを作成するための無制限のアクセス許可、またはカスタム ポリシー (API のみ)。 ゲートウェイのカスタム ポリシーは、次の要件を満たしている必要があります。

      • ファミリ: ジョブ コンピューティング

      • ポリシー ファミリのオーバーライド:

        {
          "cluster_type": {
            "type": "fixed",
            "value": "dlt"
          },
          "num_workers": {
            "type": "unlimited",
            "defaultValue": 1,
            "isOptional": true
          },
          "runtime_engine": {
            "type": "fixed",
            "value": "STANDARD",
            "hidden": true
          }
        }
        
      • Databricks では、インジェスト ゲートウェイのワーカー ノードはゲートウェイのパフォーマンスに影響しないため、可能な限り最小のワーカー ノードを指定することをお勧めします。 次のコンピューティング ポリシーを使用すると、Azure Databricks は、ワークロードのニーズに合わせてインジェスト ゲートウェイをスケーリングできます。 ドライバーノードの最低要件は8コアで、効率的かつ高性能なデータ抽出を可能にします。

        {
          "driver_node_type_id": {
            "type": "fixed",
            "value": "Standard_E64d_v4"
          },
          "node_type_id": {
            "type": "fixed",
            "value": "Standard_F4s"
          }
        }
        

      クラスター ポリシーの詳細については、「 コンピューティング ポリシーの選択」を参照してください。

  • SQL Server から取り込むには、「 Azure Databricks への取り込み用に Microsoft SQL Server を構成する」の手順を最初に完了する必要があります。

ゲートウェイとインジェスト パイプラインを作成する

Warnung

インジェスト ゲートウェイを手動で停止しないでください。 ソース データベースで変更ログが切り捨てられる前に、ゲートウェイを継続的に実行して変更をキャプチャする必要があります。 ゲートウェイが停止すると、ログの保持のために変更を削除でき、影響を受けるすべてのテーブルを完全に更新する必要があります。 ゲートウェイを停止して再起動すると、VM も再プロビジョニングされるため、起動時間が長くなります。 ゲートウェイの問題をトラブルシューティングする必要がある場合は、インジェストSQL Serverトラブルシューティングを参照するか、Databricks サポートにお問い合わせください。

Databricks ユーザーインターフェース

  1. Azure Databricks ワークスペースのサイドバーで、[ Data Ingestion をクリックします。

  2. [ データの追加 ] ページの [ Databricks コネクタ] で、[ SQL Server] をクリックします。

  3. インジェスト ウィザードの Connection ページで、SQL Server アクセス資格情報を格納する接続を選択します。 メタストアに対する CREATE CONNECTION 特権がある場合は、 Plus icon. 接続の作成 をクリックして、SQL Server接続を作成の認証の詳細を使用して新しい接続を作成できます。

  4. [次へ] をクリックします。

  5. [ インジェストのセットアップ ] ページで、インジェスト パイプラインの一意の名前を入力します。 このパイプラインは、ステージングの場所から宛先にデータを移動します。

  6. イベント ログを書き込むカタログとスキーマを選択します。 イベント ログには、監査ログ、データ品質チェック、パイプラインの進行状況、エラーが含まれます。 カタログに対するUSE CATALOG権限とCREATE SCHEMA権限がある場合は、[プラス] アイコンをクリックできます。ドロップダウン メニューでスキーマを作成し、新しいスキーマを作成します。

  7. (省略可能) すべてのテーブルの自動更新を[オン] に設定します。 自動更新がオンの場合、パイプラインは、影響を受けるテーブルを完全に更新することで、ログ クリーンアップ イベントや特定の種類のスキーマの進化などの問題を自動的に修正しようとします。 履歴の追跡が有効になっている場合、完全な更新でその履歴が消去されます。

  8. インジェスト ゲートウェイの一意の名前を入力します。 ゲートウェイは、ソースから変更を抽出し、インジェスト パイプラインが読み込まれるようステージングするパイプラインです。

  9. ステージング場所のカタログとスキーマを選択します。 この場所にボリュームが作成され、抽出されたデータがステージングされます。 カタログに対するUSE CATALOG権限とCREATE SCHEMA権限がある場合は、[プラス] アイコンをクリックできます。ドロップダウン メニューでスキーマを作成し、新しいスキーマを作成します。

  10. [パイプラインの作成] をクリックして続行します

  11. [ ソース ] ページで、取り込むテーブルを選択します。 特定のテーブルを選択した場合は、テーブル設定を構成できます。

    a. (省略可能)[ 設定] タブで、取り込まれた各テーブルの 宛先名 を指定します。 これは、オブジェクトを同じスキーマに複数回取り込むときに、変換先テーブルを区別するのに役立ちます。 「 宛先テーブルに名前を付けます」を参照してください。

    a. (省略可能)既定の 履歴追跡 設定を変更します。 「 履歴追跡の有効化 (SCD タイプ 2)」を参照してください。

  12. [ 次へ] をクリックし、[ 保存] をクリックして続行します。

  13. [ 宛先 ] ページで、データを読み込むカタログとスキーマを選択します。 カタログに対するUSE CATALOG権限とCREATE SCHEMA権限がある場合は、[プラス] アイコンをクリックできます。ドロップダウン メニューでスキーマを作成し、新しいスキーマを作成します。

  14. [ 保存] をクリックして続行します。

  15. [ データベースのセットアップ ] ページで、[ 検証 ] をクリックして、ソースが Azure Databricks インジェスト用に正しく構成されていることを確認します。 不足している構成が返されます。 解決する手順については、[ 構成の完了] をクリックします。 続けて、 [次へ] をクリックします。 または、[検証の スキップ] をクリックします。

  16. (省略可能)[ スケジュールと通知 ] ページで、[プラス] アイコンをクリック します。 スケジュールを作成します。 変換先テーブルを更新する頻度を設定します。

  17. (省略可能)[ プラス] アイコンをクリックします。 パイプライン 操作の成功または失敗の電子メール通知を設定する通知を追加し、[ パイプラインの保存と実行] をクリックします。

宣言型オートメーション バンドル

宣言型オートメーション バンドルを使用して取り込む前に、既存の接続にアクセスできる必要があります。 手順については、「SQL Server接続を作成するを参照してください。

ステージング カタログとスキーマは、移行先のカタログとスキーマと同じにすることができます。 ステージング カタログを外部カタログにすることはできません。 バンドル パイプライン YAML ファイルの gateway_definition セクションでステージングの場所を指定します。

インジェスト ゲートウェイは、ソース データベースからスナップショットと変更データを抽出し、Unity カタログステージング ボリュームに格納します。 ゲートウェイを継続的なパイプラインとして実行する必要があります。 これは、ソース データベース上にある変更ログ保持ポリシーに対応するのに役立ちます。

インジェスト パイプラインは、スナップショットを適用し、ステージング ボリュームから宛先ストリーミング テーブルにデータを変更します。

バンドルには、ジョブとタスクの YAML 定義を含め、Databricks CLI を使用して管理できます。また、さまざまなターゲット ワークスペース (開発、ステージング、運用など) で共有および実行できます。 詳細については、「 宣言型オートメーション バンドルとは」を参照してください。

  1. Databricks CLI を使用してバンドルを作成します。

    databricks bundle init
    
  2. パイプラインとジョブの構成をバンドルに追加します。 使用可能なすべてのオプション 含む完全な例については、「例」を参照してください。

  3. Databricks CLI を使用してパイプラインをデプロイします:

    databricks bundle deploy
    

Databricks ノートブック

ソース接続、ターゲット カタログ、ターゲット スキーマ、およびソースから取り込むテーブルを使用して、次のノートブックの Configuration セルを更新します。

ノートブックを入手

Terraform

Terraform を使用して、SQL Server インジェスト パイプラインをデプロイおよび管理できます。 ゲートウェイとインジェスト パイプラインを作成するための Terraform 構成を含む完全なサンプル フレームワークについては、GitHub の Lakeflow Connect Terraform サンプル リポジトリを参照してください。

データ インジェストが正常に行われたことを確認する

パイプラインの詳細ページのリスト ビューには、データが取り込まれると処理されたレコードの数が表示されます。 これらの数値は自動的に更新されます。

レプリケーションを確認する

Upserted records 列と Deleted records 列は、既定では表示されません。 有効にするには、[列の構成] [列の構成] アイコン ボタンをクリックして選択します。

例示

これらの例を使用して、パイプラインを構成します。

パイプラインの構成

宣言型オートメーション バンドル

次のバンドルは、ゲートウェイ パイプライン、インジェスト パイプライン、およびスケジュールされたジョブを定義します。 コメントアウト オプションには、使用可能なすべての構成が表示されます。 variablesセクションとtargetsセクションをソースと宛先の詳細で更新します。

bundle:
  name: lakeflow-connect-sqlserver

# Variables parameterize the bundle for different environments and sources.
# Set values here, override per-target, or pass with: databricks bundle deploy -var="key=value"
variables:
  # The name of the Unity Catalog connection to your SQL Server instance.
  # This connection must already exist and be of type SQLSERVER.
  connection_name:
    description: 'Unity Catalog connection name for the SQL Server source'
  # The SQL Server database name to ingest from.
  # In Lakeflow Connect, this maps to source_catalog in the table/schema spec.
  source_database:
    description: 'SQL Server database name (maps to source_catalog in table specs)'
  # The SQL Server schema to ingest from (for example, "dbo", "sales").
  source_schema:
    description: 'SQL Server schema name to ingest from'
  # The Unity Catalog catalog where ingested Delta tables are created.
  dest_catalog:
    description: 'Destination Unity Catalog catalog for ingested tables'
  # The Unity Catalog schema where ingested Delta tables are created.
  dest_schema:
    description: 'Destination Unity Catalog schema for ingested tables'
  # The Unity Catalog catalog for the gateway's internal staging volume.
  # Can be the same as dest_catalog. Must not be a foreign catalog.
  staging_catalog:
    description: 'Catalog for gateway staging volume'
  # The Unity Catalog schema for the gateway's internal staging volume.
  staging_schema:
    description: 'Schema for gateway staging volume'

resources:
  pipelines:
    # --- Gateway pipeline ---
    # Extracts change data from SQL Server and stages it in a Unity Catalog
    # volume. Must run continuously to capture changes before change logs are
    # truncated in the source database.
    gw_pipeline:
      name: 'lfc-sqlserver-gateway-${bundle.target}'
      # Gateway pipelines must be continuous.
      continuous: true
      # "CURRENT" (stable) or "PREVIEW" (early access).
      channel: 'CURRENT'
      # (Optional) Associate with a budget policy for cost tracking.
      # budget_policy_id: "<policy-uuid>"
      # The gateway runs on classic compute. Cluster settings are managed
      # automatically. You can optionally customize the cluster:
      # clusters:
      #   - label: "default"
      #     autoscale:
      #       min_workers: 1
      #       max_workers: 4
      #     # node_type_id: "i3.xlarge"
      #     # Restrict the cluster to an approved cluster policy.
      #     # policy_id: "<cluster-policy-id>"
      catalog: ${var.staging_catalog}
      schema: ${var.staging_schema}
      gateway_definition:
        # (Required) Unity Catalog connection name (type SQLSERVER).
        connection_name: ${var.connection_name}
        # (Required) Catalog and schema for the staging volume.
        gateway_storage_catalog: ${var.staging_catalog}
        gateway_storage_schema: ${var.staging_schema}
        # (Optional) Custom staging volume name. If not set, the system
        # auto-generates: __databricks_ingestion_gateway_staging_data-<pipeline_id>
        # gateway_storage_name: "my_custom_staging_volume"

    # --- Ingestion pipeline ---
    # Reads staged data from the gateway and applies it to Delta tables.
    mi_pipeline:
      name: 'lfc-sqlserver-ingestion-${bundle.target}'
      # Continuous mode is not supported for the ingestion pipeline.
      # Use a scheduled job to trigger runs.
      continuous: false
      channel: 'CURRENT'
      # (Optional) Associate with a budget policy for cost tracking.
      # budget_policy_id: "<policy-uuid>"
      # The ingestion pipeline runs on serverless compute only.
      serverless: true
      # (Optional) Development mode for faster iteration (no retries).
      # development: true
      catalog: ${var.dest_catalog}
      schema: ${var.dest_schema}
      # (Optional) Email notifications for pipeline events.
      # notifications:
      #   - email_recipients:
      #       - "team@example.com"
      #     alerts:
      #       - "on-update-failure"
      #       - "on-update-fatal-failure"
      #       - "on-flow-failure"
      # (Optional) Run as a service principal for production.
      # run_as:
      #   service_principal_name: "my-service-principal"
      ingestion_definition:
        # (Required) References the gateway pipeline. The connection is
        # inherited from the gateway. Do not specify connection_name here.
        ingestion_gateway_id: ${resources.pipelines.gw_pipeline.id}
        # Pipeline-level table configuration defaults. These apply to all
        # tables unless overridden at the schema or table level.
        table_configuration:
          # SCD Type: How changes are applied to destination tables.
          #   SCD_TYPE_1: Overwrites rows with latest values (default).
          #   SCD_TYPE_2: Preserves history with __START_AT/__END_AT columns.
          #               Requires CDC on source. CT does not support SCD_TYPE_2.
          scd_type: 'SCD_TYPE_1'
          # (Optional) Auto full refresh policy. Triggers a snapshot when the
          # pipeline detects issues resolvable by re-reading all source data
          # (for example, CT/CDC retention window expired).
          # auto_full_refresh_policy:
          #   enabled: true
          #   min_interval_hours: 24
        # (Optional) Schedule automatic full refreshes.
        # full_refresh_window:
        #   start_hour: 2
        #   days_of_week:
        #     - "SUNDAY"
        #   time_zone_id: "America/Los_Angeles"
        objects:
          # Option 1: Schema-level ingestion. Ingests all tables from a source
          # schema. New tables added to the schema are picked up automatically.
          - schema:
              source_catalog: ${var.source_database}
              source_schema: ${var.source_schema}
              destination_catalog: ${var.dest_catalog}
              destination_schema: ${var.dest_schema}
              # (Optional) Override table_configuration for this schema.
              # table_configuration:
              #   scd_type: "SCD_TYPE_2"
          # Option 2: Table-level ingestion. Provides granular control.
          # Replace or combine with the schema-level spec.
          # - table:
          #     source_catalog: ${var.source_database}
          #     source_schema: ${var.source_schema}
          #     source_table: "customers"
          #     destination_catalog: ${var.dest_catalog}
          #     destination_schema: ${var.dest_schema}
          #     # (Optional) Rename the table at the destination.
          #     # destination_table: "customers_v2"
          #     table_configuration:
          #       scd_type: "SCD_TYPE_1"
          #       # Include only specific columns (mutually exclusive with exclude_columns).
          #       # include_columns:
          #       #   - "customer_id"
          #       #   - "first_name"
          #       #   - "email"
          #       # Exclude specific columns. All other columns are included.
          #       # exclude_columns:
          #       #   - "internal_notes"
          #       # Override the primary key used for change detection.
          #       # primary_keys:
          #       #   - "customer_id"
          #       # Logical ordering columns for change resolution.
          #       # sequence_by:
          #       #   - "updated_at"
          #       # Auto full refresh for this table.
          #       # auto_full_refresh_policy:
          #       #   enabled: true
          #       #   min_interval_hours: 48
      # (Optional) Grant additional users or groups access.
      # permissions:
      #   - user_name: "analyst@example.com"
      #     level: "CAN_VIEW"
      #   - group_name: "data-engineers"
      #     level: "CAN_RUN"

  # --- Scheduled job ---
  # Triggers the ingestion pipeline on a schedule.
  jobs:
    mi_schedule:
      name: 'lfc-sqlserver-ingestion-schedule-${bundle.target}'
      # Quartz cron syntax: "seconds minutes hours day month day-of-week"
      # Examples: "0 0 * * * ?" (hourly), "0 0 */4 * * ?" (every 4 hours)
      schedule:
        quartz_cron_expression: '0 */30 * * * ?'
        timezone_id: 'UTC'
      tasks:
        - task_key: 'run_ingestion'
          pipeline_task:
            pipeline_id: ${resources.pipelines.mi_pipeline.id}
      # email_notifications:
      #   on_failure:
      #     - "team@example.com"

# Deploy to different workspaces with: databricks bundle deploy -t <target>
targets:
  dev:
    default: true
    workspace:
      host: https://<workspace-url>.cloud.databricks.com
    variables:
      connection_name: '<sqlserver-connection>'
      source_database: '<database-name>'
      source_schema: 'dbo'
      dest_catalog: '<dest-catalog>'
      dest_schema: '<dest-schema>'
      staging_catalog: '<staging-catalog>'
      staging_schema: '<staging-schema>'

Databricks ノートブック

パイプライン仕様の Configuration セクションの例を次に示します。

# The name of the UC connection with the credentials to access the source database
connection_name = "my_connection"

# The name of the UC catalog and schema to store the replicated tables
target_catalog_name = "main"
target_schema_name = "lakeflow_sqlserver_connector_cdc"

# The name of the UC catalog and schema to store the staging volume with intermediate
# CDC and snapshot data. Use the destination catalog/schema by default.
stg_catalog_name = target_catalog_name
stg_schema_name = target_schema_name

# The name of the Gateway pipeline to create
gateway_pipeline_name = "cdc_gateway"

# The name of the Ingestion pipeline to create
ingestion_pipeline_name = "cdc_ingestion"

# Construct the full list of tables to replicate.
# IMPORTANT: The letter case of catalog, schema, and table names must match exactly
# the case used in the source database system tables.
tables_to_replicate = replicate_full_db_schema("MY_DB", ["MY_DB_SCHEMA"])
# Append tables from additional schemas as needed:
#  + replicate_tables_from_db_schema("MY_DB", "MY_SCHEMA_2", ["table3", "table4"])

一般的なパターン

高度なパイプライン構成については、「 マネージド インジェスト パイプラインの一般的なパターン」を参照してください。

次のステップ

パイプラインを開始し、スケジュールを設定し、アラートを設定します。 一般的なパイプライン メンテナンス タスクを参照してください。

その他のリソース