Spark:基于条件生成新列new1的DataFrame处理需求
解决Spark中新增列new1的需求
首先,我先结合你给出的输入和输出推断需求:给输入DataFrame新增new1列,该列的值为当前行之前所有行中满足fn非零且小于当前行fn的fn值的最大值;如果没有符合条件的值,则取整个DataFrame中最小的非零fn值(前四行fn=0,没有小于0的非零fn,但输出new1=4,而4是输入中最小的非零fn,所以做此推断)。
下面是具体的Spark Scala实现步骤:
1. 创建输入DataFrame
先还原你的输入数据:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.expressions.Window val spark = SparkSession.builder() .appName("NewColumnCalculation") .master("local[*]") // 生产环境请移除该行 .getOrCreate() // 构建输入DataFrame val inputDF = spark.createDataFrame(Seq( (1, 0), (1, 0), (2, 0), (2, 0), (3, 4), (5, 0), (5, 0), (6, 7) )).toDF("id", "fn")
2. 计算全局最小非零fn值
先获取整个DataFrame中最小的非零fn值,作为当前行之前无符合条件数据时的默认填充值:
val globalMinNonZeroFn = inputDF .filter(col("fn") =!= 0) .agg(min("fn")) .first() .getInt(0)
3. 使用窗口函数计算new1列
定义一个覆盖当前行之前所有行的窗口,在窗口内筛选出符合fn非零且小于当前行fn的记录,取它们的fn最大值;如果没有符合条件的记录,就用之前计算的全局最小非零fn值填充:
// 定义窗口:按输入顺序排序,包含当前行之前的所有行 val windowSpec = Window .orderBy(monotonically_increasing_id()) // 严格保证行顺序与输入一致 .rowsBetween(Long.MinValue, -1) // 窗口范围:从第一行到当前行的前一行 val outputDF = inputDF // 计算窗口内符合条件的fn最大值 .withColumn("prev_valid_fn", max(when(col("fn") =!= 0 && col("fn") < current_row("fn"), col("fn"))).over(windowSpec)) // 用全局最小非零fn填充null值,得到最终的new1列 .withColumn("new1", coalesce(col("prev_valid_fn"), lit(globalMinNonZeroFn))) // 删除临时辅助列 .drop("prev_valid_fn") // 查看结果 outputDF.show()
最终输出结果
运行上述代码后,会得到如下结果:
+---+---+----+ | id| fn|new1| +---+---+----+ | 1| 0| 4| | 1| 0| 4| | 2| 0| 4| | 2| 0| 4| | 3| 4| 4| | 5| 0| 4| | 5| 0| 4| | 6| 7| 4| +---+---+----+
补充说明
- 用
monotonically_increasing_id()排序是为了严格保证行顺序和输入一致,因为Spark默认不保证DataFrame的行顺序,除非显式指定排序规则。 - 如果你的需求中,当当前行之前没有符合条件的
fn时不需要填充默认值,而是保留null,只需要去掉coalesce函数,直接将prev_valid_fn重命名为new1即可。
内容的提问来源于stack exchange,提问作者vrreddy1234
相关产品推荐
相关产品推荐

