Spark Scala实现按规则递归分配水果至Plate的方法问询
Spark DataFrame 水果分配递归实现方案
需求说明
现有Spark DataFrame,col5存储待分配的水果数量,需将这些水果分配至col1至col4四个盘子中,规则如下:
- 每次找到
col1-col4中的最小值(若多个值同为最小,优先选择左侧盘子) - 将该盘子数值加1,同时
col5减1 - 重复操作直至
col5变为0
已实现单次迭代的Scala代码,需完成递归/循环迭代实现完整分配逻辑。
预期输出示例
初始数据:
plate1 plate2 plate3 plate4 fruits 1 2 3 4 3 6 7 8 1 2 2 4 6 8 5
迭代1:
2 2 3 4 2 6 7 8 2 1 3 4 6 8 4
迭代2:
3 2 3 4 1 (若多个plate最小值相同,取左侧优先) 6 7 8 3 0 4 4 6 8 3
迭代3:
3 3 3 4 0 6 7 8 3 0 5 4 6 8 2
迭代4:
3 3 3 4 0 6 7 8 3 0 5 5 6 8 1
迭代5:
3 3 3 4 0 6 7 8 3 0 6 5 6 8 0
单次迭代代码
val data = Seq( (1.0, 2.0, 3.0, 4.0, 3.0), (6.0, 7.0, 8.0, 1.0, 2.0), (2.0, 4.0, 6.0, 8.0, 5.0) ) val columns = List("col1", "col2", "col3", "col4", "col5") val df = spark.createDataFrame(data).toDF(columns: _*) val updatedColumns = columns.map { colName => functions.when(col(colName) === functions.least(columns.map(col): _*), col(colName) + 1) .otherwise(col(colName)) .alias(colName) } val updatedDF = df.select(updatedColumns: _*) updatedDF.show()
完整迭代实现方案
Spark分布式环境下无法直接使用常规递归函数,可通过循环迭代+终止条件判断模拟递归逻辑,核心是每次迭代后检查是否还有未分配的水果,直到所有行col5为0。
实现代码
import org.apache.spark.sql.functions._ import org.apache.spark.sql.DataFrame // 定义迭代分配函数 def iterateDistribute(df: DataFrame, plateCols: List[String], fruitCol: String): DataFrame = { // 检查是否还有待分配的水果 val hasRemaining = df.select(max(col(fruitCol))).head().getDouble(0) > 0.0 if (!hasRemaining) { df } else { val leastValue = least(plateCols.map(col): _*) // 生成每个盘子列的更新逻辑:确保左侧优先选择最小值列 val updatedPlateCols = plateCols.map { colName => // 判断当前列是否是该行第一个等于最小值的列 val isFirstMin = plateCols.takeWhile(_ != colName).forall(c => col(c) > leastValue) && col(colName) === leastValue when(isFirstMin, col(colName) + 1).otherwise(col(colName)).alias(colName) } // 更新水果数量列 val updatedFruitCol = when(col(fruitCol) > 0, col(fruitCol) - 1).otherwise(col(fruitCol)).alias(fruitCol) // 生成更新后的DataFrame val updatedDF = df.select(updatedPlateCols :+ updatedFruitCol: _*) // 递归触发下一次迭代 iterateDistribute(updatedDF, plateCols, fruitCol) } } // 初始化数据 val data = Seq( (1.0, 2.0, 3.0, 4.0, 3.0), (6.0, 7.0, 8.0, 1.0, 2.0), (2.0, 4.0, 6.0, 8.0, 5.0) ) val columns = List("col1", "col2", "col3", "col4", "col5") val plateCols = columns.take(4) val fruitCol = columns.last val df = spark.createDataFrame(data).toDF(columns: _*) // 执行分配逻辑 val finalDF = iterateDistribute(df, plateCols, fruitCol) finalDF.show()
关键细节说明
- 终止条件:通过
max(col5)判断是否还有未分配的水果,避免无效迭代 - 左侧优先逻辑:原单次迭代代码的
least函数无法区分多个最小值的顺序,新增isFirstMin判断确保优先选择最左侧的最小值列 - Spark惰性优化:递归调用实际是构建执行计划,Spark会将所有迭代步骤合并为一个Job执行,避免多次提交开销
内容的提问来源于stack exchange,提问作者user22006369
相关产品推荐
相关产品推荐

