您可以按照以下步骤获得所需的输出
1.创建DataFrames并将股票列"Bike"重命名为"Bike_Temp"
2.GroupBike 的销售额dataframe 的数据,并使用aggregate 函数将剩余的列收集为array<struct<string,int>>
3.最后通过Bikecolumn加入数据框agrouped_sales和stocks_temp,如下所示
val sales=Seq(("KTM","2017-07-26",10),("KTM","2017-07-27",15),("KTM","2017-07-28",40),("KTM","2017-07-29",50),("KTM","2017-07-30",30),("RE","2017-07-27",50)).toDF("Bike","Date","Saled_units")
val stocks=Seq(("KTM","2018-07-27",80),("RE","2018-07-27",40)).toDF("Bike","Date","Units")
val stocks_temp=stocks.withColumnRenamed("Bike","Bike_Temp")
val grouped_sales=sales.groupBy("Bike").agg(collect_list(struct("Date","Saled_units")).as("date_saled_units"))
val joineddf=grouped_sales.join(stocks_temp,grouped_sales("Bike")===stocks_temp("Bike_Temp"))
4.创建一个udf并注册它并将逻辑写入countudf中的天数
import org.apache.spark.sql.functions._
import org.apache.spark.sql.Row
import java.text.SimpleDateFormat
import java.util.Calendar
import spark.implicits._
def getNoOfDays(arr:Seq[Row],stock_units:Int):Int={
val dateFormat=new SimpleDateFormat("yyyy-MM-dd")
val cal = Calendar.getInstance()
//Setting last year date
cal.add(Calendar.YEAR, -1)
val previousYearTodaysDate =dateFormat.parse(dateFormat.format(cal.getTime()))
//Iterating over collected list and converting arr<struct<>> to Seq[Date,Int]
val dateBikeArr=arr.map(row=>(dateFormat.parse(row.getAs("Date").toString),row.getAs("Saled_units").toString.toInt))
var noOfDays=0
var totalUnits=0
//In all the dates checkng for date which is greater than or equal to today's date and summing the total bikes
for(tup<-dateBikeArr){
if(tup._1.compareTo(previousYearTodaysDate)>=0)
{
totalUnits=totalUnits+tup._2
noOfDays=noOfDays+1
if(totalUnits>=stock_units)
return noOfDays
}
}
return noOfDays
}
val getNoOfDaysUDF=udf(getNoOfDays _)
5.在joineddf 上选择column Bike 和date_saled_units 并拨打udf
val outputdf=joineddf.select(joineddf.col("Bike"),getNoOfDaysUDF(joineddf.col("date_saled_units"),joineddf.col("Units")).alias("days_required"))
6.样本输出:
+----+-------------+
|Bike|days_required|
+----+-------------+
| RE| 1|
| KTM| 3|
+----+-------------+
我出于测试目的修改了日期,并假设您将始终检查当前日期。希望对你有帮助