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
相关产品推荐
相关产品推荐

