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

PySpark创建DataFrame报错:TypeError: 'JavaPackage' object is not callable

问题描述

在学习Kafka + PySpark时,尝试将测试消息转为Spark DataFrame,使用createDataFrame方法时始终报错:TypeError: 'JavaPackage' object is not callable。其他Spark操作(如读取Kafka流、CSV数据)可正常运行,但从元组或Pandas DataFrame创建DataFrame均触发该错误。

环境版本

  • PySpark版本:3.3.1
  • Scala版本:2.12.15
  • OpenJDK版本:17.0.6
  • Python版本:3.11.3

测试代码

# PYSPARK CREATE DATAFRAME FROM dictionary
import pandas as pd
from pyspark.sql.types import StructType, StructField, IntegerType, StringType
from pyspark.sql import SparkSession
from pyspark.sql.functions import from_json, udf
import findspark

findspark.init()

# Config
spark = SparkSession \
    .builder \
    .master("local[*]") \
    .appName("PySparkTest") \
    .getOrCreate()

#Creating the test dictionary I want to append to a spark dataframe
test_message = {'user_id': 19, 'recipient_id': 57, 'message': 'YbfyRHyWgjuGlzOiudEcVMLJNzqUPDvV'}


#put it into a pandas dataframe
df_pandas = pd.DataFrame([test_message])
df_pandas


##  Create a spark schema/column headers
schema = StructType([
    StructField("user_id", IntegerType(), True),
    StructField("recipient_id", IntegerType(), True),
    StructField("message", StringType(), True)
])

df_spark = spark.createDataFrame(df_pandas,schema)
df_spark.show()


# ## ALTERNATIVE WAY: Create DataFrame from a single row
# ### PARSING THE JSON COMING OUT 
# user_id = deserialized_cons.get('user_id')
# recipient_id = deserialized_cons.get('recipient_id')
# message = deserialized_cons.get('message')
# ## turn into row format and create and upload dataframe
# data = [(user_id, recipient_id, message)]
# df_spark = spark.createDataFrame(df_pandas,schema)
# df_spark.show()

报错信息

---------------------------------------------------------------------------
TypeError                                 Traceback (most recent call last)
Cell In[25], line 44
     42 ## turn into row format and create and upload dataframe
     43 data = [(user_id, recipient_id, message)]
---> 44 df_spark = spark.createDataFrame(df_pandas,schema)
     45 df_spark.show()

File ~/opt/anaconda3/envs/spark/lib/python3.11/site-packages/pyspark/sql/session.py:1273, in SparkSession.createDataFrame(self, data, schema, samplingRatio, verifySchema)
   1269     data = pd.DataFrame(data, columns=column_names)
   1271 if has_pandas and isinstance(data, pd.DataFrame):
   1272     # Create a DataFrame from pandas DataFrame.
-> 1273     return super(SparkSession, self).createDataFrame(  # type: ignore[call-overload]
   1274         data, schema, samplingRatio, verifySchema
   1275     )
   1276 return self._create_dataframe(
   1277     data, schema, samplingRatio, verifySchema  # type: ignore[arg-type]
   1278 )

File ~/opt/anaconda3/envs/spark/lib/python3.11/site-packages/pyspark/sql/pandas/conversion.py:440, in SparkConversionMixin.createDataFrame(self, data, schema, samplingRatio, verifySchema)
    438             raise
    439 converted_data = self._convert_from_pandas(data, schema, timezone)
-> 440 return self._create_dataframe(converted_data, schema, samplingRatio, verifySchema)

File ~/opt/anaconda3/envs/spark/lib/python3.11/site-packages/pyspark/sql/session.py:1320, in SparkSession._create_dataframe(self, data, schema, samplingRatio, verifySchema)
   1318     rdd, struct = self._createFromLocal(map(prepare, data), schema)
   1319 assert self._jvm is not None
-> 1320 jrdd = self._jvm.SerDeUtil.toJavaArray(rdd._to_java_object_rdd())
   1321 jdf = self._jsparkSession.applySchemaToPythonRDD(jrdd.rdd(), struct.json())
   1322 df = DataFrame(jdf, self)

