如何将JSON转换的PySpark嵌套DataFrame作为映射应用到另一个DataFrame?
如何将JSON转换的PySpark嵌套DataFrame作为映射应用到另一个DataFrame?
嘿,我懂你的需求啦,先理清楚你目前的情况:
你有一段这样的JSON数据:
{"main":{"honda":1,"toyota":2,"BMW":5,"Fiat":4}}
然后你用这段代码把它读进了PySpark:
car_map = spark.read.json('s3_path/car_map.json')
现在得到的是个嵌套结构的DataFrame,只有一个main列,里面装着各个汽车品牌对应的数值映射。你现在还有另一个已存在的DataFrame,想要把这个嵌套的映射用上去对吧?
我给你分步骤讲怎么实现:
第一步:把嵌套映射转成方便使用的键值对格式
首先得把main列里的嵌套结构拆出来,变成普通的键值对DataFrame,这样后续关联起来更顺手。用map_entries提取映射的键值对,再用explode把它展开就行:
from pyspark.sql.functions import explode, map_entries # 提取map的键值对并展开成行 car_mapping_df = car_map.select(explode(map_entries("main")).alias("car_entry")) # 把键和值拆成单独的列 car_mapping_df = car_mapping_df.select( car_mapping_df.car_entry.key.alias("car_brand"), car_mapping_df.car_entry.value.alias("brand_value") )
处理完之后,你就得到了一个两列的DataFrame:car_brand存品牌名,brand_value存对应的数值,这下就好和现有DataFrame关联了。
第二步:把映射应用到现有DataFrame
假设你的现有DataFrame叫existing_df,里面有一列是car_brand(和映射里的品牌名对应),直接用join关联就行:
# 关联现有DataFrame和映射DataFrame,用left join保留所有原始数据 result_df = existing_df.join(car_mapping_df, on="car_brand", how="left")
要是你现有DataFrame里的品牌列名字不一样,记得把on参数里的列名改成你实际的列名;how="left"是为了确保现有DataFrame里的数据不会因为没有映射值而丢失,你也可以根据需要换成inner(只保留有映射的数据)这类其他关联方式。
另一种更灵活的方式:直接用Map列查询
要是你不想拆成键值对,也可以直接用Map列来做查询,还能结合广播变量提升性能(大数据量场景特别有用):
from pyspark.sql.functions import broadcast, col, map_get # 先把映射DataFrame广播出去,减少数据传输 broadcast_car_map = broadcast(car_map.select("main")) # 交叉关联后用map_get取出对应的值 result_df = existing_df.join(broadcast_car_map, how="cross") \ .withColumn("brand_value", map_get(col("main"), col("car_brand"))) \ .drop("main")
这种方式不用拆映射结构,直接用原有的嵌套Map来查询,适合映射数据量不大的场景。
备注:内容来源于stack exchange,提问作者Chuck
相关产品推荐
相关产品推荐

