Oracle用の統合CDCパイプラインを作成
ベータ版
この機能はベータ版です。ワークスペース管理者は、 プレビュー ページからこの機能へのアクセスを制御できます。Databricksのプレビューを管理するを参照してください。
統合されたCDCパイプラインは、単一のパイプラインを使用してOracleからDatabricksへ変更データを取り込みます。統合 CDC コネクタは、抽出と適用を 1 つのパイプライン更新に統合します。
Oracle 統合 CDC コネクターは、LogMiner を未コミット トランザクション モードで使用して、オンライン redo Logs とアーカイブ Logs から変更を読み取ります。
要件
-
ワークスペースは Unity Catalog が有効になっています。
-
接続を作成する予定の場合:メタストアに対する
CREATE CONNECTIONの特権が必要です。Unity Catalog での特権の管理を参照してください。コネクタがUIベースのパイプラインオーサリングをサポートしている場合、このページのステップを完了することで、接続とパイプラインを同時に作成できます。ただし、API ベースのパイプライン オーサリングを使用する場合は、このページのステップを完了する前に Catalog Explorer で接続を作成する必要があります。「管理対象取り込みソースへの接続」を参照してください。
-
既存の接続を使用する場合:接続に対する
USE CONNECTION権限、またはALL PRIVILEGESが必要です。 -
ターゲットカタログに対する
USE CATALOG権限があります。 -
既存のスキーマに対して
USE SCHEMA、CREATE TABLE、およびCREATE VOLUMEの権限、またはターゲットカタログに対してCREATE SCHEMAの権限があります。 -
ワークスペースは統合CDCコネクター機能が有効になっていなければなりません。Databricks アカウント チームにお問い合わせください。
-
Oracle ソースデータベースのセットアップが完了しました。Databricks への取り込み用に Oracle を構成するを参照してください。
-
次の権限があります。
CREATE CONNECTIONメタストアの(新しいUnity Catalog接続を作成する場合)、または既存の接続のUSE CONNECTION。USE CATALOG宛先カタログで。USE SCHEMACREATE TABLEおよびアップグレード先スキーマでCREATE VOLUME宛先スキーマに、またはdata_staging_optionsで指定されたスキーマに
Oracle マルチテナント データベースの場合、接続ユーザーは CDB$ROOT の共通ユーザーである必要があります。詳細については、「レプリケーションユーザーの作成」を参照してください。
コンピュートの要件
統合された CDC パイプラインは、クラシック コンピュートまたはサーバレス コンピュートで実行されます。
- Classic コンピュート : Classic コンピュートプレーンは、Databricks ワークスペースのVPCまたはVNetで実行され、ネットワーク経由でOracleインスタンスにアクセスできる必要があります。コンピュートプレーンがデータベースに到達できるあらゆるネットワークパス(VPCまたはVNetピアリング、パブリックエンドポイント、およびオンプレミスのOracleに対するAWS Direct Connect、Azure ExpressRoute、またはVPNなど)がサポートされています。
- サーバレス コンピュート :Databricks サーバレス コンピュートとソースデータベース間のサーバレス ネットワーク接続を構成します。オンプレミスのソースでは、構成済みのサーバレス エグレスを経由するネットワーク パスが必要です(たとえば、トランジット ゲートウェイ、または ExpressRoute もしくは VPN を使用したピア接続済み VNet など)。
従来のコンピュートの場合、無制限のクラスター作成権限、またはcluster_typeをdltに、runtime_engineをSTANDARDに固定し、効率的な抽出のために少なくとも8コアが推奨されるカスタムクラスターポリシーを使用できます。
Oracleへの Unity Catalog接続を作成します。
パイプラインを作成する前に、Oracle への Unity Catalog 接続を作成してください。Oracle接続を作成するを参照してください。
統合CDCパイプラインを作成
データ取り込みUI、REST API、Databricks CLI、ノートブック、またはDeclarative Automation Bundlesを使用して、統合CDCパイプラインを作成します。
プログラムによるパイプライン作成リクエストにはすべて、"channel": "PREVIEW" を含める必要があります。UI を使用する場合、Databricks がチャンネルを自動的に設定します。
Oracle統合CDCパイプラインの場合、source_catalogはOracleサービス名にマッピングされます。マルチテナントデータベースの場合、これはCDB$ROOTサービス名である必要があります。
- Databricks UI
- Declarative Automation Bundles
- Databricks notebook
- Databricks CLI
- REST API
-
サイドバーで 「データ取り込み」 をクリックし、ソースタイプとして 「Oracle」 を選択します。