File ~/opt/anaconda3/envs/spark/lib/python3.11/site-packages/pyspark/rdd.py:4897, in RDD._to_java_object_rdd(self)
   4894 rdd = self._pickled()
   4895 assert self.ctx._jvm is not None
-> 4897 return self.ctx._jvm.SerDeUtil.pythonToJava(rdd._jrdd, True)

File ~/opt/anaconda3/envs/spark/lib/python3.11/site-packages/pyspark/rdd.py:5441, in PipelinedRDD._jrdd(self)
   5438 else:
   5439     profiler = None
-> 5441 wrapped_func = _wrap_function(
   5442     self.ctx, self.func, self._prev_jrdd_deserializer, self._jrdd_deserializer, profiler
   5443 )
   5445 assert self.ctx._jvm is not None
   5446 python_rdd = self.ctx._jvm.PythonRDD(
   5447     self._prev_jrdd.rdd(), wrapped_func, self.preservesPartitioning, self.is_barrier
   5448 )

File ~/opt/anaconda3/envs/spark/lib/python3.11/site-packages/pyspark/rdd.py:5243, in _wrap_function(sc, func, deserializer, serializer, profiler)
   5241 pickled_command, broadcast_vars, env, includes = _prepare_for_python_RDD(sc, command)
   5242 assert sc._jvm is not None
-> 5243 return sc._jvm.SimplePythonFunction(
   5244     bytearray(pickled_command),
   5245     env,
   5246     includes,
   5247     sc.pythonExec,
   5248     sc.pythonVer,
   5249     broadcast_vars,
   5250     sc._javaAccumulator,
   5251 )

TypeError: 'JavaPackage' object is not callable

额外异常

执行spark-submit --version时出现:

23/05/27 12:19:59 WARN Utils: Your hostname, Nics-MacBook-Pro.local resolves to a loopback address: 127.0.0.1; using 192.168.4.24 instead (on interface en0)
23/05/27 12:19:59 WARN Utils: Set SPARK_LOCAL_IP if you need to bind to another address
23/05/27 12:19:59 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
Exception in thread "main" org.apache.spark.SparkException: Failed to get main class in JAR with error 'File file:/Users/nicburkett/— does not exist'.  Please specify one with --class.

解决方案

1. 排查环境变量配置

JavaPackage错误根源是PySpark无法定位对应Java类,结合spark-submit异常,需确认以下环境变量:

  • SPARK_HOME:必须指向完整的Spark安装目录
  • PYSPARK_PYTHON/PYSPARK_DRIVER_PYTHON:指定当前conda环境的Python解释器路径,避免版本冲突
  • JAVA_HOME:指向OpenJDK 17的安装路径,Spark 3.3.1兼容Java 17

2. 修正findspark初始化

显式指定Spark路径,避免自动查找出错:

# 替换为你的Spark实际安装路径
findspark.init("/usr/local/spark")

3. 绕过Pandas直接创建DataFrame

跳过Pandas转换,直接用Python原生数据结构创建:

# 方式1:字典列表
test_message = {'user_id': 19, 'recipient_id': 57, 'message': 'YbfyRHyWgjuGlzOiudEcVMLJNzqUPDvV'}
df_spark = spark.createDataFrame([test_message], schema=schema)
df_spark.show()

# 方式2:元组列表
data = [(19, 57, 'YbfyRHyWgjuGlzOiudEcVMLJNzqUPDvV')]
df_spark = spark.createDataFrame(data, schema=schema)
df_spark.show()

4. 检查Spark安装完整性

SimplePythonFunction是Spark核心类,若找不到说明安装包损坏:

  • 重新下载对应版本的Spark(Spark 3.3.1 + Scala 2.12版本)
  • 确保解压后目录权限正常,无文件缺失

5. 修复spark-submit命令错误

spark-submit --version的异常是因为命令后存在无效字符,执行正确命令:

spark-submit --version

若仍报错,使用完整路径执行:

/path/to/spark/bin/spark-submit --version

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 08:07:01