如何在PySpark中按条件创建新列?含加权均值与平方列示例
实现Spark DataFrame的条件列生成需求
没问题,我来帮你搞定这个需求!咱们一步步来:
首先,你需要先计算age的加权平均值(以weights列为权重),这里要注意过滤掉age为null的行,避免影响计算结果。然后用Spark的when/otherwise函数根据age的条件分别生成对应的列。
完整代码实现
import findspark findspark.init() from pyspark.sql import SparkSession spark = SparkSession.builder.getOrCreate() sc = spark.sparkContext from pyspark.sql import functions as sf from pyspark.sql.functions import col, when, lit # 初始化你的DataFrame ddf = spark.createDataFrame( [[None, 'Michael',2], [30, 'Andy',3], [19, 'Justin',4], [30, 'James Dr No From Russia with Love Bond',6]], schema=['age', 'name','weights'] ) # 第一步:计算age的加权平均值(过滤age为空的行) weighted_avg_age = ddf.filter(col("age").isNotNull()) \ .agg( (sf.sum(col("age") * col("weights")) / sf.sum(col("weights"))).alias("weighted_avg") ).collect()[0]["weighted_avg"] # 第二步:添加条件列 result_df = ddf.withColumn( "weighted_age", # 当age>29时,填入加权平均值,否则为null when(col("age") > 29, lit(weighted_avg_age)).otherwise(lit(None)) ).withColumn( "age_squared", # 当age<=29时,计算age的平方,否则为null when(col("age") <= 29, col("age") ** 2).otherwise(lit(None)) ) # 查看结果 result_df.show(truncate=False)
输出结果
+----+----------------------------------------+-------+------------------+-----------+ |age |name |weights|weighted_age |age_squared| +----+----------------------------------------+-------+------------------+-----------+ |null|Michael |2 |null |null | |30 |Andy |3 |27.77777777777778 |null | |19 |Justin |4 |null |361 | |30 |James Dr No From Russia with Love Bond |6 |27.77777777777778 |null | +----+----------------------------------------+-------+------------------+-----------+
关键说明
- 加权平均值的计算:我们用
sum(age * weights)除以sum(weights),这是标准的加权平均公式; when/otherwise函数:这是Spark中实现条件逻辑的常用方法,能轻松根据列值动态生成新列;- 处理null值:因为第一行age为null,所以这一行的两个新列都会显示null,符合逻辑。
内容的提问来源于stack exchange,提问作者quant
相关产品推荐
相关产品推荐

