You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何创建带Geometry列的Apache Iceberg表?遇自定义类型不支持错误

解决方案:PySpark + Sedona 创建含原生Geometry列的Apache Iceberg表

错误原因

Sedona的Geometry类型属于Spark用户自定义类型(UDT),Iceberg默认不直接支持将UDT写入表中,因此抛出UnsupportedOperationException: User-defined types are not supported错误。要实现需求,需要将Sedona的Geometry转换为Iceberg能识别的原生地理格式(WKB,Well-Known Binary),并配置Iceberg支持地理类型解析。

修正步骤

  • 补充Iceberg扩展依赖,确保地理类型支持
  • 将Sedona Geometry对象转换为WKB二进制格式,适配Iceberg的原生地理列要求
  • 写入时明确指定表的Schema,标记geometry列类型

完整修正代码

from pyspark.sql import SparkSession
from sedona.spark import SedonaContext
from pyspark.sql.functions import expr

spark = (
    SparkSession.builder
      .appName("SedonaIcebergApp")
      .master("local[*]")
      .config(
          "spark.jars.packages",
          ",".join([
              "org.apache.sedona:sedona-spark-shaded-3.5_2.12:1.7.1",
              "org.datasyslab:geotools-wrapper:1.7.1-28.5",
              "org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.9.0",
              "org.apache.iceberg:iceberg-spark-extensions-3.5_2.12:1.9.0"
          ])
      )
      .config(
          "spark.jars.repositories",
          "https://artifacts.unidata.ucar.edu/repository/unidata-all"
      )
      .config("spark.sql.catalog.local", "org.apache.iceberg.spark.SparkCatalog")
      .config("spark.sql.catalog.local.type", "hadoop")
      .config("spark.sql.catalog.local.warehouse", "spark-warehouse")
      .config("spark.sql.extensions",
              "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions,"
              "org.apache.sedona.sql.SedonaSqlExtensions")
      .getOrCreate()
)

sedona = SedonaContext.create(spark)

# 初始化数据并转换为Sedona Geometry对象
df = sedona.createDataFrame([
    ("a", 'LINESTRING(1.0 3.0,3.0 1.0)'),
    ("b", 'LINESTRING(2.0 5.0,6.0 1.0)'),
], ["id", "geometry_str"])
df = df.withColumn("geometry", expr("ST_GeomFromText(geometry_str)"))

# 将Sedona Geometry转换为WKB二进制,适配Iceberg原生地理列
df = df.withColumn("geometry", expr("ST_AsBinary(geometry)")).drop("geometry_str")

# 写入Iceberg表,明确指定Parquet格式与地理列类型
df.writeTo("local.platlab.test_geo") \
    .using("iceberg") \
    .tableProperty("format-version", "3") \
    .tableProperty("write.format.default", "parquet") \
    .option("schema", "id string, geometry geometry") \
    .createOrReplace()

验证方法

写入完成后,可通过以下代码确认表结构与数据正确性:

# 读取Iceberg表
read_df = spark.read.format("iceberg").load("local.platlab.test_geo")
read_df.printSchema()
# 预期输出:root
#  |-- id: string (nullable = true)
#  |-- geometry: geometry (nullable = true)

# 转换回Sedona Geometry并查看数据
read_df.withColumn("geometry", expr("ST_GeomFromWKB(geometry)")).show()

内容的提问来源于stack exchange,提问作者kawa li

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.13 01:28:11