PySpark下遍历配对RDD值并替换空字段的实现方法
解决方案
实现思路
- 调用
groupByKey按键分组后,先从当前键对应的所有值中筛选出所有字段非空的参考完整记录 - 遍历当前键下的所有值列表,将空字段替换为参考记录对应位置的取值
- 最后将分组处理后的结果展开为原有的单条结构即可
方式1:自定义处理函数(更易维护)
def fill_empty_fields(kv): key, values = kv value_list = list(values) # 提取所有字段都非空的参考记录 ref_record = next(v for v in value_list if all(field.strip() for field in v)) # 遍历替换空字段 processed = [] for v in value_list: filled_v = [v[i] if v[i].strip() else ref_record[i] for i in range(len(v))] processed.append((key, filled_v)) return processed # 完整调用链 result_rdd = paired_rdd.groupByKey().flatMap(fill_empty_fields)
方式2:单lambda写法(适配你要求的单lambda写入形式)
result_rdd = paired_rdd.groupByKey().flatMap( lambda kv: [(kv[0], [v[i] if v[i].strip() else next(r for r in list(kv[1]) if all(f.strip() for f in r))[i] for i in range(3)]) for v in list(kv[1])] )
调用result_rdd.collect()即可得到你给出的预期输出结果。
内容的提问来源于stack exchange,提问作者Aaron Gonzalez
相关产品推荐
相关产品推荐

