【问题标题】:Spark number minus the same number in preview time火花数减去预览时相同的数
【发布时间】:2019-07-08 07:43:25
【问题描述】:

来源是一个csv文件:

id,sale,date
1,100,201901
1,105,201902
1,107,201904
1,108,201905
2,10,201901
2,11,201902
2,12,201904
2,13,201905

是关于一些产品的销售,1,100,201901表示从开始到日期201901,100个id为1的产品已经售出。

1,105,201902 表示从开始到日期 201902,105 个 id 为 1 的产品已售出。所以在 2019 年的第二个月,只有 5 个产品 1 售罄。

我期望的是使用 apache spark 为其添加一列,表示当月销售了多少产品。

预期结果是:

id,sale,date,inc
1,100,201901,0
1,105,201902,5
1,107,201904,2
1,108,201905,1
2,10,201901,1
2,11,201902,1
2,12,201904,1
2,13,201905,1

在实际情况下,它是批处理作业。

我已经厌倦了使用 join(下面的代码),我不确定是否应该使用 rollup、cube 或 accumulator。

如果我们每个月执行一个批处理作业,似乎还可以,问题是如果某个月我们忘记运行批处理作业,我们将在下个月运行它。

例如最后一行代码会显示:

| id|sale|  date|saleInc|
+---+----+------+-------+
|  1|2000|201901|   null|
|  1|2005|201902|      5|
|  1|2007|201903|      7|
+---+----+------+-------+

但实际上,201903 saleInc 应该是 2 而不是 7,应该是 2007 - 2005 而不是 2007-2000

这只是我的代码,你不能依赖它,你可以使用其他方式。

package incremental.test
import org.apache.spark.sql.SparkSession
import org.slf4j.LoggerFactory
import org.apache.spark.sql._
import org.apache.spark.sql.types._
import org.apache.spark.sql.functions._
import org.apache.spark.rdd.RDD
import org.apache.spark.SparkConf
//import com.qydata.stock.db._//a02z10 av1049 1yue29
import scala.reflect.api.materializeTypeTag
import org.apache.spark.SparkContext
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.internal.SQLConf.SHUFFLE_PARTITIONS
import scala.xml.dtd.Scanner
object D20190123 {
     def main(args: Array[String]){
       
       var sparkConf = new SparkConf().setMaster("local[1]")//.set("spark.default.parallelism","1").set("spark.streaming.blockInterval", "1").set("spark.shuffle.sort.bypassMergeThreshold", "1").set("spark.executor.cores", "1") .set("spark.executor.cores", "1")
           // .set("spark.cores.max", "1")
       val builder =  SparkSession.builder().config(sparkConf)//.enableHiveSupport()  
       val ss =  builder.getOrCreate()  
       import ss.implicits._
       ss.sessionState.conf.setConf(SHUFFLE_PARTITIONS, 1)
       var sc = ss.sparkContext
       sc.setLogLevel("error");


      
      var hive=Seq.empty[( Int,Int,String,Int)].toDF("id","sale","date","saleInc")
      println("====hive"); hive.show()
      
       val mongo1=Seq((1,2000,"201901")).toDF("id","sale","date");
       println("====mongo1"); mongo1.show()
       val newOfMongo1= mongo1.where('date>197001)
       println("====newOfMongo1"); newOfMongo1.show()
       val saleInHive1=hive.groupBy("id").agg('id,max('sale) as "mx").select($"id" as "hid",'mx)
       println("====saleInHive1");saleInHive1.show()
       val hiveAppend1=newOfMongo1.join(saleInHive1,'id==='hid,"left").withColumn("saleInc", 'sale-'mx)
       .select("id","sale","date","saleInc")
       println("====hiveAppend1");hiveAppend1.show()
       hive=hive.union(hiveAppend1)
       println("====hive"); hive.show()
       
      /* second batch may be missed
       * 
//       var hive=mongo1.select('id, 'sale,lit(0) as 'saleInc)//Seq((1,2000,0)).toDF("id","sale","saleInc")
       val mongo2=Seq((1,2000,"201901"),(1,2005,"201902")).toDF("id","sale","date")
       println("====mongo2"); mongo2.show()
       val newOfMongo2= mongo2.where('date>201901)
       println("====newOfMongo2"); newOfMongo2.show()
       val saleInHive2=hive.groupBy("id").agg('id,max('sale) as "mx").select($"id" as "hid",'mx)
       println("====saleInHive2");saleInHive2.show()
       val hiveAppend2=newOfMongo2.join(saleInHive2,'id==='hid,"left").withColumn("saleInc", 'sale-'mx)
       .select("id","sale","date","saleInc")
       println("====hiveAppend2");hiveAppend2.show()
       hive=hive.union(hiveAppend2)
       println("====hive"); hive.show()
       */
       
       val mongo3=Seq((1,2000,"201901"),(1,2005,"201902"),(1,2007,"201903")).toDF("id","sale","date")
       println("====mongo3"); mongo3.show()
       val newOfMongo3= mongo3.where('date>201901)//02
       println("====newOfMongo3"); newOfMongo3.show()
       val saleInHive3=hive.groupBy("id").agg('id,max('sale) as "mx").select($"id" as "hid",'mx)
       println("====saleInHive3"); saleInHive3.show()
       val hiveAppend3=newOfMongo3.join(saleInHive3,'id==='hid,"left").withColumn("saleInc", 'sale-'mx)
       .select("id","sale","date","saleInc")
       println("====hiveAppend3");hiveAppend3.show()
       hive=hive.union(hiveAppend3)
       
       println("====hive");hive.show()

     }
}

【问题讨论】:

    标签: apache-spark dataframe aggregate-functions


    【解决方案1】:

    所以你想要一个lag

    import org.apache.spark.sql.functions._
    import org.apache.spark.sql.expressions.Window
    
    df.withColumn( 
      "inc",
       $"sale" - lag($"sale", 1).over(Window.partitionBy($"id").orderBy($"date")))
    

    【讨论】:

    • 太棒了!原来可以这么简单
    猜你喜欢
    • 1970-01-01
    • 2017-05-13
    • 1970-01-01
    • 1970-01-01
    • 2017-10-02
    • 2017-11-25
    • 2022-11-12
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多