You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于另一列条件对列应用函数: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.15 08:13:15