-
使用する接続を選択します。既存の Unity Catalog 接続を選択するか、新しい接続を作成します。


-
パイプラインの名前とイベントLogsの場所を指定します。イベントLogsの場所は、Databricks がステージングデータと CDC の実行に使用されるメタデータを保存する場所です。

-
次へ をクリックします。Databricksはコンピュートをプロビジョニングし、パイプラインを作成します。このステップには時間がかかる場合があり、
Waiting for resourcesが表示されます。完了したら、取り込むソーステーブルを選択します。
-
パイプラインがソースからキャプチャしたデータを書き込む宛先スキーマを選択します。パイプラインは、選択したスキーマ内にソースと同じ名前のテーブルを自動作成します。

-
「検証」 をクリックし、検証が成功するまで待ちます。

-
パイプラインのスケジュールを設定します。パイプラインはデータが利用可能な限りランし、アイドル状態に達すると停止し、次のTriggerで同じポイントから再開されます。

-
パイプラインを確認します。リストビューには、レプリケートされたデータに関するフローと統計が表示されます。

-
パイプラインの動作を確認するには、特に更新が失敗した際の警告やエラーメッセージを確認するために、右側の [Event logs] パネルを開きます。

パイプラインのセットアップが完了し、実行中です。パイプラインが宛先スキーマに作成するテーブルをクエリーし、メダリオンアーキテクチャにおけるブロンズテーブルとして扱うことができます。
バンドルファイルでパイプラインリソースを定義します(例:resources/oracle_integrated_cdc_pipeline.yml):
variables:
pipeline_name:
description: 'Name for the integrated CDC pipeline'
connection_name:
description: 'Unity Catalog connection name'
dest_catalog:
description: 'Destination catalog for ingested data'
dest_schema:
description: 'Destination schema for ingested data'
resources:
pipelines:
oracle_integrated_cdc_pipeline:
name: ${var.pipeline_name}
channel: PREVIEW
catalog: ${var.dest_catalog}
schema: ${var.dest_schema}
ingestion_definition:
connection_name: ${var.connection_name}
connector_type: CDC
objects:
- table:
source_catalog: 'ORCL'
source_schema: 'HR'
source_table: 'EMPLOYEES'
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
destination_table: 'employees'
table_configuration:
scd_type: 'SCD_TYPE_1'
スケジュールでパイプラインを実行するには、パイプラインをトリガーするジョブを定義します。各抽出段階は少なくとも10分間実行されるため、60分以上の間隔が適切な目安となります。
resources:
jobs:
oracle_integrated_cdc_job:
name: '${var.pipeline_name}-job'
tasks:
- task_key: 'cdc_ingestion'
pipeline_task:
pipeline_id: ${resources.pipelines.oracle_integrated_cdc_pipeline.id}
schedule:
quartz_cron_expression: '0 0 * * * ?'
timezone_id: 'UTC'
Databricks CLI でバンドルをデプロイする:
databricks bundle deploy
databricks bundle run oracle_integrated_cdc_job
詳細については、「宣言型オートメーションバンドルとは」をご覧ください。
from databricks.sdk import WorkspaceClient
from databricks.sdk.service.pipelines import (
ConnectorType,
IngestionConfig,
IngestionPipelineDefinition,
TableSpec,
)
w = WorkspaceClient()
pipeline = w.pipelines.create(
name="<pipeline-name>",
channel="PREVIEW",
catalog="<destination-catalog>",
schema="<destination-schema>",
ingestion_definition=IngestionPipelineDefinition(
connection_name="<oracle-connection-name>",
connector_type=ConnectorType.CDC,
objects=[
IngestionConfig(
table=TableSpec(
source_catalog="<oracle-service-name>",
source_schema="<oracle-schema>",
source_table="<oracle-table>",
destination_catalog="<destination-catalog>",
destination_schema="<destination-schema>",
)
)
],
),
)
print(f"Pipeline created: {pipeline.pipeline_id}")
databricks pipelines create --json '{
"name": "<pipeline-name>",
"channel": "PREVIEW",
"catalog": "<destination-catalog>",
"schema": "<destination-schema>",
"ingestion_definition": {
"connection_name": "<oracle-connection-name>",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_catalog": "<oracle-service-name>",
"source_schema": "<oracle-schema>",
"source_table": "<oracle-table>"
}
}
]
}
}'
POST /api/2.0/pipelines
{
"name": "my-oracle-integrated-pipeline",
"channel": "PREVIEW",
"catalog": "main",
"schema": "ingestion",
"ingestion_definition": {
"connection_name": "my-oracle-connection",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_catalog": "ORCL",
"source_schema": "HR",
"source_table": "EMPLOYEES",
"table_configuration": {
"scd_type": "SCD_TYPE_1"
}
}
}
],
"data_staging_options": {
"catalog_name": "main",
"schema_name": "ingestion_staging"
}
}
}
ソーススキーマ内のすべてのテーブルをレプリケートするには、個別のtableオブジェクトではなく、schemaオブジェクトを使用します。
POST /api/2.0/pipelines
{
"name": "my-oracle-schema-pipeline",
"channel": "PREVIEW",
"catalog": "main",
"schema": "ingestion",
"ingestion_definition": {
"connection_name": "my-oracle-connection",
"connector_type": "CDC",
"objects": [
{
"schema": {
"source_catalog": "ORCL",
"source_schema": "HR",
"destination_catalog": "main",
"destination_schema": "ingestion"
}
}
]
}
}
パイプラインの更新を開始します:
POST /api/2.0/pipelines/<pipeline-id>/updates
{
"full_refresh": false
}
定期的な更新をスケジュールする
統合されたCDCパイプラインは、トリガーモードでのみ実行されます。繰り返しスケジュールでデータを取り込むには、パイプラインを実行する Lakeflow Jobs タスクを作成します。各更新には約30分かかり、1回の更新で変更のバックログ全体を処理しきれない場合があります。後続の更新が追いつくように、十分な頻度でパイプラインをスケジュール設定してください。60分という開始点は、ほとんどのワークロードに適しています。
構成リファレンス
パイプラインパラメーター
パラメーター | Type | 説明 |
|---|---|---|
| string | パイプラインの名前。 |
| string |
|
| Boolean | オプション。デフォルトは |
| string | デフォルトの宛先カタログです。 |
| string | デフォルトの宛先スキーマです。 |
| string | OracleへのUnity Catalog接続です。 |
| string |
|
| array | 取り込むテーブルまたはスキーマのリスト。 |
| オブジェクト | オプション。パイプラインがステージング ボリュームを作成するカタログとスキーマ。パイプラインの宛先スキーマにデフォルト設定されます。 |
テーブルの指定
パラメーター | 必須 | 説明 |
|---|---|---|
| はい | Oracleサービス名です。マルチテナントデータベースの場合、 |
| はい | Oracleスキーマ(通常、テーブルの所有者)です。 |
| はい | Oracleテーブル名。 |
| No | 宛先カタログ。デフォルトはパイプラインの |
| No | アップグレード先スキーマです。デフォルトはパイプラインの |
| No | 宛先テーブル名です。デフォルトは |
テーブル構成
パラメーター | デフォルト | 説明 |
|---|---|---|
| 自動検出 | 各行を識別する列指定されていない場合は、ソース主キーから自動的に検出されます。 |
|
|
|
| 自動検出 | CDC イベントの論理順序付けに使用される列。 |
Oracleのデータ型マッピングについては、データ型マッピングを参照してください。
Oracle 識別子の大文字と小文字の区別
Oracleは引用符なしの識別子を大文字で格納します。パイプライン構成で source_catalog、source_schema、source_table、および primary_keys を指定する場合、大文字と小文字の区別は Oracle が識別子を格納する方法と一致している必要があります。ほとんどのデータベースでは、大文字にすることになります。識別子が二重引用符で作成され、異なる大文字/小文字が保持されている場合は、その正確な大文字/小文字を使用してください。
パイプラインを監視してください
統合CDCパイプラインを作成して開始した後、以下の方法でそのステータスを監視します。
-
Databricks UI。 「 パイプライン 」セクションでパイプラインを開くと、更新ステータス、テーブルごとの取り込みメトリクス、およびリネージを表示できます。
-
REST API。
TextGET /api/2.0/pipelines/<pipeline-id> -
イベントAPI。
TextGET /api/2.0/pipelines/<pipeline-id>/events
パイプライン詳細ページのリストビューには、データの取り込み時に処理されたレコードの数が表示されます。これらの数値は自動的に更新されます。

