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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 09:37:04