0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

Iceberg(REST+S3)を様々なクライアントで操作(pyiceberg, spark, duckdb)

0
Last updated at Posted at 2026-03-01

はじめに

概要

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のコンテナを加えます。

docker-compose.yaml
# ~~~~~~~~~~~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とします。

dockerfiles/python/Dockerfile
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
dockerfiles/python/requirements.txt
pyiceberg[s3fs]==0.10.0
pandas
pyarrow

SparkイメージではSpark関係のコマンドをパスに通しておきます。Pysparkも使うためPYTHONPATH変数も定義しておきます。

dockerfiles/spark/Dockerfile
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"

これでカタログができました。

クライアントで操作してみる

ここからは作成したカタログを使って、以下の操作をしてみます。

  1. ネームスペースの作成
  2. テーブルの作成
  3. データの挿入
  4. データの閲覧
  5. テーブルとネームスペースの削除

実際に各種スクリプトやコマンドを試してみる場合は、以下のコマンドで各種コンテナに入ってから実行してください。

$ 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

https://py.iceberg.apache.org/configuration/ より

今回は1つ目と3つ目の方法を試してみます。

本記事の作成時点(2026/02/28)のpyicebergの最新バージョンは0.11.0ですが、このバージョンはRESTカタログのHttp操作に関するバグを含んでいるみたいです。(PR-3010)
なので1つ古い0.10.0を使用しています。

.pyiceberg.yamlに設定値を記載する方法

まず、Pythonスクリプトがあるディレクトリに.pyiceberg.yamlを作成します。

.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を読み込んで、カタログに接続してくれます。

mounts/python/yaml/main.py
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形で入れた後に展開しています。

mounts/python/by_passing/main.py
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と同じ操作をしてみます。

use_pyspark.py
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とはちょっと勝手が違います。

use_duckdb.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に公開されたら触ってみたいですね。

0
0
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
0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?