PySpark中Map类型字段扁平化及缺失键值Null填充方案咨询
解决PySpark Map字段扁平化时缺失键值的空值填充问题
问题背景
原始PySpark DataFrame(additional为Map类型字段):
+-------------+--------------+----+-------+-----------------------------------------------------------------------------------+ |empId |organization |h_cd|status |additional | +-------------+--------------+----+-------+-----------------------------------------------------------------------------------+ |FTE:56e662f |CATENA |0 |CURRENT|{hr_code -> 84534, bgc_val -> 170187, interviewPanel -> 6372, meetingId -> 3671} | |FTE:633e7bc |Data Science |0 |CURRENT|{hr_code -> 21036, bgc_val -> 170187, interviewPanel -> 764, meetingId -> 577} | |FTE:d9badd2 |CATENA |0 |CURRENT|{hr_code -> 60696, bgc_val -> 88770} | +-------------+--------------+----+-------+-----------------------------------------------------------------------------------+
目标扁平化结构:
+-------------+--------------+----+-------+------------+------------+-------------------+---------------+ |empId |organization |h_cd|status |hr_code |bgc_val |interviewPanel | meetingId | +-------------+--------------+----+-------+------------+------------+-------------------+---------------+ |FTE:56e662f |CATENA |0 |CURRENT|84534 |170187 |6372 |3671 | |FTE:633e7bc |Data Science |0 |CURRENT|21036 |170187 |764 |577 | |FTE:d9badd2 |CATENA |0 |CURRENT|60696 |88770 |Null |Null | +-------------+--------------+----+-------+------------+------------+-------------------+---------------+
现有RDD实现抛出KeyError: 'interviewPanel',原因是部分Map缺失指定键,需要自动填充Null的最优方案。
最优方案:使用PySpark内置函数(推荐)
直接用getItem方法访问Map字段,键不存在时自动返回Null,无需手动处理,且性能远优于RDD转换:
from pyspark.sql.functions import col flattened_df = df.select( "empId", "organization", "h_cd", "status", col("additional").getItem("hr_code").alias("hr_code"), col("additional").getItem("bgc_val").alias("bgc_val"), col("additional").getItem("interviewPanel").alias("interviewPanel"), col("additional").getItem("meetingId").alias("meetingId") ) flattened_df.show()
如果需要下划线格式的列名,修改alias参数即可,比如alias("interview_panel")。
备选方案:修正RDD逻辑处理键缺失
若必须使用RDD,用字典get方法指定默认值为None(PySpark自动转为Null),同时修正原代码中的字段名错误:
new_df = df.rdd.map(lambda x: ( x.empId, x.organization, x.h_cd, x.status, x.additional.get("hr_code"), x.additional.get("bgc_val"), x.additional.get("interviewPanel"), x.additional.get("meetingId") )).toDF([ "empId", "organization", "h_cd", "status", "hr_code", "bgc_val", "interviewPanel", "meetingId" ]) new_df.show()
注:原代码中x.data应为x.additional,x.category应为x.organization,需修正后运行。
内容的提问来源于stack exchange,提问作者1ksj8jdnu36flksf
相关产品推荐
相关产品推荐

