如何将PySpark map型DataFrame中key2的struct值赋给key1和key3
报错原因
你调用withColumn时传入的第二个参数value2是DataFrame类型,不符合该方法要求传入Column类型参数的规则,因此触发类型断言错误。
解决方法
以下提供两种适配不同场景的实现方案:
方案1:小数据量快速实现(value体积不大时使用)
直接将key2对应的value提取为Python端常量后做条件替换:
from pyspark.sql.functions import col, when, lit # 提取key2对应的value值 key2_value = myMap.filter(col('key') == 'key2').collect()[0]['value'] # 批量替换key1、key3的value result_df = myMap.withColumn( 'value', when(col('key').isin('key1', 'key3'), lit(key2_value)).otherwise(col('value')) )
验证结果:
result_df.show()
方案2:纯分布式实现(value体积大/数据量高时使用)
通过窗口函数将key2的value广播到所有行,不需要将值拉取到Driver端,避免大对象占用Driver内存:
from pyspark.sql.functions import col, when, max from pyspark.sql.window import Window # 定义全局窗口获取key2的value window_spec = Window.partitionBy() result_df = myMap.withColumn( 'key2_value', max(when(col('key') == 'key2', col('value'))).over(window_spec) ).withColumn( 'value', when(col('key').isin('key1', 'key3'), col('key2_value')).otherwise(col('value')) ).drop('key2_value')
内容的提问来源于stack exchange,提问作者doc
相关产品推荐
相关产品推荐