最初のパイプライン更新では、選択したすべてのテーブルの完全なスナップショットが実行されます。これには増分更新よりも時間がかかる場合があります。大きなテーブルの場合、初期スナップショットの完了には複数回のスケジュールされた更新が必要となる場合があります。
Unity Catalog に取り込まれたデータをクエリーできます。

フル更新および自動フル更新の動作については、ターゲットテーブルを完全に更新をご覧ください。
統合型CDCパイプラインでは、垂直オートスケールがデフォルトで有効になっています。メモリ不足が原因でパイプラインの更新が失敗した場合、次回の更新でより大きなドライバーが自動的にプロビジョニングされます。
制限事項:
一般的な制限事項
- ベータ版。 統合CDCコネクターとOracleコネクターには、ワークスペースレベルでの有効化が必要です。Databricks アカウント チームにお問い合わせください。
- トリガーモードのみです。 統合CDCパイプラインは、連続的(常時稼働)な実行には対応していません。Lakeflowジョブのタスクを使用したパイプラインのスケジュール
- チャンネルは
PREVIEWである必要があります。 プログラムによるパイプライン仕様には"channel": "PREVIEW"を含める必要があります。 - インジェストパイプラインあたりの推奨最大テーブル数は約500です。
- 統合型CDCパイプラインはスキーマ変更(DDL操作)をまだサポートしていません。
- 大きなテーブルの場合、 初期スナップショット は複数の更新にまたがる可能性があります。
- 各更新は、約 30 分間ランします。 パイプラインは、必ずしも 1 回の更新ですべての変更バックログを処理するわけではありません。後続のスケジュールされた更新は、前回の更新が中断したところから処理を再開します。このランタイムは構成できません。
- パイプライン作成後は、接続とコネクタのタイプは変更できません。
Oracle固有の制限事項
- **サポートされていないOracleのデプロイメント**:Oracle RAC、RAC構成のExadata、Physical Standby、Oracle Autonomous Databases、およびマルチテナントのAmazon RDSデータベースインスタンス。
- サポートされていないデータ型 :
XML、JSON、および空間データ型です。 - LogMinerが無視するテーブル :LogMiner
BFILEは、 、ネストされたテーブル、ID 列、期間有効性列、PKREF列、またはPKOID列を含むすべてのテーブルを無視します。LogMiner 制限事項を参照してください。 - 識別子の長さ :テーブル名および列名は30文字を超過することはできません。
- Post-12.2機能:コネクタは、Oracle Database 12c Release 2以降に追加されたデータ型や機能(
BOOLEANVECTOR、 、および を含む)をサポートしていません。JSON
トラブルシューティング
パイプラインの更新が失敗した場合:
- Databricks UI または
GET /api/2.0/pipelines/<pipeline-id>/eventsを使用してパイプラインイベントログを確認してください。 - カタログ エクスプローラーから Unity Catalog 接続をテストして、Oracle に到達できることを確認します。
- アーカイブログモードと補足ロギングが有効になっていることを確認します。ステップ 1: アーカイブLogsモードとLogs保持の確認を参照してください。
- レプリケーションユーザーが
DBX_ORACLE_SETUP_UTIL.GRANT_PERMISSIONSによって付与された権限を持っていることを確認してください。「Oracle データベースユーザーの要件」を参照してください。 - マルチテナントデータベースの場合、ユーザーが
CDB$ROOTで共通ユーザーであること、また、source_catalogがCDB$ROOTサービス名であることを確認してください。 - パイプライン仕様に
"channel": "PREVIEW"が含まれていることを確認してください。
Oracle がパイプラインがそれらを処理できる前にアーカイブログをパージする場合、影響を受けるテーブルで完全更新を実行してください。