如何在PySpark中映射RDD,将嵌套键合并为单个字符串?
PySpark RDD 嵌套键合并为字符串的实现方法
针对你提到的两种RDD结构,都可以通过map算子直接完成转换(无需flatMap,因为每个输入元素仅对应一个输出元素),具体实现如下:
情况1:键为二元组((A, B), Value)
输入RDD结构示例:
[(('DAI93865', 'FRO40251'), 1.0), (('GRO85051', 'FRO40251'), 0.999176276771005), (('GRO38636', 'FRO40251'), 0.9906542056074766), (('ELE12951', 'FRO40251'), 0.9905660377358491)]
转换代码:
# 假设原RDD名为rdd1 converted_rdd1 = rdd1.map(lambda x: (x[0][0] + x[0][1], x[1]))
转换后输出结构:
[('DAI93865FRO40251', 1.0), ('GRO85051FRO40251', 0.999176276771005), ('GRO38636FRO40251', 0.9906542056074766), ('ELE12951FRO40251', 0.9905660377358491)]
逻辑说明:通过lambda函数提取每个元素的键部分x[0](即二元组(A,B)),将A和B直接拼接成字符串,再与原元素的值x[1]组成新的键值对。
情况2:键为嵌套三元组(((A, B), C), Value)
输入RDD结构示例:
[((('DAI23334', 'ELE92920'), 'DAI62779'), 1.0), ((('DAI31081', 'GRO85051'), 'FRO40251'), 1.0)]
转换代码:
# 假设原RDD名为rdd2 converted_rdd2 = rdd2.map(lambda x: (x[0][0][0] + x[0][0][1] + x[0][1], x[1]))
转换后输出结构:
[('DAI23334ELE92920DAI62779', 1.0), ('DAI31081GRO85051FRO40251', 1.0)]
逻辑说明:逐层提取嵌套键中的A、B、C三个部分:x[0][0][0]是A,x[0][0][1]是B,x[0][1]是C,将三者拼接成单个字符串后与原数值配对。
可选:添加分隔符
如果需要在拼接时加入分隔符(比如下划线_),只需在拼接逻辑中插入即可:
# 情况1添加下划线分隔 converted_rdd1 = rdd1.map(lambda x: (x[0][0] + "_" + x[0][1], x[1])) # 情况2添加下划线分隔 converted_rdd2 = rdd2.map(lambda x: (x[0][0][0] + "_" + x[0][0][1] + "_" + x[0][1], x[1]))
内容的提问来源于stack exchange,提问作者Andrew Long
相关产品推荐
相关产品推荐

