如何修改PySpark代码从map_from_arrays生成的Map获取匹配键值
需求说明
原PySpark代码可从数组列col1中找出与col2元素共有token最多的元素,生成match列。现需将col1替换为Map类型列col1_modify(或对应同索引的键值数组对),要求在匹配Map键的基础上,新增match_value列存储匹配键对应的值,最终输出match(匹配的键)和match_value(对应的值)两列。
原有核心代码:
c1_arr = F.col('col1') c2_arr = F.split(F.trim('col2'), '\s+') arr_of_struct = F.transform( c1_arr, lambda x: F.struct( F.size(F.array_intersect(c2_arr, F.split(F.trim(x), '\s+'))).alias('cnt'), x.alias('val'), ) ) top_val = F.sort_array(arr_of_struct, False)[0]生成
match列的代码:df = df.select( '*', F.when(((top_val['cnt'] > 0)),top_val['val']).alias('match'))
修改后的实现代码
针对Map类型列col1_modify的处理
先将Map转换为键值对结构体数组,计算每个键与col2的共有token数量,再筛选出匹配度最高的键值对:
from pyspark.sql import functions as F # 处理col2,拆分为token数组 c2_arr = F.split(F.trim(F.col('col2')), '\s+') # 将Map列转换为键值对结构体数组 map_kv_arr = F.map_entries(F.col('col1_modify')) # 遍历每个键值对,生成包含匹配计数、键、值的结构体数组 arr_of_struct = F.transform( map_kv_arr, lambda kv: F.struct( F.size(F.array_intersect(c2_arr, F.split(F.trim(kv['key']), '\s+'))).alias('cnt'), kv['key'].alias('match_key'), kv['value'].alias('match_val') ) ) # 按匹配计数降序排序,取第一个匹配度最高的结果 top_match = F.sort_array(arr_of_struct, False)[0] # 生成最终的match和match_value列,仅保留有效匹配(计数>0) df = df.select( '*', F.when(top_match['cnt'] > 0, top_match['match_key']).alias('match'), F.when(top_match['cnt'] > 0, top_match['match_val']).alias('match_value') )
针对同索引键值数组的兼容处理
如果是用两个同索引数组(比如col3为键数组,col4为对应值数组),可先转换为Map类型再复用上述逻辑:
# 将键值数组转换为Map列 col1_modify = F.map_from_arrays(F.col('col3'), F.col('col4')) # 后续逻辑与Map类型列处理完全一致
内容的提问来源于stack exchange,提问作者user15649753
相关产品推荐
相关产品推荐

