はじめに
概要
Icebergを扱えるクライアントは多数あります(Spark, Trino, DuckDB, ...)。今回はその中からpyiceberg(Python), Spark, DuckDBを取り上げます。
Icebergはカスタマイズ性が非常に高い反面、クライアントから接続する際の設定値の数がとても多いし複雑です。本記事はそのうちの一部をまとめてみたものになります。
今回のIcebergの技術スタック
今回はRESTカタログ+S3ストレージでIceberg環境を用意します。
RESTカタログはApache Iceberg、S3ストレージはCephを使用しています。互換性があるほかの技術(Ozone, Minio, ...)を使っても設定は大きく変わらないはずです。
今回の作成物
今回作成したコード等はGitHub上に保管されています。記事内のコードは一部分を切り取っていたりするので、ハンズオン的に動かしたい場合などはローカルに一度保存頂くとよいかもしれません。
前準備
環境構築
たくさんの環境をつくらないといけないので今回はDocker上に環境構築を行います。
以前docker-composeでCeph環境を作る記事を書いたので、こちらに付け足す感じで進めます。
ここに、Polaris, Python, Spark, DuckDBのコンテナを加えます。
# ~~~~~~~~~~~ceph用コード~~~~~~~~~~~
polaris:
image: apache/polaris:1.3.0-incubating
ports:
# API port
- "8181:8181"
depends_on:
rgw1:
condition: service_started
environment:
AWS_REGION: ${S3_REGION}
AWS_ACCESS_KEY_ID: ${RGW_ACCESS_KEY}
AWS_SECRET_ACCESS_KEY: ${RGW_SECRET_KEY}
POLARIS_BOOTSTRAP_CREDENTIALS: POLARIS,root,s3cr3t
polaris.realm-context.realms: POLARIS
quarkus.otel.sdk.disabled: "true"
healthcheck:
test: ["CMD", "curl", "--fail", "http://localhost:8182/q/health"]
interval: 2s
timeout: 10s
retries: 10
start_period: 10s
networks:
ceph_network:
python:
build: dockerfiles/python
tty: true
volumes:
- ./mounts/python:/work
networks:
ceph_network:
spark:
build: dockerfiles/spark
user: spark
command: /opt/spark/bin/spark-class org.apache.spark.deploy.master.Master
networks:
ceph_network:
duckdb:
image: duckdb/duckdb:1.4.4
tty: true
networks:
ceph_network:
Dockerfile達は以下の通りです。
Pythonコンテナは必要なpipライブラリをインストールします。pyicebergは使用するストレージによって若干コマンドが変わりますが、今回はS3互換のストレージを使うので、s3fsとします。
FROM python:3.14
RUN apt-get update && apt-get -y upgrade
COPY ./requirements.txt .
RUN pip install -U pip
RUN pip install -r requirements.txt
pyiceberg[s3fs]==0.10.0
pandas
pyarrow
SparkイメージではSpark関係のコマンドをパスに通しておきます。Pysparkも使うためPYTHONPATH変数も定義しておきます。
FROM apache/spark:4.1.1-scala2.13-java21-python3-r-ubuntu
USER root
RUN apt-get update && apt-get -y upgrade
ENV PATH=${SPARK_HOME}/bin:$PATH
ENV PYTHONPATH="${SPARK_HOME}/python:${SPARK_HOME}/python/lib/py4j-0.10.9.9-src.zip:${PYTHONPATH}"
RUN mkdir /nonexistent
RUN chown spark /nonexistent
カタログ作成
クライアントで操作する前に、Polaris Catalog REST APIを使って新しいカタログを生成します。
まず、AWS-CLIを使ってバケットを作成しておきます。
$ aws s3 mb s3://test-bucket
ここからはPolarisのAPIを使います。
Polaris APIはDockerが入っているホストPCあるいは、Polarisコンテナ内からたたくので、ホスト名がコンテナ名ではなくlocalhostになっています。他コンテナからたたく場合はpolarisに書き換えてください。
まずは、Bootstrap用のアカウントでOauth2認証します。
レスポンスのうち、TOKENの値は変数に入れておきます。
$ curl http://localhost:8181/api/catalog/v1/oauth/tokens \
--user root:s3cr3t \
-H "Polaris-Realm: POLARIS" \
-d grant_type=client_credentials \
-d scope=PRINCIPAL_ROLE:ALL
カタログを作成します。
body='{
"catalog": {
"name": "test_catalog",
"type": "INTERNAL",
"readOnly": false,
"properties": {
"default-base-location": "s3://test-bucket"
},
"storageConfigInfo": {
"storageType":"S3",
"allowedLocations": ["s3://test-bucket"],
"endpoint":"http://rgw1:7480",
"endpointInternal":"http://rgw1:7480",
"stsUnavailable": true,
"kmsUnavailable": true,
"pathStyleAccess": true
}
}
}'
curl http://localhost:8181/api/management/v1/catalogs \
-H "Authorization: Bearer $TOKEN" \
-H "Accept: application/json" \
-H "Content-Type: application/json" \
-H "Polaris-Realm: POLARIS" \
-d "$body"
これでカタログができました。
クライアントで操作してみる
ここからは作成したカタログを使って、以下の操作をしてみます。
- ネームスペースの作成
- テーブルの作成
- データの挿入
- データの閲覧
- テーブルとネームスペースの削除
実際に各種スクリプトやコマンドを試してみる場合は、以下のコマンドで各種コンテナに入ってから実行してください。
$ docker exec -it {コンテナID} bash
pyiceberg
pyicebergはカタログ接続に必要な値の設定方法が3種類あります。
There are three ways to pass in configuration:
- Using the .pyiceberg.yaml configuration file (Recommended)
- Through environment variables
- By passing in credentials through the CLI or the Python API
今回は1つ目と3つ目の方法を試してみます。
本記事の作成時点(2026/02/28)のpyicebergの最新バージョンは0.11.0ですが、このバージョンはRESTカタログのHttp操作に関するバグを含んでいるみたいです。(PR-3010)
なので1つ古い0.10.0を使用しています。
.pyiceberg.yamlに設定値を記載する方法
まず、Pythonスクリプトがあるディレクトリに.pyiceberg.yamlを作成します。
catalog:
test_catalog:
# Storage
s3.endpoint: http://rgw1:7480
s3.access-key-id: POLARIS123ACCESS
s3.secret-access-key: POLARIS456SECRET
# OAuth2
auth:
type: oauth2
oauth2:
client_id: root
client_secret: s3cr3t
token_url: http://polaris:8181/api/catalog/v1/oauth/tokens
scope: PRINCIPAL_ROLE:ALL
# Catalog
type: rest
uri: http://polaris:8181/api/catalog
warehouse: test_catalog
header.X-Iceberg-Access-Delegation: ""
以下のスクリプトを実行すると、同ディレクトリ内にある.pyiceberg.yamlを読み込んで、カタログに接続してくれます。
import pyarrow as pa
from pyiceberg.catalog import load_catalog
from pyiceberg.expressions import GreaterThanOrEqual
from pyiceberg.schema import Schema
from pyiceberg.types import IntegerType, LongType, NestedField, StringType
def main():
# Reads all settings from .pyiceberg.yaml automatically
catalog = load_catalog("test_catalog")
# write
catalog.create_namespace_if_not_exists("public")
table = catalog.create_table_if_not_exists(
identifier="public.raw",
schema=Schema(
NestedField(field_id=1, name="id", field_type=LongType(), required=False),
NestedField(field_id=2, name="name", field_type=StringType(), required=False),
NestedField(field_id=3, name="age", field_type=IntegerType(), required=False),
),
)
records = pa.table({
"id": pa.array([1, 2, 3], type=pa.int64()),
"name": pa.array(["Alice", "Bob", "Charlie"], type=pa.string()),
"age": pa.array([30, 25, 35], type=pa.int32()),
})
table.append(records)
# read
table = catalog.load_table("public.raw")
print(table.scan(row_filter=GreaterThanOrEqual("age", 30)).to_pandas())
# delete
catalog.drop_table("public.raw")
catalog.drop_namespace("public")
if __name__ == "__main__":
main()
.pyiceberg.yamlには複数のカタログを設定することができます、1つ目のネストに追記してあげれば大丈夫です。ファイル内のどの設定を使うかは、load_catalog()関数の第一引数の文字列を参照しています。
Pythonスクリプト内で設定値を直接指定する場合
.pyiceberg.yamlではなく、Pythonスクリプトで直接設定値を記載することもできます。ユーザーから入力を受け取って設定したかったりする場合は便利そうです。
1つめのスクリプトと大差はなく、load_catalog()へ可変長引数として設定値を渡してあげればよいです。本スクリプトではdict形で入れた後に展開しています。
import pyarrow as pa
from pyiceberg.catalog import load_catalog
from pyiceberg.expressions import GreaterThanOrEqual
from pyiceberg.schema import Schema
from pyiceberg.types import IntegerType, LongType, NestedField, StringType
def main():
catalog = load_catalog(
"test_catalog",
**{
# Storage
"s3.endpoint": "http://rgw1:7480",
"s3.access-key-id": "POLARIS123ACCESS",
"s3.secret-access-key": "POLARIS456SECRET",
# Oauth2 (new AuthManager — pass nested dict to match YAML structure)
"auth": {
"type": "oauth2",
"oauth2": {
"client_id": "root",
"client_secret": "s3cr3t",
"token_url": "http://polaris:8181/api/catalog/v1/oauth/tokens",
"scope": "PRINCIPAL_ROLE:ALL",
},
},
# Catalog
"type": "rest",
"uri": "http://polaris:8181/api/catalog",
"warehouse": "test_catalog",
"header.X-Iceberg-Access-Delegation": "",
}
)
# write
catalog.create_namespace_if_not_exists("public")
table = catalog.create_table_if_not_exists(
identifier="public.raw",
schema=Schema(
NestedField(field_id=1, name="id", field_type=LongType(), required=False),
NestedField(field_id=2, name="name", field_type=StringType(), required=False),
NestedField(field_id=3, name="age", field_type=IntegerType(), required=False),
),
)
records = pa.table({
"id": pa.array([1, 2, 3], type=pa.int64()),
"name": pa.array(["Alice", "Bob", "Charlie"], type=pa.string()),
"age": pa.array([30, 25, 35], type=pa.int32()),
})
table.append(records)
# read
table = catalog.load_table("public.raw")
print(table.scan(row_filter=GreaterThanOrEqual("age", 30)).to_pandas())
# delete
catalog.drop_table("public.raw")
catalog.drop_namespace("public")
if __name__ == "__main__":
main()
load_catalog()は可変長引数で設定を定義しても、同ディレクトリに.pyiceberg.yamlがあるとそちらの設定も読み込んでマージしてしまいます。片方の設定に寄せたい場合はディレクトリを分けましょう。なお、環境変数設定と可変長引数設定はマージしないっぽいです。
Spark
Sparkは様々なクライアントがありますが(sql, r, python, ...)、今回はsqlとpythonを使ってみます。
といっても、設定値の渡し方が異なるだけでどちらも接続さえしてしまえば同一のsqlステートメントで操作ができます。
spark-sql
spark-sqlコマンド時に設定値を渡してあげます。
spark-sql \
--packages org.apache.iceberg:iceberg-spark-runtime-4.0_2.13:1.10.1,org.apache.iceberg:iceberg-aws-bundle:1.10.1 \
--conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
--conf spark.sql.catalog.test_catalog=org.apache.iceberg.spark.SparkCatalog \
--conf spark.sql.catalog.test_catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO \
--conf spark.sql.catalog.test_catalog.s3.endpoint=http://rgw1:7480 \
--conf spark.sql.catalog.test_catalog.s3.access-key-id=POLARIS123ACCESS \
--conf spark.sql.catalog.test_catalog.s3.secret-access-key=POLARIS456SECRET \
--conf spark.sql.catalog.test_catalog.rest.auth.type=oauth2 \
--conf spark.sql.catalog.test_catalog.credential=root:s3cr3t \
--conf spark.sql.catalog.test_catalog.oauth2-server-uri=http://polaris:8181/api/catalog/v1/oauth/tokens \
--conf spark.sql.catalog.test_catalog.scope=PRINCIPAL_ROLE:ALL \
--conf spark.sql.catalog.test_catalog.token-refresh-enabled=true \
--conf spark.sql.catalog.test_catalog.type=rest \
--conf spark.sql.catalog.test_catalog.uri=http://polaris:8181/api/catalog \
--conf spark.sql.catalog.test_catalog.warehouse=test_catalog \
--conf spark.sql.catalog.test_catalog.client.region=irrelevant
client.region=irrelevantの設定はローカルのs3を使っているため、適当な文字列で動くだけでAWS-S3を使う場合は適切な文字列を設定してあげてください。またs3.reagionという設定も必要になるかもしれません。
あとはSQLステートメントで操作できます。
-- create
CREATE NAMESPACE IF NOT EXISTS test_catalog.public;
CREATE TABLE IF NOT EXISTS test_catalog.public.raw (
id BIGINT,
name STRING,
age INT
)
USING iceberg;
-- write
INSERT INTO test_catalog.public.raw VALUES
(1, 'Alice', 30),
(2, 'Bob', 25),
(3, 'Charlie', 35);
-- read (age >= 30)
SELECT * FROM test_catalog.public.raw
WHERE age >= 30;
-- delete
DROP TABLE IF EXISTS test_catalog.public.raw;
DROP NAMESPACE IF EXISTS test_catalog.public;
pyspark
SparkSessionの作成時に設定値を渡してあげます。その後はspark.sql()関数を使って、spark-sqlと同じ操作をしてみます。
from pyspark.sql import SparkSession
spark = (SparkSession.builder
.config("spark.jars.packages", "org.apache.iceberg:iceberg-spark-runtime-4.0_2.13:1.10.1,org.apache.iceberg:iceberg-aws-bundle:1.10.1")
.config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
.config("spark.sql.catalog.test_catalog", "org.apache.iceberg.spark.SparkCatalog")
# S3
.config("spark.sql.catalog.test_catalog.io-impl", "org.apache.iceberg.aws.s3.S3FileIO")
.config("spark.sql.catalog.test_catalog.client.region", "irrelevant")
.config("spark.sql.catalog.test_catalog.s3.access-key-id", "POLARIS123ACCESS")
.config("spark.sql.catalog.test_catalog.s3.secret-access-key", "POLARIS456SECRET")
.config("spark.sql.catalog.test_catalog.s3.endpoint", "http://rgw1:7480")
# oauth2
.config("spark.sql.catalog.test_catalog.rest.auth.type", "oauth2")
.config("spark.sql.catalog.test_catalog.credential", "root:s3cr3t")
.config("spark.sql.catalog.test_catalog.oauth2-server-uri", "http://polaris:8181/api/catalog/v1/oauth/tokens")
.config("spark.sql.catalog.test_catalog.scope", 'PRINCIPAL_ROLE:ALL')
.config("spark.sql.catalog.test_catalog.token-refresh-enabled", "true")
# catalog
.config("spark.sql.catalog.test_catalog.type", "rest")
.config("spark.sql.catalog.test_catalog.uri", "http://polaris:8181/api/catalog")
.config("spark.sql.catalog.test_catalog.warehouse", "test_catalog")
).getOrCreate()
spark.sql("CREATE NAMESPACE IF NOT EXISTS test_catalog.public;")
spark.sql(
"CREATE TABLE IF NOT EXISTS test_catalog.public.raw ("
"id BIGINT,"
"name STRING,"
"age INT"
") USING iceberg;"
)
spark.sql(
"INSERT INTO test_catalog.public.raw VALUES "
"(1, 'Alice', 30),"
"(2, 'Bob', 25),"
"(3, 'Charlie', 35);"
)
spark.sql(
"SELECT * FROM test_catalog.public.raw "
"WHERE age >= 30;"
).show()
spark.sql("DROP TABLE IF EXISTS test_catalog.public.raw;")
spark.sql("DROP NAMESPACE IF EXISTS test_catalog.public;")
ちなみに、pysparkをつかったスクリプトの実行方法は以下の二通がありますが、
$ python3 use_pyspark.py
$ spark-submit use_pyspark.py
上記のスクリプトは前者のコマンドでないと動きません。
原因がわからなくて色々調べたのですが、後者のコマンドはspark.jars.packages設定でMavenから持ってくることができないみたいです。回避するにはspark-submit --packagesコマンドで回避するか、$SPARK_HOME/jarsフォルダに自分でファイルを設置するかです。
DuckDB
DuckDBはSQLステートメントで各種設定を記載します。Iceberg拡張機能のインストールもSQLステートメントでかけます。spark-sqlとはちょっと勝手が違います。
INSTALL aws;
INSTALL httpfs;
INSTALL iceberg;
LOAD iceberg;
LOAD httpfs;
-- oauth2
CREATE SECRET token (
TYPE iceberg,
CLIENT_ID 'root',
CLIENT_SECRET 's3cr3t',
OAUTH2_SERVER_URI 'http://polaris:8181/api/catalog/v1/oauth/tokens'
);
-- s3
CREATE SECRET s3_secret (
TYPE s3,
KEY_ID 'POLARIS123ACCESS',
SECRET 'POLARIS456SECRET',
ENDPOINT 'rgw1:7480',
URL_STYLE 'path',
USE_SSL false
);
-- catalog
ATTACH 'test_catalog' AS test_catalog (
TYPE iceberg,
SECRET token,
ENDPOINT 'http://polaris:8181/api/catalog',
ACCESS_DELEGATION_MODE 'none'
);
-- create namespace
CREATE SCHEMA IF NOT EXISTS test_catalog.public;
-- create table
CREATE TABLE IF NOT EXISTS test_catalog.public.raw (
id BIGINT,
name VARCHAR,
age INTEGER
);
-- write
INSERT INTO test_catalog.public.raw VALUES
(1, 'Alice', 30),
(2, 'Bob', 25),
(3, 'Charlie', 35);
-- read
SELECT * FROM test_catalog.public.raw
WHERE age >= 30;
-- delete (table first, then namespace)
DROP TABLE IF EXISTS test_catalog.public.raw;
DROP SCHEMA IF EXISTS test_catalog.public;
s3に設定において、http://(場合によってはhttps://)を省略して書くのがお作法のようです。
最後に
無駄な設定があったら申し訳ないです。Icebergはクライアントだけでも様々なバリデーションがあって面白いですよね。
最近はpg_lakeというSnowflake社製のPostgreSQLカタログが気になってます。Docker hubに公開されたら触ってみたいですね。