PySpark嵌套when结合isin方法抛出数据类型不匹配异常咨询
问题结论
这不是PySpark对嵌套when的功能限制,属于写法错误,核心是对isinAPI的参数规则理解偏差。
- 报错根因:
col("SUBCAT")是字符串类型,你传入isin的参数是when...otherwise返回的数组类型列。isin的执行逻辑是「判断左侧值是否和传入的任意一个参数相等」,不会自动把单个数组类型参数拆成独立元素做匹配,因此触发了string != array<string>的类型不匹配错误。 - 第一版代码正常运行的原因:直接传入Python原生列表
["AA", "AB", "BA", "BB"]时,PySpark会在本地直接把Python列表拆分为多个独立的字面量参数传给isin,等价于isin("AA","AB","BA","BB"),所有参数都是字符串类型,和左侧字段类型匹配,因此不会报错。
修复方案
两种常用实现都可以满足按CAT值动态匹配SUBCAT白名单的需求:
方案1:用array_contains替代isin
直接判断SUBCAT是否包含在动态生成的目标数组中,逻辑最简洁:
df_grouped = df \ .groupby("ID") \ .agg( collect_set( when( array_contains( when(col("CAT") == lit("A"), array(lit("AA"), lit("AB"))) .otherwise(array(lit("BA"), lit("BB"))), col("SUBCAT") ), struct(*[df[columnName] for columnName in restOfColumns]) ) ) )
方案2:把匹配逻辑拆分到when的分支内
不需要生成中间数组,直接在每个when分支里写对应场景下的isin判断,性能和方案1一致:
subcat_match_rule = when( col("CAT") == "A", col("SUBCAT").isin("AA", "AB") ).otherwise( col("SUBCAT").isin("BA", "BB") ) df_grouped = df \ .groupby("ID") \ .agg( collect_set( when(subcat_match_rule, struct(*[df[columnName] for columnName in restOfColumns])) ) )
注意:不推荐为了用
isin强行把数组列拆成多行做关联,这种写法shuffle开销更大,性能远低于上面两种方案。
内容的提问来源于stack exchange,提问作者gorrch
相关产品推荐
相关产品推荐

