Apache Sedona中PySpark DataFrame调用toPandas()报错问题
Spark DataFrame调用
.toPandas()报错 报错场景
执行以下代码时触发报错:
pm_table_bej_test.toPandas()
完整报错信息
Py4JJavaError Traceback (most recent call last) ~\AppData\Local\Temp\ipykernel_25776\2227465720.py in ----> 1 pm_table_bej_test.toPandas() c:\Users\user\Anaconda3\envs\python_ds_gis\lib\site-packages\pyspark\sql\pandas\conversion.py in toPandas(self) 203 204 # Below is toPandas without Arrow optimization. --> 205 pdf = pd.DataFrame.from_records(self.collect(), columns=self.columns) 206 column_counter = Counter(self.columns) 207 c:\Users\user\Anaconda3\envs\python_ds_gis\lib\site-packages\pyspark\sql\dataframe.py in collect(self) 815 """ 816 with SCCallSiteSync(self._sc): --> 817 sock_info = self._jdf.collectToPython() 818 return list(_load_from_socket(sock_info, BatchedSerializer(CPickleSerializer()))) 819 c:\Users\user\Anaconda3\envs\python_ds_gis\lib\site-packages\py4j\java_gateway.py in __call__(self, *args) 1319 1320 answer = self.gateway_client.send_command(command) --> 1321 return_value = get_return_value( 1322 answer, self.gateway_client, self.target_id, self.name) 1323 ... at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:657) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1144) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:642) ... 1 more
当前运行代码
# 导入依赖 import pandas as pd import geopandas as gpd from pyspark.sql import SparkSession from pyspark.sql.functions import col, expr, when from sedona.register import SedonaRegistrator from sedona.utils import SedonaKryoRegistrator, KryoSerializer from sedona.core.formatMapper.shapefileParser import ShapefileReader from sedona.utils.adapter import Adapter # 初始化SparkContext spark = ( SparkSession .builder .master("local[*]") .appName("Sedona App") .config("spark.serializer", KryoSerializer.getName) .config("spark.kryo.registrator", SedonaKryoRegistrator.getName) .getOrCreate() ) # 从CSV导入表 pm_table_bej = ( spark .read .option("delimiter", ",") .option("header", "true") .csv(CSV_PATH) ) # 创建临时视图(原代码缺失此关键步骤,补充后可正常运行) pm_table_bej.createOrReplaceTempView("pm_table_bej_temp_view") # 查询语句 pm_table_bej_test = spark.sql( """ SELECT OID_, ST_GeomFromText( CONCAT( 'POINT(', Lon_of_Observation_Point, ' ', Lat_of_Observation_Point, ')' ) ) AS geometry, Local_Time_of_Day_of_Observation_Point, Unix_Timestamp_of_Observation_Point, from_unixtime(Unix_Timestamp_of_Observation_Point) AS timestamp FROM pm_table_bej_temp_view WHERE TRUE AND Local_Day_of_Week_of_Observation_Point = "Fri" """ ) pm_table_bej_test.printSchema() pm_table_bej_test.show(5)
代码执行输出结果
root |-- OID_: string (nullable = true) |-- geometry: geometry (nullable = true) |-- Local_Time_of_Day_of_Observation_Point: string (nullable = true) |-- Unix_Timestamp_of_Observation_Point: string (nullable = true) |-- timestamp: string (nullable = true) +----+--------------------+--------------------------------------+-----------------------------------+-------------------+ |OID_| geometry|Local_Time_of_Day_of_Observation_Point|Unix_Timestamp_of_Observation_Point| timestamp| +----+--------------------+--------------------------------------+-----------------------------------+-------------------+ | 14|POINT (106.783285...| 00:39:16| 1675359556|2023-02-03 00:39:16| | 16|POINT (106.78329 ...| 00:39:16| 1675359556|2023-02-03 00:39:16| | 21|POINT (106.78349 ...| 06:01:30| 1675378890|2023-02-03 06:01:30| | 22|POINT (106.78349 ...| 06:01:45| 1675378905|2023-02-03 06:01:45| | 23|POINT (106.783517...| 06:01:08| 1675378868|2023-02-03 06:01:08| +----+--------------------+--------------------------------------+-----------------------------------+-------------------+ only showing top 5 rows
问题分析与解决方案
核心原因
报错根源在于:Spark DataFrame中的geometry列是Sedona专属的空间几何类型,默认无法直接被Python的序列化机制解析,调用toPandas()时会触发跨语言序列化失败,进而抛出Py4JJavaError。
解决方案
方案1:将几何类型转为WKT字符串后转Pandas
在SQL查询中把几何对象转为WKT文本格式,让Pandas可以正常解析:
pm_table_bej_test = spark.sql( """ SELECT OID_, ST_AsText(ST_GeomFromText( CONCAT( 'POINT(', Lon_of_Observation_Point, ' ', Lat_of_Observation_Point, ')' ) )) AS geometry, Local_Time_of_Day_of_Observation_Point, Unix_Timestamp_of_Observation_Point, from_unixtime(Unix_Timestamp_of_Observation_Point) AS timestamp FROM pm_table_bej_temp_view WHERE TRUE AND Local_Day_of_Week_of_Observation_Point = "Fri" """ ) # 正常转换为Pandas DataFrame pdf = pm_table_bej_test.toPandas() # 如需转为GeoDataFrame,再手动解析WKT列 gdf = gpd.GeoDataFrame(pdf, geometry=gpd.GeoSeries.from_wkt(pdf['geometry']))
方案2:使用Sedona Adapter直接转GeoDataFrame
Sedona提供了专门的工具类,可直接将Spark空间DataFrame转为GeoPandas的GeoDataFrame:
from sedona.utils.adapter import Adapter # 一步转换为GeoDataFrame gdf = Adapter.to_geopandas(pm_table_bej_test)
方案3:开启Arrow优化(针对大数据量)
如果数据量较大,可开启Spark的Arrow优化加速转换,同时确保Sedona序列化配置正确:
# 初始化SparkSession时添加Arrow配置 spark = ( SparkSession .builder .master("local[*]") .appName("Sedona App") .config("spark.serializer", KryoSerializer.getName) .config("spark.kryo.registrator", SedonaKryoRegistrator.getName) .config("spark.sql.execution.arrow.pyspark.enabled", "true") .config("spark.sql.execution.arrow.pyspark.fallback.enabled", "true") .getOrCreate() ) # 配合方案1的WKT转换后,再调用toPandas pdf = pm_table_bej_test.toPandas()
内容的提问来源于stack exchange,提问作者Amri Rasyidi
相关产品推荐
相关产品推荐

