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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 14:20:27