使用Spark实现JSON扁平化映射:关联客户ID与客户名
Spark扁平化处理JSON对应列表数据的解决方案
问题描述
现有如下结构的JSON数据,data字段下包含两个对应列表cust_name和cust_id,需要用Spark将其处理成客户ID与客户名称一一对应的行数据。
Input JSON:
{ "data":{ "cust_name":["cust_1","cust_2","cust_3","cust_4","cust_5"], "cust_id":[1,2,3,4,5] } }
Expected Output:
------------------------------- | customer_id | customer_name | ------------------------------- | 1 | cust_1 | | 2 | cust_2 | | 3 | cust_3 | | 4 | cust_4 | | 5 | cust_5 | -------------------------------
解决方案
核心思路是用arrays_zip将两个列表按位置配对,再用explode扁平化数组,以下是Scala和Python两种实现方式:
Scala 实现
import org.apache.spark.sql.functions.{arrays_zip, explode, col} // 读取JSON数据(示例用字符串,实际可替换为文件路径) val jsonStr = """{"data":{"cust_name":["cust_1","cust_2","cust_3","cust_4","cust_5"],"cust_id":[1,2,3,4,5]}}""" val df = spark.read.json(spark.sparkContext.parallelize(Seq(jsonStr))) // 处理逻辑 val resultDf = df .select("data.cust_id", "data.cust_name") .withColumn("zipped", arrays_zip(col("cust_id"), col("cust_name"))) .select(explode(col("zipped")).as("zipped_col")) .select( col("zipped_col.cust_id").alias("customer_id"), col("zipped_col.cust_name").alias("customer_name") ) // 查看结果 resultDf.show()
Python 实现
from pyspark.sql import SparkSession from pyspark.sql.functions import arrays_zip, explode, col # 初始化SparkSession spark = SparkSession.builder.appName("FlattenCustomerData").getOrCreate() # 读取JSON数据(示例用字符串,实际可替换为文件路径) json_str = """{"data":{"cust_name":["cust_1","cust_2","cust_3","cust_4","cust_5"],"cust_id":[1,2,3,4,5]}}""" df = spark.read.json(spark.sparkContext.parallelize([json_str])) # 处理逻辑 result_df = df \ .select("data.cust_id", "data.cust_name") \ .withColumn("zipped", arrays_zip(col("cust_id"), col("cust_name"))) \ .select(explode(col("zipped")).alias("zipped_col")) \ .select( col("zipped_col.cust_id").alias("customer_id"), col("zipped_col.cust_name").alias("customer_name") ) # 展示结果 result_df.show()
关键函数说明
arrays_zip:将多个数组按索引位置配对,生成包含对应元素的结构体数组,比如把cust_id[0]和cust_name[0]组合成一个结构体explode:将数组中的每个结构体元素拆分成单独的行,完成扁平化操作
内容的提问来源于stack exchange,提问作者Abhishek Jain
相关产品推荐
相关产品推荐

