如何编写可接受可变数量分区列参数的PySpark函数?
解决Spark get_recent_date函数分区列支持多值的问题
问题根源
你的代码存在三个关键问题:
- 参数拼写错误:函数定义里写的是
*partion_col,但后续引用的是partition_col,参数名不匹配 order_col是关键字参数,调用时必须用关键字传参,不能直接用位置参数Window.partitionBy()需要接收单个或多个列名参数,而*partition_col得到的是元组,直接传入会被Spark当成单个列名处理,触发字符串拼接类型错误
修正后的函数代码
from pyspark.sql.window import Window from pyspark.sql.functions import dense_rank, desc def get_recent_date(input_df, *partition_col, order_col): # 展开partition_col元组,传递多个列名给partitionBy w = Window().partitionBy(*partition_col)\ .orderBy(desc(order_col)) output_df = input_df.withColumn('DenseRank', dense_rank().over(w)) return output_df
正确调用方式
- 单个分区列:
get_recent_date(input, 'event_category', order_col='event_date') - 多个分区列:
get_recent_date(input, 'event_category', 'participant_category', order_col='event_date')
错误原因说明
- 参数拼写错误会导致
partition_col未被正确赋值,Spark尝试处理空值或错误的参数类型时引发异常 - 如果调用时不指定
order_col=关键字,最后一个位置参数会被归入*partition_col的元组中,导致order_col缺失,触发参数错误 partitionBy(*partition_col)通过*操作符展开元组,将每个元素作为独立的列名参数传递,避免了元组被当成单个列名处理的问题
内容的提问来源于stack exchange,提问作者SunflowerParty
相关产品推荐
相关产品推荐

