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

Spark DataFrame转pandas-on-spark后创建GeoDataFrame报ArrowInvalid错误求助

解决pandas-on-spark DataFrame转GeoDataFrame的ArrowInvalid错误

问题原因

pandas-on-spark DataFrame底层基于Arrow存储数据,而geopandas通过apply生成的shapely Point对象无法被Arrow的类型推断机制识别,从而抛出ArrowInvalid错误。

解决方案

方案一:转普通pandas DataFrame后处理

将pandas-on-spark DataFrame转换为普通pandas DataFrame(原生支持shapely对象),再构建GeoDataFrame:

import geopandas as gpd
from shapely.geometry import Point

# 转换为普通pandas DataFrame
df_pandas = df_new.to_pandas()

# 构建GeoDataFrame
model_gdf = gpd.GeoDataFrame(
    df_pandas,
    crs=sweden_gdf.crs,
    geometry=df_pandas.apply(lambda row: Point(row['x']*1000, row['y']*1000), axis=1)
)

适用场景:数据量较小,内存足以容纳普通pandas DataFrame。

方案二:通过WKT字符串过渡(大数据量友好)

利用Spark UDF生成WKT格式的点字符串,再通过geopandas解析为几何对象,避免全量转普通pandas的内存压力:

from pyspark.sql.functions import udf
from pyspark.sql.types import StringType
import geopandas as gpd

# 定义UDF生成WKT格式的点字符串
def generate_point_wkt(x, y):
    return f"POINT({x*1000} {y*1000})"

point_wkt_udf = udf(generate_point_wkt, StringType())

# 在pandas-on-spark DataFrame中添加WKT列
df_with_wkt = df_new.withColumn("geometry_wkt", point_wkt_udf(df_new['x'], df_new['y']))

# 转普通pandas后解析WKT生成几何列
df_pandas = df_with_wkt.to_pandas()
model_gdf = gpd.GeoDataFrame(
    df_pandas,
    crs=sweden_gdf.crs,
    geometry=gpd.GeoSeries.from_wkt(df_pandas['geometry_wkt'])
)

适用场景:数据量较大,需要利用Spark分布式处理能力,减少内存占用。

注意事项

目前geopandas对pandas-on-spark的原生支持有限,需通过普通pandas DataFrame或WKT等中间格式完成转换。

内容的提问来源于stack exchange,提问作者Nafisah Abidemi Abdulkadir

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 10:07:31