はじめに
Oracle AI Data Catalog は Oracle が提供する Iceberg REST catalog サービスで、Oracle Autonomous AI Database 内で動作します。
Iceberg REST Catalog は、Apache Iceberg プロジェクトが公式に定義する OpenAPI 仕様で、これに準拠すれば、どのエンジンからも同じ REST クライアントで任意の準拠カタログに接続できる、というのが設計思想です。
ということで、この記事では、Apache Spark で Oracle AI Data Catalog を使いながら OCI Object Storage 上の Iceberg テーブルを操作してみようと思います。
Apache Spark と AI Data Catalog との役割分担
Spark は、
- Parquet データファイルの物理書き込み (Object Storage へ)
- manifest / manifest-list / 新しい metadata.json の生成
- AI Data Catalog に commit リクエスト
を行います。
Spark はデータとメタデータファイルを作成しますが、"どれが公式の最新版か"を決定するのは AI Data Catalog です。Catalog が commit を受理して初めて、その書き込みが他エンジンから見えるようになります。
AI Data Catalog は。
- テーブルの登録・発見・一覧
- どの metadata.json が最新かの確定 (commit)
- 楽観ロック (並行 commit の衝突検出)
AI Data Catalog は Spark が生成した metadata を受け取ってアトミック(原始的)にポインタを進める役割を担っています。
AI Data Catalog の作成
こちらの記事の最初の部分に、作成手順が詳しく書いてありますので、参照して下さい。
カタログ用のバケット ai-data-catalog を事前に準備しておきます。
begin
oracle_ai_data_catalog.register_storage_oci(
p_warehouse => 's3://ai-data-catalog',
p_endpoint => 'https://{{namespace}}.compat.objectstorage.{region}.oci.customer-oci.com',
p_region => '{{region}}',
p_access_key => 'xxxxxxxxxxxx',
p_secret_key => 'xxxxxxxxxxxx'
);
end;
/
ここまで準備できたら データベース側の作業は完了です。
念のため、カタログの状態を確認しておきましょう (python)。
import requests, json
SERVER_URI = "https://xxx.adb.us-ashburn-1.oraclecloudapps.com/catalog"
CLIENT_ID = "xxxxx"
CLIENT_PW = "xxxxx"
# get access token
response = requests.post(
url = f"{SERVER_URI}/v1/auth/token",
headers = {
"Content-Type": "application/x-www-form-urlencoded",
},
data = {
"grant_type": "client_credentials",
"client_id": CLIENT_ID,
"client_secret": CLIENT_PW,
"scope": "PRINCIPAL_ROLE:ALL",
},
verify=False
)
response.raise_for_status()
token = response.json()["access_token"]
# get catalog config
response = requests.get(
url = f"{SERVER_URI}/v1/config",
headers = {
"Authorization": f"Bearer {token}",
},
verify = False
)
response.raise_for_status()
print(json.dumps(response.json()["defaults"], indent=2))
出力結果
{
"default-base-location": "s3://ai-data-catalog"
}
Apache Spark の準備
Apache Spark インストール
この記事を参考にして Apache Spark をインストールして下さい。
(Spark 3.5.6 / Scala 2.12 / Hadoop 3.3)
VS Code Notebook 環境の準備
今回は VS Code 上の Notebook で色々と試してみたいと思います。
Notebook のカーネルに使う Python の仮想環境を作成して、pyspark と findspark をインストールしておいて下さい。Spark 3.5.6 がサポートするのは Python 3.8 以上です。
(pyspark のバージョンは 3.5.6 に合わせて下さい)
$ pip list | grep spark
findspark 2.0.1
pyspark 3.5.6
Apache Spark で Iceberg テーブルを操作する
基本事項の確認: catalog, namespace, schema, database, table
- 3階層の識別子:
catalog.namespace.table(例:aicat.testns.testtab) - schema と database は同義語 です
- namespaceは Iceberg/Spark Catalog API での呼び方
- つまり namespace = schema = database
上記をまずはしっかりと把握しておきましょう。
Spark セッションの開始
では、Notebook から操作を開始しましょう。
$\small \textsf{script}$
import findspark
findspark.init() # pyspark を Notebook で動かすための準備
from pyspark.sql import SparkSession
cat_server_uri = "https://xxxxx.adb.us-ashburn-1.oraclecloudapps.com/catalog"
oauth2_server_uri = f"{cat_server_uri}/v1/auth/token"
access_key_id = "xxxxx"
secret_access_key = "xxxxx"
cat_db_user = "xxxxx"
cat_db_pw = "xxxxx"
spark = (
SparkSession.builder
.appName("MyIcebergApp")
.config("spark.jars.packages", "org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.11.0,org.apache.iceberg:iceberg-aws-bundle:1.11.0")
.config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
# work around AWS SDK v2 compatibility problem
.config("spark.driver.extraJavaOptions",
"-Daws.requestChecksumCalculation=when_required -Daws.responseChecksumValidation=when_required") \
.config("spark.executor.extraJavaOptions",
"-Daws.requestChecksumCalculation=when_required -Daws.responseChecksumValidation=when_required") \
# カタログ 'aicat' の構成
.config("spark.sql.catalog.aicat", "org.apache.iceberg.spark.SparkCatalog")
.config("spark.sql.catalog.aicat.type", "rest")
.config("spark.sql.catalog.aicat.uri", cat_server_uri)
.config("spark.sql.catalog.aicat.s3.access-key-id", access_key_id)
.config("spark.sql.catalog.aicat.s3.secret-access-key", secret_access_key)
.config("spark.sql.catalog.aicat.rest.auth.type", "oauth2")
.config("spark.sql.catalog.aicat.oauth2-server-uri", oauth2_server_uri)
.config("spark.sql.catalog.aicat.credential", f"{cat_db_user}:{cat_db_pw}")
.config("spark.sql.catalog.aicat.scope", "PRINCIPAL_ROLE:ALL")
.getOrCreate()
)
ポイント
- セッションの config でカタログ 'aicat' を定義しています
- AI Data Catalog はファイルの操作に S3互換API を使います
- カタログの操作に必要な認証 (OAuth2) のための設定を行なっています
- 追加パッケージとして
iceberg-spark-runtimeとiceberg-aws-bundleを指定しています (Spark, Scala, Iceberg のバージョンを合わせるように!)
次に catalog と namespace (= database) の状況を確認します。
$\small \textsf{script}$
spark.sql(f"select current_catalog(), current_schema()").show()
spark.sql(f"use aicat")
spark.sql(f"select current_catalog(), current_schema()").show()
Spark は初期状態では spark_catalog というデフォルトの catalog があり、その中に default という namespace (= database) があります。このまま完全修飾せずにテーブルの操作を続けると、この空間で作業してしまいます。
"use aicat" で catalog を切り替えましたが、この catalog には まだ何も namespace がありません。
testns という namespace を作成します。
$\small \textsf{script}$
spark.sql("create namespace if not exists testns")
spark.sql("use testns")
spark.sql(f"select current_catalog(), current_schema()").show()
Iceberg テーブルの操作
まずテストデータを作ります。
$\small \textsf{script}$
from pyspark.sql import functions as F
df = spark.range(3).withColumn("data", F.concat(F.lit("data_"), F.col("id")))
df.show()
新規にテーブルを作成して、データを書き込みます。
$\small \textsf{script}$
table_name = "aicat.testns.testtab"
# 最初にテーブルを作成してから追加する
ddl_str = df._jdf.schema().toDDL()
sql = f"create table if not exists {table_name} ({ddl_str}) using iceberg"
print(sql)
spark.sql(sql) # create table
df.writeTo(table_name).append() # append
以下の方法でテーブルを作成しないでください
# これをやると metadata.json に不整合が生じる
df.writeTo("aicat.testns.testtable").create()
CREATE TABLE で明示的にテーブルを作ってから .append() しましょう。
作成したテーブルを見てみましょう。
$\small \textsf{script}$
spark.table(table_name).printSchema()
spark.sql(f"select * from {table_name}").show()
テーブルの詳細情報も確認しておきましょう。
$\small \textsf{script}$
spark.sql(f"describe table extended {table_name}").show(truncate=False)
次に、テーブルを更新してみます。
$\small \textsf{script}$
spark.sql(f"update {table_name} set data = 'data_X' where id = 0")
spark.sql(f"select * from {table_name}").show()
最後に、AI Data Catalog にアクセスして、管理されているテーブル情報を確認してみましょう。
url = f"{SERVER_URI}/v1/namespaces/testns/tables/testtab"
response = requests.get(url, headers={"Authorization": f"Bearer {token}"}, verify=False)
print(json.dumps(response.json(), indent=2))
$\small \textsf{出力結果}$
{
"metadata-location": "s3://ai-data-catalog/testns/testtab/metadata/00002-833ab224-37b9-4078-bf53-4d97ab1043dd.metadata.json",
"metadata": {
"format-version": 2,
"table-uuid": "03d62593-5a24-43ff-88a2-9ad521bd97cc",
"location": "s3://ai-data-catalog/testns/testtab",
"last-sequence-number": 2,
"last-updated-ms": 1783928760639,
"last-column-id": 2,
"current-schema-id": 0,
"schemas": [
{
"type": "struct",
"schema-id": 0,
"fields": [
{
"id": 1,
"name": "id",
"required": true,
"type": "long"
},
{
"id": 2,
"name": "data",
"required": true,
"type": "string"
}
]
}
],
"default-spec-id": 0,
"partition-specs": [
{
"spec-id": 0,
"fields": []
}
],
"last-partition-id": 999,
"default-sort-order-id": 0,
"sort-orders": [
{
"order-id": 0,
"fields": []
}
],
"properties": {
"owner": "opc",
"created-at": "2026-07-12T02:21:56.255220851Z",
"write.parquet.compression-codec": "zstd"
},
"current-snapshot-id": 5790970030471432630,
"refs": {
"main": {
"snapshot-id": 5790970030471432630,
"type": "branch"
}
},
"snapshots": [
{
"sequence-number": 1,
"snapshot-id": 7761524200876133388,
"timestamp-ms": 1783823489940,
"summary": {
"operation": "append",
"spark.app.id": "local-1783823446279",
"manifests-created": "1",
"manifests-kept": "0",
"manifests-replaced": "0",
"added-data-files": "3",
"added-records": "3",
"added-files-size": "2029",
"changed-partition-count": "1",
"total-records": "3",
"total-files-size": "2029",
"total-data-files": "3",
"total-delete-files": "0",
"total-position-deletes": "0",
"total-equality-deletes": "0",
"engine-version": "3.5.6",
"app-id": "local-1783823446279",
"engine-name": "spark",
"iceberg-version": "Apache Iceberg 1.11.0 (commit 6976e020b894f6a6777704df2b8c4458cb291ae9)",
"app-name": "MyIcebergApp"
},
"manifest-list": "s3://ai-data-catalog/testns/testtab/metadata/snap-7761524200876133388-1-a70f7b62-b9ed-4a9f-b828-ebfc581afa81.avro",
"schema-id": 0
},
{
"sequence-number": 2,
"snapshot-id": 5790970030471432630,
"parent-snapshot-id": 7761524200876133388,
"timestamp-ms": 1783928760639,
"summary": {
"operation": "overwrite",
"spark.app.id": "local-1783928415574",
"manifests-created": "2",
"manifests-kept": "0",
"manifests-replaced": "1",
"added-data-files": "1",
"deleted-data-files": "1",
"added-records": "1",
"deleted-records": "1",
"added-files-size": "675",
"removed-files-size": "675",
"changed-partition-count": "1",
"total-records": "3",
"total-files-size": "2029",
"total-data-files": "3",
"total-delete-files": "0",
"total-position-deletes": "0",
"total-equality-deletes": "0",
"engine-version": "3.5.6",
"app-id": "local-1783928415574",
"engine-name": "spark",
"iceberg-version": "Apache Iceberg 1.11.0 (commit 6976e020b894f6a6777704df2b8c4458cb291ae9)",
"app-name": "MyIcebergApp"
},
"manifest-list": "s3://ai-data-catalog/testns/testtab/metadata/snap-5790970030471432630-1-a61cefa9-f2fc-4ea0-8b16-244b94598694.avro",
"schema-id": 0
}
],
"statistics": [],
"partition-statistics": [],
"snapshot-log": [
{
"timestamp-ms": 1783823489940,
"snapshot-id": 7761524200876133388
},
{
"timestamp-ms": 1783928760639,
"snapshot-id": 5790970030471432630
}
],
"metadata-log": [
{
"timestamp-ms": 1783822916291,
"metadata-file": "s3://ai-data-catalog/testns/testtab/metadata/00000-068f1a63-61e0-4ee9-9eca-887b348a73a7.metadata.json"
},
{
"timestamp-ms": 1783823489940,
"metadata-file": "s3://ai-data-catalog/testns/testtab/metadata/00001-b9d860b0-2b7b-4ace-a34c-9a6c966ccb8b.metadata.json"
}
]
},
"config": {
"s3.path-style-access": "true",
"s3.endpoint": "https://xxxxxx.compat.objectstorage.us-ashburn-1.oci.customer-oci.com",
"io-impl": "org.apache.iceberg.aws.s3.S3FileIO",
"client.region": "us-ashburn-1"
}
}
yaml 形式にして、主要部分だけ抜き出すとこんな感じです。
metadata-location: s3://ai-data-catalog/testns/testtab/metadata/00002-833ab224-37b9-4078-bf53-4d97ab1043dd.metadata.json
metadata:
location: s3://ai-data-catalog/testns/testtab
last-sequence-number: 2
last-updated-ms: 1783928760639
current-snapshot-id: 5790970030471432630
refs:
main:
snapshot-id: 5790970030471432630
type: branch
snapshots:
- sequence-number: 1
snapshot-id: 7761524200876133388
timestamp-ms: 1783823489940
manifest-list: s3://ai-data-catalog/testns/testtab/metadata/snap-7761524200876133388-1-a70f7b62-b9ed-4a9f-b828-ebfc581afa81.avro
- sequence-number: 2
snapshot-id: 5790970030471432630
parent-snapshot-id: 7761524200876133388
timestamp-ms: 1783928760639
manifest-list: s3://ai-data-catalog/testns/testtab/metadata/snap-5790970030471432630-1-a61cefa9-f2fc-4ea0-8b16-244b94598694.avro
snapshot-log:
- timestamp-ms: 1783823489940
snapshot-id: 7761524200876133388
- timestamp-ms: 1783928760639
snapshot-id: 5790970030471432630
metadata-log:
- timestamp-ms: 1783822916291
metadata-file: s3://ai-data-catalog/testns/testtab/metadata/00000-068f1a63-61e0-4ee9-9eca-887b348a73a7.metadata.json
- timestamp-ms: 1783823489940
metadata-file: s3://ai-data-catalog/testns/testtab/metadata/00001-b9d860b0-2b7b-4ace-a34c-9a6c966ccb8b.metadata.json
config:
s3.endpoint: https://xxxxxxx.compat.objectstorage.us-ashburn-1.oci.customer-oci.com
snapshot が2つあることがわかりますね。
まとめ
前回の Iceberg の操作を扱った記事では Hadoop Catalog を使っていました。
Hadoop Catalog では、「現在の metadata.json はどれか」を version-hint.text や metadata ディレクトリ内のファイル名のバージョン番号で管理し、コミットは新しいmetadata.json をアトミックに rename することで実現する設計となっています。これは POSIX ファイルシステムや HDFS のアトミック rename が前提です。
ところが S3 や今回の OCI Object Storage のような S3 互換ストレージには、アトミック rename が存在しません。そのため Hadoop Catalog を S3 上で使うと、複数クライアントが同時書き込みした際のコミットの競合検出が弱くなり、更新のロスト(lost update)が起きるリスクがあります。
REST Catalog は、コミットの整合性検証をカタログサーバー側でトランザクショナルに行うため、ストレージの特性に依存せず、安全に楽観的ロックがかけられます。Iceberg の特長を活かしながら複数のクライアントからトランザクショナルに Object Storage 上のテーブルの読み書きしたい場合に AI Data Catalog を有効活用することができます。
何よりも Autonomous AI Database がこの AI Data Catalog を活用する訳ですけどね 😀
参考
※ Spark の設定に関して、ドキュメンテーション記載の設定ではうまく接続できなかったので、Try & Error して接続確認した結果がこの記事の設定内容となっています。







