注: 以下の翻訳の正確性は検証されていません。AIPを利用して英語版の原文から機械的に翻訳されたものです。
Code Workspaces の Jupyter® ノートブックから、PyIceberg、SQL、または Spark を使用して Iceberg テーブルを操作できます。
仮想 Iceberg テーブルの場合、SQL ではソースをワークスペースに追加する必要がないため、ソースでコードインポートを有効にする必要もありません。ソースをコードで使用すべきでない場合は、この方法が適している可能性があります。
Jupyter® ノートブックで PyIceberg を使用して Iceberg テーブルの読み取りと書き込みを行うには、以下の手順に沿って進めます。


Jupyter® ノートブックで SQL を使用して Iceberg テーブルをクエリするには、上記と同じ手順に沿って進めます。データの読み取りを選択して Iceberg テーブルをワークスペースに追加した後、モードを transforms-table から containers-sql に切り替えます。

Code Workspaces で SQL を使用して表形式データセットをクエリする方法について、詳しくはこちら。
Jupyter® ノートブックで Spark を使用して Iceberg テーブルの読み取りと書き込みを行うには、以下のセクションの手順に沿って進めます。
Spark 3.5 with Scala 2.12 と aws-bundle の JAR をダウンロードします。/libs という名前の新しいフォルダーを作成し、このフォルダーに JAR をアップロードします。まず、Spark セッションを作成します。このコードを実行すると、ユーザートークンの入力を求めるプロンプトが表示されます。トークンはアカウント設定で生成できます。トークンを作成する詳しい手順については、ユーザー生成トークンを参照してください。
Copied!1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17from pyspark.sql import SparkSession from getpass import getpass spark = ( SparkSession.builder .master("local[*]") .appName("foundry") .config("spark.jars", "file:///home/user/repo/libs/iceberg-spark-runtime-3.5_2.12-1.9.1.jar,file:///home/user/repo/libs/iceberg-aws-bundle-1.9.1.jar") .config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") .config("spark.sql.catalog.foundry", "org.apache.iceberg.spark.SparkCatalog") .config("spark.sql.catalog.foundry.type", "rest") .config("spark.sql.catalog.foundry.uri", "https://<your_foundry_url>/iceberg") .config("spark.sql.catalog.foundry.default-namespace", "foundry") .config("spark.sql.catalog.foundry.token", getpass("Foundry token:")) .config("spark.sql.defaultCatalog", "foundry") .getOrCreate() )
Iceberg のドキュメント ↗では、Iceberg カタログへの接続を確立するために使用する上記のパラメーターについて、より詳しく説明されています。ステップ2でアップロードした JAR の名前を使用して、spark.jars のファイルパスを必ず更新してください。
これらのうち、Foundry 固有の Iceberg カタログパラメーター ↗は次のとおりです。
| パラメーター | 値 | 説明 |
|---|---|---|
spark.sql.catalog.foundry | org.apache.iceberg.spark.SparkCatalog | カタログの実装クラス ↗。 |
spark.sql.catalog.foundry.type | rest | 基盤となるカタログの実装の種類(REST) |
spark.sql.catalog.foundry.uri | https://<your_foundry_url>/iceberg | REST カタログの URL |
spark.sql.catalog.foundry.default-namespace | foundry | カタログのデフォルトの名前空間 |
spark.sql.catalog.foundry.token | getpass("Foundry token:") | トークンアクセス資格情報の入力を求めるプロンプトを表示します |
パスに空白が含まれる場合は、空白が正しくエスケープされていることを確認する必要があります。Spark では、バッククォート(`)を使用して空白をエスケープできます。例:`/.../My folder/Iceberg table`。PyIceberg では、URL エンコーディングを使用できます。
これで、Spark セッションを使用して Iceberg テーブルの読み取りと書き込みを行えます。たとえば、Iceberg ドキュメントのクイックスタートガイド ↗に沿って、テーブルを作成して行を挿入できます。
Copied!1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21from pyspark.sql.types import DoubleType, FloatType, LongType, StructType, StructField, StringType schema = StructType([ StructField("vendor_id", LongType(), True), StructField("trip_id", LongType(), True), StructField("trip_distance", FloatType(), True), StructField("fare_amount", DoubleType(), True), StructField("store_and_fwd_flag", StringType(), True) ]) df = spark.createDataFrame([], schema) df.writeTo("`/.../taxis`").create() schema = spark.table("`/.../taxis`").schema data = [ (1, 1000371, 1.8, 15.32, "N"), (2, 1000372, 2.5, 22.15, "N"), (2, 1000373, 0.9, 9.01, "N"), (1, 1000374, 8.4, 42.13, "Y") ] df = spark.createDataFrame(data, schema) df.writeTo("`/.../taxis`").append()
Copied!1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16spark.sql(""" CREATE TABLE `/.../taxis` ( vendor_id bigint, trip_id bigint, trip_distance float, fare_amount double, store_and_fwd_flag string ) PARTITIONED BY (vendor_id); """) spark.sql(""" INSERT INTO `/.../taxis` VALUES (1, 1000371, 1.8, 15.32, 'N'), (2, 1000372, 2.5, 22.15, 'N'), (2, 1000373, 0.9, 9.01, 'N'), (1, 1000374, 8.4, 42.13, 'Y'); """)
Jupyter®、JupyterLab®、および Jupyter® のロゴは、NumFOCUS の商標または登録商標です。
ここで言及されている第三者の商標(ロゴおよびアイコンを含む)はすべて、それぞれの所有者に帰属します。提携または推奨を示唆するものではありません。