以下示例演示如何使用 PySpark ↗ 连接到 R2 Data Catalog。
- 注册 Cloudflare 账户 ↗。
- 创建 R2 存储桶并启用数据目录。
- 创建 R2 API 令牌,并授予 R2 和数据目录权限。
- 安装 PySpark ↗ 库。
from pyspark.sql import SparkSession
# 定义目录连接详情(替换变量)
WAREHOUSE = "<WAREHOUSE>"
TOKEN = "<TOKEN>"
CATALOG_URI = "<CATALOG_URI>"
# 构建具有 Iceberg 配置的 Spark 会话
spark = SparkSession.builder \
.appName("R2DataCatalogExample") \
.config('spark.jars.packages', 'org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.6.1,org.apache.iceberg:iceberg-aws-bundle:1.6.1') \
.config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
.config("spark.sql.catalog.my_catalog", "org.apache.iceberg.spark.SparkCatalog") \
.config("spark.sql.catalog.my_catalog.type", "rest") \
.config("spark.sql.catalog.my_catalog.uri", CATALOG_URI) \
.config("spark.sql.catalog.my_catalog.warehouse", WAREHOUSE) \
.config("spark.sql.catalog.my_catalog.token", TOKEN) \
.config("spark.sql.catalog.my_catalog.header.X-Iceberg-Access-Delegation", "vended-credentials") \
.config("spark.sql.catalog.my_catalog.s3.remote-signing-enabled", "false") \
.config("spark.sql.defaultCatalog", "my_catalog") \
.getOrCreate()
spark.sql("USE my_catalog")
# 如果命名空间不存在则创建它
spark.sql("CREATE NAMESPACE IF NOT EXISTS default")
# 使用 Iceberg 在命名空间中创建表
spark.sql("""
CREATE TABLE IF NOT EXISTS default.my_table (
id BIGINT,
name STRING
)
USING iceberg
""")
# 创建简单的 DataFrame
df = spark.createDataFrame(
[(1, "Alice"), (2, "Bob"), (3, "Charlie")],
["id", "name"]
)
# 将 DataFrame 写入 Iceberg 表
df.write \
.format("iceberg") \
.mode("append") \
.save("default.my_table")
# 从 Iceberg 表中读回数据
result_df = spark.read \
.format("iceberg") \
.load("default.my_table")
result_df.show()