跳转到内容
搜索文档

Spark (PySpark)

最后更新 查看 MarkdownAgent 设置

以下示例演示如何使用 PySpark 连接到 R2 Data Catalog。

前提条件

示例用法

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()

这篇文档对您有帮助吗?