【问题标题】:how to operate on two dataframes in spark scala如何在 spark scala 中对两个数据帧进行操作
【发布时间】:2019-01-03 15:49:16
【问题描述】:

我有两张桌子
1. 库存:那天我拥有的自行车数量
2. 销量:去年同一天售出的自行车数量

我想根据去年的数据预测我可以卖出多少天的股票。
例如,如果我有80 units of KTM bikes on <strong>2018-07-26</strong> ,那么我可以在3 days 中出售它们。
我该怎么做?

Ignore 2017-07-25 as this is more than one year old
on 2017-07-26 (80-15=65)
on 2017-07-27 (55-40=15)
on 2017-07-28 (15-50=-35)

so i can sell them in 3 days 

我可以在 1 天内卖出 RE

stocks table
+----+----------+-----+
|Bike|      Date|Units|
+----+----------+-----+
| KTM|2018-07-26|   80|
|  RE|2018-07-26|   40|
+----+----------+-----+

第二张表

sales table
+----+----------+-----------+
|Bike|      Date|Saled_units|
+----+----------+-----------+
| KTM|2017-07-25|         10|
| KTM|2017-07-26|         15|
| KTM|2017-07-27|         40|
| KTM|2017-07-28|         50|
| KTM|2017-07-29|         30|
|  RE|2017-07-26|         50|
+----+----------+-----------+

如何使用 spark SQL 做到这一点?

【问题讨论】:

    标签: scala apache-spark apache-spark-sql


    【解决方案1】:

    您可以按照以下步骤获得所需的输出

    1.创建DataFrames并将股票列"Bike"重命名为"Bike_Temp"

    2.GroupBike 的销售额dataframe 的数据,并使用aggregate 函数将剩余的列收集为array&lt;struct&lt;string,int&gt;&gt;

    3.最后通过Bikecolumn加入数据框agrouped_salesstocks_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|
    +----+-------------+
    

    我出于测试目的修改了日期,并假设您将始终检查当前日期。希望对你有帮助

    【讨论】:

      猜你喜欢
      • 2018-04-17
      • 2023-04-07
      • 1970-01-01
      • 1970-01-01
      • 2020-08-20
      • 2020-08-20
      • 2021-06-06
      • 2016-05-30
      • 1970-01-01
      相关资源
      最近更新 更多