解决PySpark中将字典转换为DataFrame时的报错问题
问题:Spark 3.4 + Python 3.10 转换字典数据为PySpark DataFrame失败
我正在使用Spark 3.4与Python 3.10,尝试将值为列表的Python字典转换为PySpark DataFrame,但多种方法均失败。
尝试的方法及报错
方法1:通过JSON读取器创建DataFrame
步骤1:将值为列表的字典转换为字典列表
raw_data_size = len(raw_data['origin']) raw_data_keys = raw_data.keys() mapping_list=[] mapping = {} for i in range(raw_data_size): for key in raw_data_keys: mapping[key]=raw_data[key][i] mapping_list.append(mapping) mapping={}
mapping_list[:1]输出:
[{'regiment': 'Nighthawks', 'company': '1st', 'deaths': 523, 'battles': 5, 'size': 1045, 'veterans': 1, 'readiness': 1, 'armored': 1, 'deserters': 4, 'origin': 'Arizona'}]
步骤2:定义Schema
schema = StructType([ StructField('regiment', StringType(), True), StructField('company', StringType(), True), StructField('deaths', IntegerType(), True), StructField('battles', IntegerType(), True), StructField('size', IntegerType(), True), StructField('veterans', IntegerType(), True), StructField('readiness', IntegerType(), True), StructField('armored', IntegerType(), True), StructField('deserters', IntegerType(), True), StructField('origin', StringType(), True) ])
步骤3:创建DataFrame时报错
执行代码:
army=spark.read.schema(schema).json(mapping_list)
报错堆栈:
--------------------------------------------------------------------------- Py4JJavaError Traceback (most recent call last) Cell In[23], line 1 ----> 1 army=spark.read.schema(schema).json(mapping_list) File C:\software\programming\spark-3.4.0-bin-hadoop3\python\pyspark\sql\readwriter.py:418, in DataFrameReader.json(self, path, schema, primitivesAsString, prefersDecimal, allowComments, allowUnquotedFieldNames, allowSingleQuotes, allowNumericLeadingZero, allowBackslashEscapingAnyCharacter, mode, columnNameOfCorruptRecord, dateFormat, timestampFormat, multiLine, allowUnquotedControlChars, lineSep, samplingRatio, dropFieldIfAllNull, encoding, locale, pathGlobFilter, recursiveFileLookup, modifiedBefore, modifiedAfter, allowNonNumericNumbers) 416 if type(path) == list: 417 assert self._spark._sc._jvm is not None ---> 418 return self._df(self._jreader.json(self._spark._sc._jvm.PythonUtils.toSeq(path))) 419 elif isinstance(path, RDD): 421 def func(iterator: Iterable) -> Iterable: File C:\ProgramData\anaconda3\lib\site-packages\py4j\java_gateway.py:1321, in JavaMember.__call__(self, *args) 1315 command = proto.CALL_COMMAND_NAME +\ 1316 self.command_header +\ 1317 args_command +\ 1318 proto.END_COMMAND_PART 1320 answer = self.gateway_client.send_command(command) -> 1321 return_value = get_return_value( 1322 answer, self.gateway_client, self.target_id, self.name) 1324 for temp_arg in temp_args: 1325 temp_arg._detach() File C:\software\programming\spark-3.4.0-bin-hadoop3\python\pyspark\errors\exceptions\captured.py:169, in capture_sql_exception.<locals>.deco(*a, **kw) 167 def deco(*a: Any, **kw: Any) -> Any: 168 try: ---> 169 return f(*a, **kw) 170 except Py4JJavaError as e: 171 converted = convert_exception(e.java_exception) File C:\ProgramData\anaconda3\lib\site-packages\py4j\protocol.py:326, in get_return_value(answer, gateway_client, target_id, name) 324 value = OUTPUT_CONVERTER[type](answer[2:], gateway_client) 325 if answer[1] == REFERENCE_TYPE: -> 326 raise Py4JJavaError( 327 "An error occurred while calling {0}{1}{2}.\n". 328 format(target_id, ".", name), value) 329 else: 330 raise Py4JError( 331 "An error occurred while calling {0}{1}{2}. Trace:\n{3}\n". 332 format(target_id, ".", name, value)) Py4JJavaError: An error occurred while calling o99.json. : java.lang.ClassCastException: class java.util.HashMap cannot be cast to class java.lang.String (java.util.HashMap and java.lang.String are in module java.base of loader 'bootstrap') at scala.collection.immutable.List.map(List.scala:293) at org.apache.spark.sql.execution.datasources.DataSource$.checkAndGlobPathIfNecessary(DataSource.scala:722) at org.apache.spark.sql.execution.datasources.DataSource.checkAndGlobPathIfNecessary(DataSource.scala:551) at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:404) at org.apache.spark.sql.DataFrameReader.loadV1Source(DataFrameReader.scala:229) at org.apache.spark.sql.DataFrameReader.$anonfun$load$2(DataFrameReader.scala:211) at scala.Option.getOrElse(Option.scala:189) at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:211) at org.apache.spark.sql.DataFrameReader.json(DataFrameReader.scala:362) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:568) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:374) at py4j.Gateway.invoke(Gateway.java:282) at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) at py4j.commands.CallCommand.execute(CallCommand.java:79) at py4j.ClientServerConnection.waitForCommands(ClientServerConnection.java:182) at py4j.ClientServerConnection.run(ClientServerConnection.java:106) at java.base/java.lang.Thread.run(Thread.java:833)
其他尝试的方法及报错
- 方法2:直接使用
spark.createDataFrame(mapping_list).show()
报错:Py4JJavaError: An error occurred while calling o174.showString. : org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 5.0 failed 1 times, most recent failure: Lost task 0.0 in stage 5.0 (TID 27) (host.docker.internal executor driver): java.io.IOException: Cannot run program "C:\ProgramData\anaconda3": CreateProcess error=5, Access is denied - 方法3:使用
Row构造数据:spark.createDataFrame(Row(**x) for x in mapping_list).show()
报错:Py4JJavaError: An error occurred while calling o153.showString. : org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 4.0 failed 1 times, most recent failure: Lost task 0.0 in stage 4.0 (TID 26) (host.docker.internal executor driver): java.io.IOException: Cannot run program "C:\ProgramData\anaconda3": CreateProcess error=5, Access is denied - 方法4:转换为Pandas DataFrame后再加载,同样失败。
问题分析
1. JSON读取器报错原因
spark.read.json()方法接收的参数是文件路径(字符串或字符串列表),而非Python字典列表。传入字典列表会导致Spark将其当作文件路径处理,从而抛出ClassCastException(HashMap转String失败)。
2. Access Denied报错原因
这是Spark运行时的权限问题:Spark尝试调用Python解释器时,路径指向了C:\ProgramData\anaconda3目录而非具体的python.exe可执行文件,导致权限不足。通常是PYSPARK_PYTHON环境变量配置错误所致。
解决方法
方法1:修复PYSPARK_PYTHON环境变量
在启动Spark前,设置正确的Python可执行文件路径:
import os os.environ['PYSPARK_PYTHON'] = 'C:/ProgramData/anaconda3/python.exe' # 替换为你的Python实际路径
也可以直接在系统环境变量中配置PYSPARK_PYTHON为Python可执行文件的完整路径。
方法2:正确使用createDataFrame创建DataFrame
修复环境变量后,使用定义好的Schema创建DataFrame:
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, IntegerType # 初始化SparkSession(如果未初始化) spark = SparkSession.builder.appName("ArmyData").getOrCreate() # 你的mapping_list数据 mapping_list = [{'regiment': 'Nighthawks', 'company': '1st', 'deaths': 523, 'battles': 5, 'size': 1045, 'veterans': 1, 'readiness': 1, 'armored': 1, 'deserters': 4, 'origin': 'Arizona'}] # 定义Schema schema = StructType([ StructField('regiment', StringType(), True), StructField('company', StringType(), True), StructField('deaths', IntegerType(), True), StructField('battles', IntegerType(), True), StructField('size', IntegerType(), True), StructField('veterans', IntegerType(), True), StructField('readiness', IntegerType(), True), StructField('armored', IntegerType(), True), StructField('deserters', IntegerType(), True), StructField('origin', StringType(), True) ]) # 创建DataFrame army = spark.createDataFrame(mapping_list, schema=schema) army.show()
方法3:直接转换值为列表的字典(无需手动转字典列表)
可以利用Spark的createDataFrame直接处理值为列表的字典,更高效:
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, IntegerType spark = SparkSession.builder.appName("ArmyData").getOrCreate() # 原始的raw_data(值为列表的字典) raw_data = { 'regiment': ['Nighthawks'], 'company': ['1st'], 'deaths': [523], 'battles': [5], 'size': [1045], 'veterans': [1], 'readiness': [1], 'armored': [1], 'deserters': [4], 'origin': ['Arizona'] } # 定义Schema schema = StructType([ StructField('regiment', StringType(), True), StructField('company', StringType(), True), StructField('deaths', IntegerType(), True), StructField('battles', IntegerType(), True), StructField('size', IntegerType(), True), StructField('veterans', IntegerType(), True), StructField('readiness', IntegerType(), True), StructField('armored', IntegerType(), True), StructField('deserters', IntegerType(), True), StructField('origin', StringType(), True) ]) # 直接转换 army = spark.createDataFrame(zip(*raw_data.values()), schema=schema) army.show()
内容的提问来源于stack exchange,提问作者AJ22
相关产品推荐
相关产品推荐

