4
2

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

Oracle AI Data Catalog を使って Apache Spark から Iceberg テーブルを操作する

4
Last updated at Posted at 2026-07-13

はじめに

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 の作成

こちらの記事の最初の部分に、作成手順が詳しく書いてありますので、参照して下さい。

Oracle AI データ・カタログの有効化

ストレージ資格証明の登録

カタログ用のバケット 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;
/

ネットワークACLを設定する

ここまで準備できたら データベース側の作業は完了です。

念のため、カタログの状態を確認しておきましょう (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 の仮想環境を作成して、pysparkfindspark をインストールしておいて下さい。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 から操作を開始しましょう。

image.png

$\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-runtimeiceberg-aws-bundle を指定しています (Spark, Scala, Iceberg のバージョンを合わせるように!)

次に catalog と namespace (= database) の状況を確認します。

image.png

$\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 を作成します。

image.png

$\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 テーブルの操作

まずテストデータを作ります。

image.png

$\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()

新規にテーブルを作成して、データを書き込みます。

image.png

$\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() しましょう。

作成したテーブルを見てみましょう。

image.png

$\small \textsf{script}$
spark.table(table_name).printSchema()
spark.sql(f"select * from {table_name}").show()

テーブルの詳細情報も確認しておきましょう。

image.png

$\small \textsf{script}$
spark.sql(f"describe table extended {table_name}").show(truncate=False)

次に、テーブルを更新してみます。

image.png

$\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 して接続確認した結果がこの記事の設定内容となっています。

4
2
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
4
2

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?