如何不依赖pandas将PySpark DataFrame指定两列转为键值对字典
PySpark DataFrame转键值对字典的原生实现方案
核心实现
你可以直接通过Spark RDD的原生算子完成转换,全程分布式并行处理,无需引入pandas依赖,示例代码如下:
# 假设你的源DataFrame命名为 df director_map = df.rdd \ .map(lambda row: (row["id"], row["director_name"])) \ .collectAsMap()
方案说明
- 中间的
map算子会在各个Executor节点并行执行,将每行数据转换为(id, director_name)的二元组,只有最后执行collectAsMap时才会将结果汇总到Driver节点生成字典,完全符合你并行化处理的要求 - 针对你提供的样例数据执行上述代码后,得到的输出为:
注:Python字典中字符串类型的值默认会带引号,如果你需要无引号的输出格式,可以自行遍历字典做字符串格式化处理{123: 'james cameron', 32: 'gore verbinski', 34: 'sam mendes', 2345: 'christopher nolan', 987: 'doug walker'}
注意事项
- 如果id列存在重复值,
collectAsMap会默认保留最后出现的对应director_name,你可以提前执行df = df.dropDuplicates(["id"])做去重处理,避免结果不符合预期 - 该操作最终会将全量数据拉取到Driver节点,如果你的数据量极大,建议先做过滤筛选再执行转换,避免Driver内存溢出
内容的提问来源于stack exchange,提问作者Meike
相关产品推荐
相关产品推荐

