はじめに
Apache Icebergの特徴(複数エンジンでのテーブル共有、ACIDトランザクション、タイムトラベルなど)は、解説記事でよく見かけます。ただ、文字で見るだけでは中々その利点も腹落ちしずらいところ。
この記事では、OCI上に実機を組んで、Icebergの代表的な特徴をすべて"目で見て"確かめてみます。1つのIcebergテーブルを ADB・Spark・Trinoの3エンジン で共有し、Sparkが更新している最中に1秒ごとに読み続け、最後は過去のデータに戻ってみます。
今回確かめるIcebergの4つの特徴
先に、この記事で実機確認する特徴を整理しておきます。
Icebergに触れたことがない方は、ここを押さえておくと本編を追いやすくなります。
| # | 特徴 | ひとことで言うと | 確かめる場所 |
|---|---|---|---|
| ① | オープンなテーブル共有 | 1つのテーブルを複数エンジン(ADB/Spark/Trino)がコピーなしで読み書きできる | 実験1 |
| ② | スナップショット分離 | 更新の最中でも、読む人は常に「完全な状態」だけを見る | 実験2 |
| ③ | 原子的コミット | 切り替わるのはコミットの一瞬。どのエンジンから見ても同時 | 実験2 |
| ④ | タイムトラベル | 過去の状態にSQL 1行で戻れる | 実験3 |
鍵になるのはIcebergのスナップショットの仕組みです。
Icebergテーブルは、書き込みが確定(コミット)するたびに「その時点の完全な姿=スナップショット」を1つ記録します。読み手は常に「コミット済みのスナップショット」だけを読むので、書きかけのデータは構造的に見えません。そして古いスナップショットは消えずに残るので、過去にも戻れます。
Icebergそのものの入門や特徴、メタデータの内部構造は、過去記事にまとめています。基礎から始める場合は、ぜひ以下の記事も覗いてみてください。
参考:
- Open Table Format(Apache Iceberg)を"まるっと理解する" - AI時代のデータ基盤を考えるための基礎技術- まとめ
- Apache Iceberg入門:誕生の背景から特徴、アーキテクチャまとめ
- Apache Iceberg入門:作られるファイルから理解するメタデータとマニフェスト
では、この4つを順番に実機で見ていきます。
今回の構成とシナリオ
カタログには、Autonomous AI Database(ADB)に内蔵されたIceberg RESTカタログ「Oracle AI Data Catalog(AICAT)」を使います。3エンジンはすべてこのカタログ1本を参照し、実データはObject Storage上に置かれます。
▪️ 利用するサービス
- Autonomous AI Database(AICATタグ付きで作成)
- OCI Data Flow(Spark 3.5)
- Compute(Trino 482をpodmanで起動)
- Object Storage(Icebergウェアハウス用バケット+資材用バケット)
(カタログはIceberg REST Catalogであれば何でも構いません。要件は「3つのエンジンが同じカタログを参照できること」なので、他のRESTカタログに読み替えても話は変わりません。今回はADBを作るだけで用意できるAICATを使っています。)
データはKaggleのOlistデータセット(ブラジルECの実注文データ)から、1週間分の注文 1,071件 を使います。シナリオは「注文はいったん全件 ACCEPTED で確定(=S0)。その後、後着のキャンセル12件+配送遅延71件、計83件を1回のMERGEで反映(=S1)」です。
| 状態 | ACCEPTED | CANCELED | DELIVERY_DELAYED | 合計 |
|---|---|---|---|---|
| S0(受付時点) | 1,071 | 0 | 0 | 1,071 |
| S1(MERGE後) | 988 | 12 | 71 | 1,071 |
実験2では、次の問いを確かめます。
MERGE中に読み続けたとき、この2つ以外の数字(一部だけ更新された中途半端な状態)が1回でも見えるか?
前提条件
-
AICATが有効なADBが作成済みであること(フリーフォームタグ
ADB$TOOLS=AI_CATを付けて作成し、ウェアハウス用バケットを登録)。手順はこちら → ADBからIcebergレイクハウスを扱えます!Oracle AI Data Catalog(AICAT)を試してみた -
OCI Data FlowでSparkアプリを動かせること → OCIでApache Iceberg入門:Object StorageとData Flowで動かしてみる
-
Object StorageのS3互換APIで使う顧客秘密キーを作成済みであること → 【OCI】顧客秘密キーの作成手順
-
Trinoを動かすCompute(本記事ではpodmanを使用)
実験1: 特徴①を見る — ADBで作ったテーブルを、Trinoがそのまま読めるか
1-1. ADBでIcebergテーブルを作ってS0をコミット
Olistの注文CSVは、Database Actionsの「データ・ロード」でヒープ表 OLIST_ORDERS に取り込んでおきます。そこからIcebergテーブルを作成します。ADBのSQLだけで、Object Storage上にParquetとIcebergメタデータが直接書かれます。
-- テーブル作成
CREATE ICEBERG TABLE "demo"."orders_demo" (
order_id STRING,
order_date TIMESTAMP,
order_state STRING
)
WITHIN CATALOG "AICAT"
STORAGE LOCATION "s3://<ウェアハウス用バケット>/demo/orders_demo/";
次にデータを挿入していきます。
-- データ挿入
INSERT INTO "demo"."orders_demo"@AICAT
SELECT order_id,
order_purchase_timestamp,
'ACCEPTED'
FROM admin.olist_orders
WHERE order_purchase_timestamp >= TIMESTAMP '2018-08-20 00:00:00'
AND order_purchase_timestamp < TIMESTAMP '2018-08-27 00:00:00';
COMMIT;
1,071 rows inserted. と表示されればS0のコミット完了です。
注意点:
STORAGE LOCATIONは二重引用符で囲みます。単一引用符だと s3:// がADB内部で https:// 形式に変換され、not in allowed locations エラーになります。
ADBのIcebergテーブルへの書き込みは INSERT INTO … SELECT(+COMMIT)のみ対応で、COMMITのたびにスナップショットが1つ積まれます。集計して確認しておきます。
SELECT SUM(CASE WHEN order_state='ACCEPTED' THEN 1 ELSE 0 END) accepted,
SUM(CASE WHEN order_state='CANCELED' THEN 1 ELSE 0 END) canceled,
SUM(CASE WHEN order_state='DELIVERY_DELAYED' THEN 1 ELSE 0 END) delayed,
COUNT(*) total
FROM "demo"."orders_demo"@AICAT;
1071 / 0 / 0 / 1071 が返ってきました。
カタログとObject Storageの確認
さらにここで、カタログとObject Storageではどのように反映されているかみておきます。
まずAICATのカタログをリフレッシュすると、先ほどまで何もなかったところにテーブルが登録されていることが確認できます。
データも挿入されています。
Object Storageのバケットには、metadata(metadata.json等)とdata(Parquet)の実体ができています。第3章で見たファイル構成が、実際に並んでいるところです。
SQLを実行しただけですが、この時点でObject Storage上にオープンな形式のテーブルが存在し、他のエンジンと共有できる状態になっています。
1-2. Trinoから同じテーブルを読む
次に、まったく別のエンジンであるTrinoから同じテーブルを読みます。エクスポートもコピーもしません。TrinoのIcebergコネクタにAICATをRESTカタログとして登録するだけです。
podman run -d --name trino --userns=keep-id -p 8080:8080 \
-v ~/trino/catalog:/etc/trino/catalog:Z docker.io/trinodb/trino:latest
カタログ定義はこの1ファイルです。
connector.name=iceberg
iceberg.catalog.type=rest
iceberg.rest-catalog.uri=https://<ADBのホスト名>.adb.<リージョン>.oraclecloudapps.com/catalog
iceberg.rest-catalog.security=OAUTH2
iceberg.rest-catalog.oauth2.token=<Bearerトークン>
iceberg.rest-catalog.oauth2.token-refresh-enabled=false
iceberg.rest-catalog.vended-credentials-enabled=false
fs.s3.enabled=true
s3.endpoint=https://<テナンシのnamespace>.compat.objectstorage.<リージョン>.oci.customer-oci.com
s3.region=<リージョン>
s3.path-style-access=true
s3.aws-access-key=<顧客秘密キーのアクセスキー>
s3.aws-secret-key=<顧客秘密キーのシークレット>
BearerトークンはAICATのトークンエンドポイントから、DBユーザーの資格情報(OAuth2 client_credentials)で取得します。有効期限が60分なので、更新をスクリプトにしておきます。
トークン取得&Trino再起動スクリプトはこちら
#!/bin/bash
# AICATのBearerトークン(60分有効)を取得してTrinoのカタログ定義に反映する
set -euo pipefail
ADB_HOST=<ADBのホスト名>.adb.<リージョン>.oraclecloudapps.com
ADMIN_PASS='<ADMINのパスワード>'
TOKEN=$(curl -s -X POST "https://${ADB_HOST}/catalog/v1/auth/token" \
--data-urlencode "grant_type=client_credentials" \
--data-urlencode "client_id=ADMIN" \
--data-urlencode "client_secret=${ADMIN_PASS}" \
--data-urlencode "scope=PRINCIPAL_ROLE:ALL" | jq -r '.access_token')
sed -i "s|^iceberg.rest-catalog.oauth2.token=.*|iceberg.rest-catalog.oauth2.token=${TOKEN}|" \
~/trino/catalog/aicat.properties
podman restart trino
echo "OK: トークン更新・Trino再起動済み(有効期限: 約60分)"
読んでみます。Trinoでは $snapshots というメタデータ表からスナップショットID(テーブルの"版番号")も見えるので、一緒に取得します。
SELECT
(SELECT snapshot_id FROM aicat.demo."orders_demo$snapshots"
ORDER BY committed_at DESC LIMIT 1) AS snapshot,
count_if(order_state='ACCEPTED') AS accepted,
count_if(order_state='CANCELED') AS canceled,
count_if(order_state='DELIVERY_DELAYED') AS delayed,
count(*) AS total
FROM aicat.demo.orders_demo
snapshot=7025329549986002255 accepted=1071 canceled=0 delayed=0 total=1071
特徴①が確認できました。ADBで作ったテーブルが、コピーなしでTrinoからそのまま読めています。このスナップショットID(70253295…)は、実験3のタイムトラベルでもう一度使います。
実験2: 特徴②③を見る — MERGE中に1秒ごとに読み続ける
次は、更新中の見え方です。書き手役として、Data Flow上のSparkアプリで後着更新83件を1回の MERGE INTO で反映します。
ここで、実験用の工夫を加えます。
実データ1,071行のMERGEは一瞬で終わってしまい、「書き込み中に読む」時間が取れません。そこで、1行処理するごとに数百ミリ秒待つUDFをUPDATE句に仕込み、書き込みを4並列(4ファイル)に分割しています。
これでMERGE処理がまるごと約2分かかるようになります。
※「約2分」はSparkの起動やcommit処理も含めた、Runが始まってからの実測です。なお、この遅延は実験用の演出ですが、スナップショット分離の保証自体は遅延の有無と無関係に成立します。
# 1行ごとに指定ミリ秒待つUDF(書き込みを遅くする実験用の仕掛け)
def slow_pass(v):
time.sleep(row_delay_ms / 1000.0)
return v
spark.udf.register("slow_pass", slow_pass, "string")
spark.sql("""
MERGE INTO aicat.demo.orders_demo t
USING updates s
ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET t.order_state = slow_pass(s.new_state)
""")
名前:iceberg-demo-merge
Spark構成プロパティ:
spark.sql.catalog.aicat:org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.aicat.type:rest
spark.sql.catalog.aicat.uri:https://<ADBのホスト名>.adb.<リージョン>.oraclecloudapps.com/catalog
spark.sql.catalog.aicat.credential:ADMIN:<ADMINのパスワード>
spark.sql.catalog.aicat.s3.access-key-id:<顧客秘密キーのアクセスキー>
spark.sql.catalog.aicat.s3.secret-access-key:<顧客秘密キーのシークレット>
spark.sql.shuffle.partitions:4(4ファイル分割書き込み用)
spark.sql.adaptive.coalescePartitions.enabled:false
spark.executorEnv.AWS_REQUEST_CHECKSUM_CALCULATION:when_required(ハマりポイント2参照)
spark.executorEnv.AWS_RESPONSE_CHECKSUM_VALIDATION:when_required
引数:--run-id ${run_id} --row-delay-ms 300
Sparkアプリ全文はこちら ※長いです(ADBが書いたマニフェストの修復処理込み。後述のハマりポイント1参照)
#!/usr/bin/env python3
"""Data Flow App 1: Iceberg MERGE (Step 4).
処理:
1) 対象table現snapshot(S0)のmanifest/manifest list を検査し、ADB(AICAT)ライター特有の
スキーマ属性欠落(field-id / element-id / logicalType:map)をその場で修復する。
同一パス上書きで、データ・snapshot IDは不変(標準リーダーが読めるようになるだけ)。
2) update manifest(83件)を読み、1回のMERGE INTO(update-only)でS1をcommitする。
3) S1の親=S0、件数988/12/71/1071を検証する。
引数は --run-id (と任意の --pre-merge-wait) のみ。接続情報・秘密情報・パス規約は
すべてData Flow ApplicationのSpark構成プロパティから読む:
spark.sql.catalog.aicat.uri / .credential / .s3.access-key-id / .s3.secret-access-key
spark.demo.assets (例 oci://iceberg-demo-asset@<namespace>)
spark.demo.s3.endpoint (customer-oci.com形式)
spark.demo.region
update manifestの場所は <spark.demo.assets>/runs/<run_id>/update_manifest.csv 規約。
"""
import argparse
import io
import json
import sys
import time
import urllib.parse
import urllib.request
from pyspark.sql import SparkSession
def log(msg):
print(f"[merge_app] {msg}", flush=True)
# ====================================================================
# manifest修復ロジック (scripts/repair_manifests.py と同一)
# ====================================================================
MANIFEST_FILE_FIELD_IDS = {
"manifest_path": 500, "manifest_length": 501, "partition_spec_id": 502,
"content": 517, "sequence_number": 515, "min_sequence_number": 516,
"added_snapshot_id": 503,
"added_data_files_count": 504, "added_files_count": 504,
"existing_data_files_count": 505, "existing_files_count": 505,
"deleted_data_files_count": 506, "deleted_files_count": 506,
"added_rows_count": 512, "existing_rows_count": 513, "deleted_rows_count": 514,
"partitions": 507, "key_metadata": 519,
}
R508_FIELD_IDS = {"contains_null": 509, "contains_nan": 518, "lower_bound": 510, "upper_bound": 511}
MANIFEST_ENTRY_FIELD_IDS = {
"status": 0, "snapshot_id": 1, "sequence_number": 3, "file_sequence_number": 4, "data_file": 2,
}
DATA_FILE_FIELD_IDS = {
"content": 134, "file_path": 100, "file_format": 101, "partition": 102,
"record_count": 103, "file_size_in_bytes": 104,
"column_sizes": 108, "value_counts": 109, "null_value_counts": 110,
"nan_value_counts": 137, "lower_bounds": 125, "upper_bounds": 128,
"key_metadata": 131, "split_offsets": 132, "equality_ids": 135,
"sort_order_id": 140, "referenced_data_file": 143,
}
DATA_FILE_MAPS = {
"column_sizes": (117, 118), "value_counts": (119, 120), "null_value_counts": (121, 122),
"nan_value_counts": (138, 139), "lower_bounds": (126, 127), "upper_bounds": (129, 130),
}
DATA_FILE_ARRAYS = {"split_offsets": 133, "equality_ids": 136}
def _nonnull_branch(t):
if isinstance(t, list):
for b in t:
bt = b.get("type") if isinstance(b, dict) else b
if bt != "null":
return b
return None
return t
def _schema_has_field_ids(schema):
if isinstance(schema, dict) and schema.get("type") == "record":
return any("field-id" in f for f in schema.get("fields", []))
return False
def _fix_manifest_list_schema(schema):
for f in schema["fields"]:
name = f["name"]
if name in MANIFEST_FILE_FIELD_IDS:
f["field-id"] = MANIFEST_FILE_FIELD_IDS[name]
if name == "partitions":
arr = _nonnull_branch(f["type"])
if isinstance(arr, dict) and arr.get("type") == "array":
arr["element-id"] = 508
items = arr["items"]
if isinstance(items, dict) and items.get("type") == "record":
for sf in items["fields"]:
if sf["name"] in R508_FIELD_IDS:
sf["field-id"] = R508_FIELD_IDS[sf["name"]]
return schema
def _fix_manifest_entry_schema(schema):
for f in schema["fields"]:
name = f["name"]
if name in MANIFEST_ENTRY_FIELD_IDS:
f["field-id"] = MANIFEST_ENTRY_FIELD_IDS[name]
if name == "data_file":
df = f["type"] if isinstance(f["type"], dict) else _nonnull_branch(f["type"])
for sf in df["fields"]:
sn = sf["name"]
if sn in DATA_FILE_FIELD_IDS:
sf["field-id"] = DATA_FILE_FIELD_IDS[sn]
if sn in DATA_FILE_MAPS:
arr = _nonnull_branch(sf["type"])
if isinstance(arr, dict) and arr.get("type") == "array":
arr["logicalType"] = "map"
kid, vid = DATA_FILE_MAPS[sn]
items = arr["items"]
if isinstance(items, dict) and items.get("type") == "record":
for kv in items["fields"]:
if kv["name"] == "key":
kv["field-id"] = kid
elif kv["name"] == "value":
kv["field-id"] = vid
if sn in DATA_FILE_ARRAYS:
arr = _nonnull_branch(sf["type"])
if isinstance(arr, dict) and arr.get("type") == "array":
arr["element-id"] = DATA_FILE_ARRAYS[sn]
if sn == "partition":
prec = sf["type"] if isinstance(sf["type"], dict) else _nonnull_branch(sf["type"])
if isinstance(prec, dict) and prec.get("fields"):
raise SystemExit("partitioned table repair not implemented")
return schema
def _rewrite_avro(raw, fixer):
import fastavro
src = io.BytesIO(raw)
reader = fastavro.reader(src)
schema = reader.writer_schema
if _schema_has_field_ids(schema):
return None, False
records = list(reader)
meta = {k: v for k, v in (reader.metadata or {}).items() if not k.startswith("avro.")}
fixed_schema = fixer(json.loads(json.dumps(schema)))
out = io.BytesIO()
fastavro.writer(out, fixed_schema, records, codec=reader.codec or "null", metadata=meta)
return out.getvalue(), True
def _with_retry(what, fn, attempts=6, base_wait=5):
"""ゲートウェイの一過性エラー(401/5xx/接続断)に耐えるリトライ。"""
for i in range(1, attempts + 1):
try:
return fn()
except Exception as e: # noqa: BLE001
if i == attempts:
raise
wait = min(base_wait * (2 ** (i - 1)), 60)
log(f"{what} failed (try {i}/{attempts}): {e} — retrying in {wait}s")
time.sleep(wait)
def catalog_token(host, user, pw):
def _fetch():
data = urllib.parse.urlencode({
"grant_type": "client_credentials", "client_id": user,
"client_secret": pw, "scope": "PRINCIPAL_ROLE:ALL"}).encode()
req = urllib.request.Request(f"https://{host}/catalog/v1/auth/token", data=data, method="POST")
with urllib.request.urlopen(req, timeout=60) as r:
return json.loads(r.read())["access_token"]
return _with_retry("catalog_token", _fetch)
def load_table_metadata(host, token, namespace, table):
def _fetch():
req = urllib.request.Request(
f"https://{host}/catalog/v1/namespaces/{namespace}/tables/{table}",
headers={"Authorization": f"Bearer {token}"})
with urllib.request.urlopen(req, timeout=60) as r:
return json.loads(r.read())["metadata"]
return _with_retry("load_table_metadata", _fetch)
def repair_snapshot_manifests(cfg, snapshot_id, ml_uri):
"""指定snapshotのmanifest list/manifestsを修復(S3のみ使用、RESTは呼ばない)"""
import fastavro as fa
import boto3
from botocore.config import Config
s3 = boto3.client("s3", region_name=cfg["region"], endpoint_url=cfg["s3_endpoint"],
aws_access_key_id=cfg["access_key"], aws_secret_access_key=cfg["secret_key"],
config=Config(s3={"addressing_style": "path"},
request_checksum_calculation="when_required",
response_checksum_validation="when_required"))
def s3_split(uri):
u = urllib.parse.urlparse(uri)
return u.netloc, u.path.lstrip("/")
for snap_id, ml in [(snapshot_id, ml_uri)]:
bucket, key = s3_split(ml)
log(f"repair: snapshot {snap_id} manifest-list {ml}")
ml_raw = s3.get_object(Bucket=bucket, Key=key)["Body"].read()
ml_reader = fa.reader(io.BytesIO(ml_raw))
ml_needs_schema_fix = not _schema_has_field_ids(ml_reader.writer_schema)
entries = list(ml_reader)
length_changed = False
for e in entries:
mb, mk = s3_split(e["manifest_path"])
m_raw = s3.get_object(Bucket=mb, Key=mk)["Body"].read()
fixed, did = _rewrite_avro(m_raw, _fix_manifest_entry_schema)
new_len = len(fixed) if did else len(m_raw)
if did:
s3.put_object(Bucket=mb, Key=mk, Body=fixed)
log(f"repair: fixed manifest {mk} ({len(m_raw)} -> {new_len})")
if e["manifest_length"] != new_len:
e["manifest_length"] = new_len
length_changed = True
if ml_needs_schema_fix or length_changed:
meta = {k: v for k, v in (ml_reader.metadata or {}).items() if not k.startswith("avro.")}
schema = json.loads(json.dumps(ml_reader.writer_schema))
if ml_needs_schema_fix:
schema = _fix_manifest_list_schema(schema)
out = io.BytesIO()
fa.writer(out, schema, entries, codec=ml_reader.codec or "null", metadata=meta)
s3.put_object(Bucket=bucket, Key=key, Body=out.getvalue())
log(f"repair: rewrote manifest list ({len(ml_raw)} -> {out.getbuffer().nbytes})")
else:
log("repair: nothing to do (already standard)")
def main():
import os
os.environ.setdefault("AWS_REQUEST_CHECKSUM_CALCULATION", "when_required")
os.environ.setdefault("AWS_RESPONSE_CHECKSUM_VALIDATION", "when_required")
p = argparse.ArgumentParser()
p.add_argument("--run-id", required=True)
p.add_argument("--pre-merge-wait", type=int, default=10)
p.add_argument("--row-delay-ms", type=int, default=300,
help="更新83行あたりの書込み時遅延(UPDATE SET式で評価)。"
"書込みステージが約1.5〜2分の実トランザクションになり、"
"『書込みトランザクション進行中もreaderはS0だけを見る』を観測可能にする")
args = p.parse_args()
spark = SparkSession.builder.appName(f"iceberg-merge-{args.run_id}").getOrCreate()
# OCI S3互換APIのaws-chunked非対応への対策(ドライバーJVM)。
# extraJavaOptionsでの指定はData Flow環境の既定JVM設定を上書きしてしまうため、
# 実行時にSystemプロパティとして設定する(S3クライアントは遅延初期化なので有効)。
# エグゼキュータ側はApp構成の spark.executorEnv.AWS_* で対応済み。
jsys = spark._jvm.java.lang.System
jsys.setProperty("aws.requestChecksumCalculation", "WHEN_REQUIRED")
jsys.setProperty("aws.responseChecksumValidation", "WHEN_REQUIRED")
conf = spark.sparkContext.getConf()
cat_uri = conf.get("spark.sql.catalog.aicat.uri") # https://<host>/catalog
cred = conf.get("spark.sql.catalog.aicat.credential") # user:pass
host = urllib.parse.urlparse(cat_uri).netloc
db_user, db_pass = cred.split(":", 1)
cfg = {
"host": host,
"db_user": db_user,
"db_pass": db_pass,
"namespace": "demo",
"table": f"orders_{args.run_id.lower()}",
"s3_endpoint": conf.get("spark.demo.s3.endpoint"),
"region": conf.get("spark.demo.region"),
"access_key": conf.get("spark.sql.catalog.aicat.s3.access-key-id"),
"secret_key": conf.get("spark.sql.catalog.aicat.s3.secret-access-key"),
}
assets = conf.get("spark.demo.assets").rstrip("/")
manifest_uri = f"{assets}/runs/{args.run_id}/update_manifest.csv"
full = f'aicat.`{cfg["namespace"]}`.`{cfg["table"]}`'
log(f"target={full} manifest={manifest_uri}")
# 診断のみ(非致死): Data Flow環境からのトークンエンドポイント到達性を記録
for ua in ("Python-urllib/3.11", "curl/8.4.0"):
try:
data = urllib.parse.urlencode({
"grant_type": "client_credentials", "client_id": cfg["db_user"],
"client_secret": cfg["db_pass"], "scope": "PRINCIPAL_ROLE:ALL"}).encode()
req = urllib.request.Request(
f"https://{cfg['host']}/catalog/v1/auth/token", data=data, method="POST",
headers={"User-Agent": ua})
with urllib.request.urlopen(req, timeout=30) as r:
log(f"diag: token endpoint UA='{ua}' -> HTTP {r.status}")
except Exception as e: # noqa: BLE001
body = ""
try:
body = e.read().decode()[:200] if hasattr(e, "read") else ""
except Exception: # noqa: BLE001
pass
log(f"diag: token endpoint UA='{ua}' -> {e} body={body}")
# 1) S0の特定とmanifestパス解決はSparkカタログ経由(metadata.jsonのみ参照、
# データ/manifestは読まないため未修復でも安全)。
# 初回アクセスはカタログ初期化(OAuthトークン取得)を含むため、ゲートウェイの
# 断続的な401(経路依存で観測)に備えてリトライする。
snaps = _with_retry("catalog first access (snapshots)", lambda: spark.sql(
f"SELECT snapshot_id, manifest_list FROM {full}.`snapshots` ORDER BY committed_at DESC"
).collect(), attempts=6, base_wait=10)
if len(snaps) == 0:
raise SystemExit("S0がまだ作られていない(snapshotが存在しない)")
if len(snaps) != 1:
raise SystemExit(f"snapshotが{len(snaps)}個ある — S0のみの状態でMERGEを実行すること")
s0_id, ml_uri = snaps[0]["snapshot_id"], snaps[0]["manifest_list"]
# 2) 修復 (Sparkがtableの実データを読む前に実施)
repair_snapshot_manifests(cfg, s0_id, ml_uri)
log(f"repair done. S0 = {s0_id}")
# 2) update manifest検証
updates = spark.read.option("header", True).csv(manifest_uri)
n = updates.count()
states = {r["new_state"]: r["c"] for r in
updates.groupBy("new_state").count().withColumnRenamed("count", "c").collect()}
log(f"manifest rows={n} breakdown={states}")
if n != 83 or states.get("CANCELED") != 12 or states.get("DELIVERY_DELAYED") != 71:
raise SystemExit(f"manifest validation failed: rows={n} breakdown={states}")
if updates.select("order_id").distinct().count() != 83:
raise SystemExit("manifest validation failed: order_id not unique")
updates.createOrReplaceTempView("updates")
pre = {r["order_state"]: r["c"] for r in
spark.sql(f"SELECT order_state, count(*) c FROM {full} GROUP BY order_state").collect()}
log(f"pre-merge counts: {pre}")
if pre != {"ACCEPTED": 1071}:
raise SystemExit(f"pre-merge state is not S0: {pre}")
log(f"waiting {args.pre_merge_wait}s before MERGE (readers keep seeing S0)...")
time.sleep(args.pre_merge_wait)
# 3) 1回のMERGE (update-only) = 1 commit
# ソース行にパススルーUDFで遅延を入れ、join+書込みフェーズを本物の
# 長時間トランザクションにする(commitは従来どおり最後に1回だけ)。
# broadcast joinを無効化し、UDF評価が書込みステージ内で行われるようにする。
from pyspark.sql import functions as F
from pyspark.sql import types as T
delay_s = args.row_delay_ms / 1000.0
def _slow_pass(x):
time.sleep(delay_s)
return x
# MERGEはjoin条件に非決定的式を許さないため、UDFは決定的として登録する
# (joinキーを包むパススルーなので、決定的でも行ごとの実行は省略されない)
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")
spark.udf.register("slow_pass", _slow_pass, T.StringType())
# 書込みを4ファイルに分割し、書込みの早い段階からdata fileがstorageに着地するようにする
# (=「ファイルは存在するのにreaderはS0のまま」の観測窓を数秒→数十秒に広げる)。
# AQEの小パーティション併合と、Icebergのhash分散(1タスク化)を無効にする。
spark.conf.set("spark.sql.shuffle.partitions", "4")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "false")
spark.sql(f"ALTER TABLE {full} SET TBLPROPERTIES('write.merge.distribution-mode'='none')")
# 遅延はUPDATE SET式に入れる: shuffleの下流=書込みステージの投影で行ごとに
# 評価されるため、data fileの書込みと遅延が確実に交錯する
# (ソース側に入れるとshuffle上流のスキャン段で消化され、書込みは一瞬で終わる)。
est = 83 * delay_s / 2 # 4タスク×2並列
log(f"executing single MERGE INTO (update-only), write phase slowed by "
f"{args.row_delay_ms}ms/updated-row (~{est:.0f}s of real write-stage time)...")
t_start = time.time()
spark.sql(f"""
MERGE INTO {full} t
USING (SELECT order_id, new_state FROM updates) s
ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET t.order_state = slow_pass(s.new_state)
""")
log(f"MERGE committed. (transaction wall time: {time.time() - t_start:.1f}s)")
# 4) 検証
post = {r["order_state"]: r["c"] for r in
spark.sql(f"SELECT order_state, count(*) c FROM {full} GROUP BY order_state").collect()}
total = sum(post.values())
expected = {"ACCEPTED": 988, "CANCELED": 12, "DELIVERY_DELAYED": 71}
log(f"post-merge: {post} total={total}")
if post != expected or total != 1071:
raise SystemExit(f"post-merge validation failed: {post} (expected {expected})")
# S1と親子関係の検証もSparkカタログ経由(RESTは呼ばない)
rows = spark.sql(
f"SELECT snapshot_id, parent_id FROM {full}.`snapshots` ORDER BY committed_at DESC"
).collect()
if len(rows) != 2:
raise SystemExit(f"snapshot数が2ではない: {len(rows)}")
cur, parent = rows[0]["snapshot_id"], rows[0]["parent_id"]
log(f"S1 = {cur}, parent = {parent}")
if parent != s0_id:
raise SystemExit(f"S1 parent mismatch: parent={parent} expected S0={s0_id}")
log("SUCCESS: S1 committed as single snapshot with parent S0 (988/12/71/1071)")
spark.stop()
if __name__ == "__main__":
main()
OCI Data Flowの基本的な使い方は以下の記事で扱っていますので参照してください。
- OCIでApache Iceberg入門:Object StorageとData Flowで動かしてみる
- 【ハンズオン】OCI Data Flowで始めるApache Spark|ETLからMLまで体験してみよう!
読み手側は、1秒ごとに「現在のスナップショットID+状態別の件数+Object Storage上のデータファイル数」を1行で表示し続けるループを回します。「ストレージに新しいファイルが増えていく様子」と「読者に見える結果」を同時に観察するためです。
読み取りループ(reader_loop.sh)はこちら
#!/bin/bash
# Trino reader loop (Step 4用): current snapshotの集計と、Object Storage上の
# data file数を1行で繰り返し表示する。
#
# 表示の工夫:
# - 最初に観測したsnapshotを「S0:」、commit後の新snapshotを「S1:」とラベル表示
# - snapshotが切り替わった行の直前に「◀◀◀ COMMIT」マーカー行を自動挿入
# (切替の一瞬が録画で見逃されないように)
#
# 見せ場: MERGEの書込み中は snapshot=S0 のまま data-files が増えていき
# (=storage上に未commitのdata fileが存在するのに読まれない)、
# COMMITマーカーの次の行から S1(988/12/71) に一斉に切り替わる。
#
# 件数の組として表示されてよいのは S0(1071/0/0) と S1(988/12/71) の2状態のみ。
# 出力は logs/reader_<table>.log にも自動保存される(step4_verify.shが参照)。
#
# 使い方: ./reader_loop.sh <table> [interval_sec]
# 例: ./reader_loop.sh orders_rec1 1
TABLE="${1:?usage: reader_loop.sh <table> [interval_sec]}"
INTERVAL="${2:-1}"
BASE=~/iceberg-demo
LOG="$BASE/logs/reader_${TABLE}.log"
: > "$LOG"
export SUPPRESS_LABEL_WARNING=True
SQL="SELECT
(SELECT snapshot_id FROM aicat.demo.\"${TABLE}\$snapshots\" ORDER BY committed_at DESC LIMIT 1),
count_if(order_state='ACCEPTED'),
count_if(order_state='CANCELED'),
count_if(order_state='DELIVERY_DELAYED'),
count(*)
FROM aicat.demo.${TABLE}"
datafiles() {
local n
n=$(oci os object list -ns <namespace> -bn <ウェアハウス用バケット> \
--prefix "demo/${TABLE}/data/" --query 'length(data)' --output json 2>/dev/null)
[ -z "$n" ] && n=0
echo "$n"
}
TMPQ=$(mktemp "$BASE/tmp/reader_q_XXXXXX")
trap 'rm -f "$TMPQ"' EXIT
FIRST_SNAP="" # 最初に観測したsnapshot = S0
PREV_SNAP="" # 直前の行のsnapshot(切替検知用)
GEN=0 # 世代番号(S0, S1, S2…)
while true; do
# Trinoクエリとstorage一覧を並行実行して1周期を短く保つ
podman exec trino trino --execute "$SQL" 2>&1 | head -1 > "$TMPQ" &
QPID=$!
DF=$(datafiles)
wait "$QPID" 2>/dev/null
OUT=$(cat "$TMPQ")
if [[ "$OUT" == \"* ]]; then
IFS=',' read -r SNAP ACC CAN DEL TOT <<< "$(echo "$OUT" | tr -d '"')"
if [ -z "$FIRST_SNAP" ]; then
FIRST_SNAP="$SNAP"; PREV_SNAP="$SNAP"; GEN=0
elif [ "$SNAP" != "$PREV_SNAP" ]; then
GEN=$((GEN+1))
MARKER=$(printf '%s ◀◀◀ COMMIT — snapshotが切り替わりました (S%s → S%s)' \
"$(date '+%H:%M:%S')" "$((GEN-1))" "$GEN")
echo "$MARKER" | tee -a "$LOG"
PREV_SNAP="$SNAP"
fi
LINE=$(printf '%s snapshot=S%s:%s accepted=%s canceled=%s delayed=%s total=%s | data-files=%s' \
"$(date '+%H:%M:%S')" "$GEN" "$SNAP" "$ACC" "$CAN" "$DEL" "$TOT" "$DF")
else
LINE=$(printf '%s ERROR: %s' "$(date '+%H:%M:%S')" "$OUT")
fi
echo "$LINE" | tee -a "$LOG"
sleep "$INTERVAL"
done
ループを流したままData Flowコンソールからiceberg-demo-mergeをRunします。
※Data Flowの初回Runは起動に6分ほどかかりました。2回目以降は1〜2分です。
実験結果
MERGEが動いている間、ループの出力はこうなりました(見やすさのため、実際のログから一部の行を抜粋しています)。
09:28:39 snapshot=S0:7025329549986002255 accepted=1071 canceled=0 delayed=0 total=1071 | data-files=3
09:28:52 snapshot=S0:7025329549986002255 accepted=1071 canceled=0 delayed=0 total=1071 | data-files=3
09:29:02 snapshot=S0:7025329549986002255 accepted=1071 canceled=0 delayed=0 total=1071 | data-files=4
09:29:13 snapshot=S0:7025329549986002255 accepted=1071 canceled=0 delayed=0 total=1071 | data-files=4
09:29:16 ◀◀◀ COMMIT — snapshotが切り替わりました (S0 → S1)
09:29:16 snapshot=S1:345263046526741820 accepted=988 canceled=12 delayed=71 total=1071 | data-files=5
09:29:18 snapshot=S1:345263046526741820 accepted=988 canceled=12 delayed=71 total=1071 | data-files=5
実験ログから、次の3点が確認できます。
-
data-filesが1→2→3→4と増えている間も、読み手の結果はS0のまま1行も揺れません。MERGEが書いた新しいParquetファイルはObject Storage上に物理的に存在するのに、読み手からは見えていないということです。書き込み途中のファイルはまだどのスナップショットにも属していない(どのマニフェストからも参照されていない)ので、読みようがありません。これが特徴②スナップショット分離の正体です - コミットの一瞬でS1(988/12/71)へ切り替わります。1050/3/18のような「一部だけ更新された中途半端な数字」は一度も観測されませんでした。別の検証走行でRunの開始から終了まで観測し続けたときも、読み取り324回の内訳は「S0が290回、S1が33回、それ以外0回・エラー0回」でした
- コミット後は
data-files=5(新4+旧1)になります。旧ファイルが消されずに残っている点は、実験3でまた触れます
MERGEの最中に、ADBのDatabase Actionsからも同じ集計SQLを投げてみました。結果は 1071 / 0 / 0 / 1071。Trinoだけでなく、どのエンジンから見てもS0です。切り替わるのは、カタログ上のコミットの一瞬だけです。特徴③(原子的コミット)も確認できました! 書き込み中でもSELECTは待たされず即返ってきます(読み書きが互いにブロックしません)。
事後検証として、新しいデータファイルの書き込みが完了してからコミットされるまでの94秒間の読み取り35回を集計したところ、いずれもS0でした(S1・混在は0回)。
実験3: 特徴④を見る — タイムトラベルで更新前のデータに戻る
コミット前のデータは消えていません。まず、テーブル自身が持っている履歴を見てみます。
SELECT snapshot_id, parent_id, operation, committed_at
FROM aicat.demo."orders_demo$snapshots" ORDER BY committed_at;
7025329549986002255 (null) append ... ← S0
345263046526741820 7025329549986002255 overwrite ... ← S1(親=S0)
S1の parent_id がS0を指しています。「S1はS0から生まれた」という系譜が、テーブル自身のメタデータとして残っています。
実験1の集計SQLに、FOR VERSION AS OF を1行加えます。
SELECT count_if(order_state='ACCEPTED') AS accepted,
count_if(order_state='CANCELED') AS canceled,
count_if(order_state='DELIVERY_DELAYED') AS delayed,
count(*) AS total
FROM aicat.demo.orders_demo
FOR VERSION AS OF 7025329549986002255 -- ← 追加したのはこの1行(S0のID)
accepted=1071 canceled=0 delayed=0 total=1071
MERGE前のデータがそのまま返ってきました。バックアップからのリストアも、コピーの保管もしていません。実験2の「コミット後も旧データファイルが残る(data-files=5)」の伏線が、ここで回収されました。
個別の注文でも見てみます。MERGEでキャンセルになった注文を、現在とS0で並べると:
SELECT 'S1 (current)' AS label, order_id, order_state
FROM aicat.demo.orders_demo
WHERE order_id = '0623dbbd06c0d4059dbfefab8748a491'
UNION ALL
SELECT 'S0 (time travel)', order_id, order_state
FROM aicat.demo.orders_demo FOR VERSION AS OF 7025329549986002255
WHERE order_id = '0623dbbd06c0d4059dbfefab8748a491';
S1 (current) 0623dbbd… CANCELED
S0 (time travel) 0623dbbd… ACCEPTED
同じ注文の「過去の姿」が見えます。特徴④(タイムトラベル)も確認できました。 問い合わせ対応や監査で「あのとき、どうだったか」に正確に答えられるということです。
ハマりポイントまとめ
すんなり動いたわけではありません。踏んだ罠を残しておきます。
| # | ハマったこと | 原因と対処 |
|---|---|---|
| 1 | SparkのMERGEが Not a list type: map<int,long> 等で失敗 |
ADB(AICAT)が書いたAvroマニフェストに field-id 等の標準スキーマ属性が一部欠落していた。アプリ冒頭でマニフェストを標準形に修復(同一パス上書き。データもスナップショットIDも不変)してからMERGEすることで回避(アプリ全文参照) |
| 2 | SparkからObject StorageへのPUTが501エラー | OCIのS3互換APIが aws-chunked 転送を拒否する。aws.requestChecksumCalculation=WHEN_REQUIRED を指定して解決 |
| 3 | Trinoが約1時間後に Failed to load table |
AICATのBearerトークンが60分で失効。トークンの自動リフレッシュは unauthorized_client で使えず、静的トークン+スクリプト再取得の運用に |
| 4 | Trinoで一部のWHERE句が Error processing metadata で失敗 |
ADBが書いたマニフェストの列統計とTrinoの述語プッシュダウンの相性問題らしく、フィルタが全データファイルを誤って除外することがある。同じ等値フィルタでも動く列(本文の WHERE order_id = … は正常)と失敗する列があった。GROUP BYのみの集計と FOR VERSION AS OF は常に正常。失敗する場合は `列 |
| 5 | STORAGE LOCATIONで権限エラー | 単一引用符だと s3:// がADB内部でhttps形式に変換される。二重引用符が必須 |
| 6 | 遅延UDFがMERGEの前に実行されてしまう | UDFを検索側(USING句)に置くとシャッフル前に評価される。UPDATE SET句側に置くことで書き込みステージ内で評価されるように |
「1つのテーブルを複数エンジンで共有できる」オープンさの裏で、エンジン間のメタデータ互換にはまだ発展途上な部分がありました。一方で、回避策を入れた後の読み書きは安定しており、3エンジン間で切り替えの一貫性が崩れることはありませんでした。
おわりに
Icebergの4つの特徴、
① 複数エンジンでのテーブル共有
② スナップショット分離
③ 原子的コミット
④ タイムトラベル
を、いずれも実機のログで確認できました。新しいファイルがストレージに増えていく間も読み取り結果が1行も揺れない様子は、実機で動かすことで確認できます。ぜひ試してみてください。
次回は、このIcebergテーブルの上にベクトル検索と生成AIを載せてみます。本文をObject Storageに置いたまま、ADBにはベクトルと索引だけを持たせる構成です。









