基于另一列条件对列应用函数:Pandas转PySpark实现
用PySpark实现Pandas条件列处理逻辑
针对你提供的Pandas代码逻辑,在PySpark中可以通过内置函数高效实现(避免使用UDF以提升大数据量下的性能),具体代码如下:
from pyspark.sql import functions as F df = df.withColumn( "cid_cm11_new", F.when( F.col("code") == "O", # 处理id列:取.前的部分,左补0到15位(对应Python的rjust(15,'0')) F.lpad(F.split(F.col("id"), "\\.").getItem(0), 15, "0") ).otherwise( # 处理id2列:取.前的部分,左补0到9位(对应Python的rjust(9,'0')) F.lpad(F.split(F.col("id2"), "\\.").getItem(0), 9, "0") ) )
代码说明:
F.when(...).otherwise(...):对应Pandas lambda中的条件判断逻辑,当code等于"O"时执行第一个分支,否则执行第二个分支F.split(col, "\\."):按.拆分字符串(PySpark中需转义.,所以用\\.),返回数组类型结果.getItem(0):取拆分后数组的第一个元素,等价于Pandas里的split('.')[0]F.lpad(str, length, pad):在字符串左侧填充指定字符(此处为"0")到目标长度,效果和Python的str.rjust(length, pad)完全一致
注意事项:
- 如果
id或id2不是字符串类型,需要先通过F.col("id").cast("string")转换为字符串后再执行拆分操作 - 内置函数比自定义UDF(
F.udf)性能高出很多,尤其是处理超大规模数据集时,能充分利用Spark的分布式优化能力
内容的提问来源于stack exchange,提问作者ggupta
相关产品推荐
相关产品推荐

