如何用Python Spark Core合并JSON与CSV文件中的数据?
解决Spark Core合并JSON与CSV数据的方案
问题修正与合并步骤
首先需要修正CSV数据处理中的错误,再通过Spark的join操作完成数据合并,最终生成目标格式的元组。
1. 修正CSV数据处理逻辑
原CSV处理代码存在两处问题:
- 遗漏了CSV的第二个字段(如示例中的
SDS111),而目标元组要求包含该字段 reduceByKey后的map操作索引越界(kv[1][2]超出元组长度)
修正后的CSV处理代码:
csv_data = lines.map(lambda line: line.split(';')) \ .map(lambda values: (values[0], values[1], float(values[2]), float(values[3]))) \ .map(lambda kv: (kv[0], (kv[1], (kv[2], kv[3])))) \ .reduceByKey(lambda a, b: (a[0], (max(a[1][0] + b[1][0]), max(a[1][1] + b[1][1])))) \ .sortByKey()
说明:
- 保留CSV的第二个字段
values[1],将其与数值元组绑定为值部分 - 修正
reduceByKey后的逻辑,确保合并后仍保留业务字段,同时避免索引错误 - 移除冗余的
map操作,简化处理流程
2. 调整JSON数据的键值结构
为适配join操作,将JSON的RDD转换为以item1为键,(item2, item3)为值的结构:
json_data = rdd.map(lambda item: (item["item1"], (item["item2"], item["item3"])))
3. 合并RDD并生成目标元组
使用Spark的join操作,基于共同的code字段(即item1/CSV首列)合并数据,再通过map转换为目标格式:
# 合并JSON与CSV数据 combined_rdd = json_data.join(csv_data) # 转换为目标元组格式:(item2, item3, code, SDS字段, 数值元组) result_rdd = combined_rdd.map(lambda x: ( x[1][0][0], # JSON的item2 x[1][0][1], # JSON的item3 x[0], # 共同的code字段 x[1][1][0], # CSV的第二个业务字段 x[1][1][1] # 合并后的数值元组 )) # 可选:打印结果或保存到存储 result_rdd.foreach(print)
完整修正后的代码
import json import pyspark sc = pyspark.SparkContext('local[*]') try: # 处理JSON数据 with open("/content/drive/../file.json") as f: data = json.load(f) rdd = sc.parallelize(data) json_data = rdd.map(lambda item: (item["item1"], (item["item2"], item["item3"]))) # 处理CSV数据 lines = sc.textFile('/content/drive/../file*.csv', 5) \ .map(lambda line: line.strip()) \ .filter(lambda line: len(line.split(';')) == 4) csv_data = lines.map(lambda line: line.split(';')) \ .map(lambda values: (values[0], values[1], float(values[2]), float(values[3]))) \ .map(lambda kv: (kv[0], (kv[1], (kv[2], kv[3])))) \ .reduceByKey(lambda a, b: (a[0], (max(a[1][0] + b[1][0]), max(a[1][1] + b[1][1])))) \ .sortByKey() # 合并数据并生成目标元组 combined_rdd = json_data.join(csv_data) result_rdd = combined_rdd.map(lambda x: ( x[1][0][0], x[1][0][1], x[0], x[1][1][0], x[1][1][1] )) # 示例:输出结果 result_rdd.foreach(print) sc.stop() except Exception as e: print(e) sc.stop()
关键补充说明
join操作仅保留两个RDD中都存在的code记录,若需保留JSON中所有记录(即使CSV无对应数据),可替换为leftOuterJoin- 若CSV中同一
code有多条记录,reduceByKey的合并逻辑可根据实际业务需求调整(如直接取最大值、求和等)
内容的提问来源于stack exchange,提问作者KennyBerga
相关产品推荐
相关产品推荐

