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

解决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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 06:07:03