PySpark读取JSON转RDD/DataFrame报错:TypeError: 'RDD' object is not iterable
错误原因分析
- RDD直接迭代错误:
parse_dataframe函数接收的是RDD对象,但convert_single_object_per_line里直接用for line in json_list遍历RDD——RDD是分布式数据集,不能在Driver端直接迭代,这直接触发了TypeError: 'RDD' object is not iterable。 - 原始数据解析无效:你对变量
a做了json.dumps(a)再json.loads(a)的操作,实际上只是把原始字符串转了一圈,并没有把它解析成字典结构,后续无法正确提取OUT字段的值。 - 数据格式不匹配:你的原始数据不是标准JSON格式,强行用
jsonRDD解析无法得到期望的DataFrame结构。
修正后的代码实现
from pyspark.sql import SparkSession from ast import literal_eval # 初始化SparkSession(Spark 2.0+推荐用这个替代旧的SparkContext+sqlContext) spark = SparkSession.builder.appName("Simple App").getOrCreate() # 原始数据 a = "{'OUT': '000011@@@;1-2=239=0;1-3=271=0;1-4=266=0;1-5=103=0;1-6=682=0;1-7=81=0;1-8=389=0;1-9=155=0;1-10=28=0;1-11=64=0;1-12=268=0;1-13=335=0;1-14=999=0;2-3=257=0;2-4=294=0;2-5=96=0;2-6=999=0;2-7=81=0;2-8=372=0;2-9=163=0'}" # 1. 解析原始字符串为字典(用literal_eval处理单引号的字典字符串) data_dict = literal_eval(a) out_value = data_dict['OUT'] # 2. 拆分OUT的值,提取a字段和各个条目 parts = out_value.split('@@@') a_value = parts[0] items = parts[1].split(';') # 3. 构造DataFrame的行数据 rows = [] for item in items: if not item: # 跳过空条目 continue b, c, e = item.split('=') rows.append({ 'a': a_value, 'b': b, 'c': int(c), 'e': int(e) }) # 4. 生成DataFrame df = spark.createDataFrame(rows) # 显示前2行结果,匹配你的期望格式 df.show(n=2, truncate=False)
代码说明
- 替换旧API:用SparkSession整合SparkContext和SQLContext的功能,是Spark 2.0+的标准入口,无需单独初始化多个上下文。
- 正确解析原始数据:用
ast.literal_eval处理单引号包裹的字典字符串,避免JSON解析报错,精准提取OUT字段内容。 - 手动拆分数据:根据原始数据的分隔规则(
@@@分割a字段和条目,;分割条目,=分割字段),构造符合期望格式的行数据,再通过createDataFrame直接生成DataFrame,彻底避免错误的RDD迭代操作。 - 类型转换:将
c和e字段转为整数类型,保证数据类型符合业务逻辑。
运行结果
+------+----+---+---+ |a |b |c |e | +------+----+---+---+ |000011|1-2 |239|0 | |000011|1-3 |271|0 | +------+----+---+---+ only showing top 2 rows
内容的提问来源于stack exchange,提问作者Mary
相关产品推荐
相关产品推荐

