如何在PySpark中基于字符计数创建派生属性,统计f/g/h生成新列
PySpark实现方案
完全对齐你原有pandas逻辑,使用PySpark原生内置函数实现,全程无自定义UDF,大数据量下执行效率远高于UDF方案:
- 先导入依赖的函数模块
from pyspark.sql import functions as F
- 执行列处理逻辑
# 此处填写你需要批量处理的所有目标列名 collist = ['A'] for col_name in collist: # 定义判断条件:列值非空 且 包含f/g/h任意字符 has_target = F.col(col_name).isNotNull() & F.col(col_name).rlike("[fgh]") # 统计f数量的新列 df = df.withColumn( f"{col_name}_f", F.when( has_target, F.length(F.regexp_replace(F.col(col_name), "[^f]", "")) ).otherwise( F.when(F.col(col_name).isNull(), None).otherwise(F.col(col_name).cast("int")) ) ) # 统计g数量的新列 df = df.withColumn( f"{col_name}_g", F.when( has_target, F.length(F.regexp_replace(F.col(col_name), "[^g]", "")) ).otherwise( F.when(F.col(col_name).isNull(), None).otherwise(F.col(col_name).cast("int")) ) ) # 统计h数量的新列 df = df.withColumn( f"{col_name}_h", F.when( has_target, F.length(F.regexp_replace(F.col(col_name), "[^h]", "")) ).otherwise( F.when(F.col(col_name).isNull(), None).otherwise(F.col(col_name).cast("int")) ) )
逻辑说明
- 空值处理和原有逻辑完全一致,空值直接保留
- 只要原列值包含f/g/h任意一个字符,就统计对应字符的出现次数
- 原列值不包含f/g/h任意字符时,自动转成整数保留原值
- 所有运算均为Spark原生算子,分布式执行效率高,适合超大数据量批量处理
内容的提问来源于stack exchange,提问作者Abhi
相关产品推荐
相关产品推荐

