データ接続と統合Iceberg tablesFoundryでのIcebergの使用Code WorkspacesでのJupyter®

注: 以下の翻訳の正確性は検証されていません。AIPを利用して英語版の原文から機械的に翻訳されたものです。

Code Workspaces の Jupyter® ノートブックで Iceberg を使用する

Code Workspaces の Jupyter® ノートブックから、PyIceberg、SQL、または Spark を使用して Iceberg テーブルを操作できます。

仮想 Iceberg テーブルの場合、SQL ではソースをワークスペースに追加する必要がないため、ソースでコードインポートを有効にする必要もありません。ソースをコードで使用すべきでない場合は、この方法が適している可能性があります。

PyIceberg

Jupyter® ノートブックで PyIceberg を使用して Iceberg テーブルの読み取りと書き込みを行うには、以下の手順に沿って進めます。

  1. Code Workspaces で Jupyter® ワークスペースを開きます。
  2. データパネルで追加 > データの読み取りを選択して、Iceberg テーブルをワークスペースに追加します。

Jupyter ワークスペースのデータパネルに、データの読み取りの選択肢がある追加ドロップダウンが表示されています。

  1. データパネルのテーブルのインポートセクションに記載されている手順に沿って、必要なクライアントライブラリをインストールし、テーブルのエイリアスを追加して、データを操作するためのコードスニペットをコピーします。

Jupyter ワークスペースのデータパネルにあるテーブルのインポートセクションに、transforms-tables をインストールする手順が表示されています。

SQL

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

テーブルのインポートの表示に、containers-sql の選択肢があるモードドロップダウンが表示されています。

Code Workspaces で SQL を使用して表形式データセットをクエリする方法について、詳しくはこちら。

Spark(高度な使用方法)

Jupyter® ノートブックで Spark を使用して Iceberg テーブルの読み取りと書き込みを行うには、以下のセクションの手順に沿って進めます。

Iceberg を使用するために Code Workspaces をセットアップする

  1. PySpark のセットアップ:Code Workspaces の FAQ ドキュメントの手順に沿って、PySpark を使用するためにコードワークスペースをセットアップします。
  2. Iceberg JAR をアップロードする:Iceberg のリリースページ ↗から Spark 3.5 with Scala 2.12 と aws-bundle の JAR をダウンロードします。/libs という名前の新しいフォルダーを作成し、このフォルダーに JAR をアップロードします。
  3. ネットワークポリシー: Iceberg ストレージバケットのネットワークポリシーをコードワークスペースにインポートします。

Jupyter® ノートブックのコード例

まず、Spark セッションを作成します。このコードを実行すると、ユーザートークンの入力を求めるプロンプトが表示されます。トークンはアカウント設定で生成できます。トークンを作成する詳しい手順については、ユーザー生成トークンを参照してください。

Copied!
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 from 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.foundryorg.apache.iceberg.spark.SparkCatalogカタログの実装クラス ↗。
spark.sql.catalog.foundry.typerest基盤となるカタログの実装の種類(REST)
spark.sql.catalog.foundry.urihttps://<your_foundry_url>/icebergREST カタログの URL
spark.sql.catalog.foundry.default-namespacefoundryカタログのデフォルトの名前空間
spark.sql.catalog.foundry.tokengetpass("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 21 from 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 16 spark.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 の商標または登録商標です。

ここで言及されている第三者の商標(ロゴおよびアイコンを含む)はすべて、それぞれの所有者に帰属します。提携または推奨を示唆するものではありません。