This is an automated email from the ASF dual-hosted git repository.

iwasakims pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/bigtop.git


The following commit(s) were added to refs/heads/master by this push:
     new e0b7f9e31 BIGTOP-4368. Fix bigpetstore-spark to be built with the 
recent version of Spark and related libraries. (#1330)
e0b7f9e31 is described below

commit e0b7f9e310e7a1de238961c2df08c9cb74f111b0
Author: Kengo Seki <[email protected]>
AuthorDate: Fri Feb 21 11:26:41 2025 +0900

    BIGTOP-4368. Fix bigpetstore-spark to be built with the recent version of 
Spark and related libraries. (#1330)
---
 bigtop-bigpetstore/bigpetstore-spark/README.md     | 16 ++---
 bigtop-bigpetstore/bigpetstore-spark/build.gradle  | 68 +++++-----------------
 .../bigpetstore-spark/settings.gradle              |  0
 .../spark/analytics/PetStoreStatistics.scala       | 35 +++++------
 .../spark/analytics/RecommendProducts.scala        |  9 ++-
 .../bigpetstore/spark/datamodel/DataModel.scala    |  6 +-
 .../bigpetstore/spark/datamodel/IOUtils.scala      | 22 +++----
 .../org/apache/bigpetstore/spark/etl/ETL.scala     | 17 ++----
 .../bigpetstore/spark/generator/SparkDriver.scala  | 26 ++++-----
 .../bigpetstore/spark/TestFullPipeline.scala       | 18 ++----
 .../spark/analytics/AnalyticsSuite.scala           | 16 ++---
 .../bigpetstore/spark/datamodel/IOUtilsSuite.scala | 15 ++---
 .../apache/bigpetstore/spark/etl/ETLSuite.scala    | 29 +++++----
 .../spark/generator/SparkDriverSuite.scala         | 13 ++---
 14 files changed, 99 insertions(+), 191 deletions(-)

diff --git a/bigtop-bigpetstore/bigpetstore-spark/README.md 
b/bigtop-bigpetstore/bigpetstore-spark/README.md
index 881718680..93897cca5 100644
--- a/bigtop-bigpetstore/bigpetstore-spark/README.md
+++ b/bigtop-bigpetstore/bigpetstore-spark/README.md
@@ -23,7 +23,7 @@ providing generators for synthetic transaction data and 
pipelines for
 processing that data.  Each ecosystems has its own version of the
 application.
 
-The Spark application currently builds against Spark 1.3.0.
+The Spark application currently builds against Spark 3.5.4.
 
 Architecture
 ------------
@@ -49,7 +49,7 @@ The data generator creates a dirty CSV file containing the 
following fields:
 * Customer City: String
 * Customer State: String
 * Transaction ID: Int
-* Transation Date Time: String (e.g., "Tue Nov 03 01:08:11 EST 2014")
+* Transaction Date Time: String (e.g., "Tue Nov 03 01:08:11 EST 2014")
 * Transaction Product: String (e.g., "category=dry cat food;brand=Feisty 
Feline;flavor=Chicken & Rice;size=14.0;per_unit_cost=2.14;")
 
 Note that the transaction ID is unique only per customer -- the customer and 
transaction IDs form a unique composite key.
@@ -64,7 +64,7 @@ internal structured data model is defined as input for the 
analytics components:
 * Transaction(customerId: Long, transactionId: Long, storeId: Long, dateTime: 
java.util.Calendar, productId: Long)
 
 The ETL stage parses and cleans up the dirty CSV and writes out RDDs for each 
data type in the data model, serialized using
-the `saveAsObjectFile()` method.  The analytics components can use the 
`IOUtils.load()` method to de-serialize the structured
+the `saveAsObjectFile()` method. The analytics components can use the 
`IOUtils.load()` method to de-serialize the structured
 data.
 
 Running Tests
@@ -84,14 +84,14 @@ Build a fat jar as follows:
 gradle clean shadowJar
 ```
 
-This will produce a jar file under `build/libs` (referred to as 
`bigpetstore-spark-X.jar`).  You can then
+This will produce a jar file under `build/libs` (referred to as 
`bigpetstore-spark-X.jar`). You can then
 use this jar to run a Spark job as follows:
 
 ```
 spark-submit --master local[2] --class 
org.apache.bigtop.bigpetstore.spark.generator.SparkDriver 
bigpetstore-spark-X.jar generated_data/ 10 1000 365.0 345
 ```
 
-You will need to change the master if you want to run on a cluster.  The last 
five parameters control the output directory,
+You will need to change the master if you want to run on a cluster. The last 
five parameters control the output directory,
 the number of stores, the number of customers, simulation length (in days), 
and the random seed (which is optional).
 
 
@@ -117,7 +117,7 @@ spark-submit --master local[2] --class 
org.apache.bigtop.bigpetstore.spark.etl.S
 Running the SparkSQL component
 -------------------------------
 
-Once ETL'd we can now process the data and do analytics on it.  The 
DataModel.scala class itself is used to read/write classes
+Once ETL'd we can now process the data and do analytics on it. The 
DataModel.scala class itself is used to read/write classes
 from files.  To run the analytics job, which outputs a JSON file at the end, 
you now will run the following:
 
 ```
@@ -137,7 +137,7 @@ This will output a JSON file to the current directory, 
which has formatting (app
 {
    "totalTransaction":34586,
    "transactionsByZip":[
-  
{"count":64,"productId":54,"zipcode":"94583"},{"count":38,"productId":18,"zipcode":"34761"},
+   
{"count":64,"productId":54,"zipcode":"94583"},{"count":38,"productId":18,"zipcode":"34761"},
    
{"count":158,"productId":14,"zipcode":"11368"},{"count":66,"productId":46,"zipcode":"33027"},
    
{"count":52,"productId":27,"zipcode":"94583"},{"count":84,"productId":19,"zipcode":"33027"},
    
{"count":143,"productId":0,"zipcode":"94583"},{"count":58,"productId":41,"zipcode":"72715"},
@@ -161,7 +161,7 @@ This will output a JSON file to the current directory, 
which has formatting (app
 ```
 
 Of course, the above data is for a front end web app which will display 
charts/summary stats of the transactions.
-Keep tracking Apache BigTop for updates on this front !
+Keep tracking Apache BigTop for updates on this front!
 
 Running the Product Recommendation Component
 --------------------------------------------
diff --git a/bigtop-bigpetstore/bigpetstore-spark/build.gradle 
b/bigtop-bigpetstore/bigpetstore-spark/build.gradle
index 6d5a9574b..34174a992 100644
--- a/bigtop-bigpetstore/bigpetstore-spark/build.gradle
+++ b/bigtop-bigpetstore/bigpetstore-spark/build.gradle
@@ -24,13 +24,12 @@ apply plugin: "scala"
 apply plugin: 'com.github.johnrengelman.shadow'
 
 buildscript {
-  repositories { jcenter() }
+  repositories { gradlePluginPortal() }
   dependencies {
-    classpath 'com.github.jengelman.gradle.plugins:shadow:1.0.2'
+    classpath 
"com.github.johnrengelman.shadow:com.github.johnrengelman.shadow.gradle.plugin:5.2.0"
   }
 }
 
-
 // Read the groupId and version properties from the "parent" bigtop project.
 // It would be better if there was some better way of doing this. Howvever,
 // at this point, we have to do this (or some variation thereof) since gradle
@@ -45,17 +44,13 @@ def setProjectProperties() {
 setProjectProperties()
 description = """"""
 
-// We are using 1.7 as gradle can't play well when java 8 and scala are 
combined.
-// There is an open issue here: http://issues.gradle.org/browse/GRADLE-3023
-// There is talk of this being resolved in the next version of gradle. Till 
then,
-// we are stuck with java 7. But we do have scala if we want more syntactic 
sugar.
-sourceCompatibility = 1.7
-targetCompatibility = 1.7
+sourceCompatibility = 1.8
+targetCompatibility = 1.8
 
 // Specify any additional project properties.
 ext {
-    sparkVersion = "1.3.0"
-    scalaVersion = "2.10"
+    sparkVersion = "3.5.4"
+    scalaVersion = "2.12"
 }
 
 shadowJar {
@@ -63,10 +58,8 @@ shadowJar {
 }
 
 repositories {
+    mavenLocal()
     mavenCentral()
-    maven {
-        url "http://dl.bintray.com/rnowling/bigpetstore";
-    }
 }
 
 tasks.withType(AbstractCompile) {
@@ -74,12 +67,6 @@ tasks.withType(AbstractCompile) {
     options.compilerArgs << "-Xlint:all"
 }
 
-tasks.withType(ScalaCompile) {
-    // Enables incremental compilation.
-    // 
http://www.gradle.org/docs/current/userguide/userguide_single.html#N12F78
-    scalaCompileOptions.useAnt = false
-}
-
 tasks.withType(Test) {
     testLogging {
         // Uncomment this if you want to see the console output from the tests.
@@ -100,43 +87,20 @@ sourceSets {
     }
 }
 
-
-// To see the API that is being used here, consult the following docs
-// 
http://www.gradle.org/docs/current/dsl/org.gradle.api.artifacts.ResolutionStrategy.html
-def updateDependencyVersion(dependencyDetails, dependencyString) {
-    def parts = dependencyString.split(':')
-    def group = parts[0]
-    def name = parts[1]
-    def version = parts[2]
-    if (dependencyDetails.requested.group == group
-            && dependencyDetails.requested.name == name) {
-        dependencyDetails.useVersion version
-    }
-}
-
-
 dependencies {
-    compile "org.apache.spark:spark-core_${scalaVersion}:${sparkVersion}"
-    compile "org.apache.spark:spark-mllib_${scalaVersion}:${sparkVersion}"
-    compile 
"org.apache.spark:spark-network-shuffle_${scalaVersion}:${sparkVersion}"
-    compile "org.apache.spark:spark-sql_${scalaVersion}:${sparkVersion}"
-    compile "org.apache.spark:spark-graphx_${scalaVersion}:${sparkVersion}"
-    compile "org.apache.spark:spark-hive_${scalaVersion}:${sparkVersion}"
+    compile("org.apache.spark:spark-core_${scalaVersion}:${sparkVersion}")
+    compile("org.apache.spark:spark-mllib_${scalaVersion}:${sparkVersion}")
+    compile("org.apache.spark:spark-sql_${scalaVersion}:${sparkVersion}")
     compile "com.github.rnowling.bigpetstore:bigpetstore-data-generator:0.2.1"
-    compile "joda-time:joda-time:2.7"
-    compile "org.json4s:json4s-jackson_2.10:3.1.0"
+    compile "joda-time:joda-time:2.13.1"
+    compile "org.json4s:json4s-jackson_${scalaVersion}:3.6.12"
 
-    testCompile "junit:junit:4.11"
-    testCompile "org.hamcrest:hamcrest-all:1.3"
-    testCompile "org.scalatest:scalatest_${scalaVersion}:2.2.1"
-    testCompile "joda-time:joda-time:2.7"
+    testCompile "junit:junit:4.13.2"
+    testCompile "org.scalatest:scalatest_${scalaVersion}:3.2.19"
+    testCompile "org.scalatestplus:junit-4-13_${scalaVersion}:3.2.19.0"
+    testCompile "joda-time:joda-time:2.13.1"
 }
 
-task listJars << {
-    configurations.shadow.each { println it.name }
-}
-
-
 eclipse {
     classpath {
         // Comment out the following two lines if you want to generate an 
eclipse project quickly.
diff --git a/bigtop-bigpetstore/bigpetstore-spark/settings.gradle 
b/bigtop-bigpetstore/bigpetstore-spark/settings.gradle
new file mode 100644
index 000000000..e69de29bb
diff --git 
a/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/analytics/PetStoreStatistics.scala
 
b/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/analytics/PetStoreStatistics.scala
index e7e5b087e..98bbcb42e 100644
--- 
a/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/analytics/PetStoreStatistics.scala
+++ 
b/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/analytics/PetStoreStatistics.scala
@@ -18,23 +18,18 @@
 package org.apache.bigtop.bigpetstore.spark.analytics
 
 import java.io.File
-import java.sql.Timestamp
 
-import scala.Nothing
+import scala.language.postfixOps
 
 import org.apache.spark.sql._
 import org.apache.spark.{SparkContext, SparkConf}
-import org.apache.spark.SparkContext._
 import org.apache.spark.rdd._
 
-import org.joda.time.DateTime
-import org.json4s.JsonDSL.WithBigDecimal._
-
 import org.apache.bigtop.bigpetstore.spark.datamodel._
 
 object PetStoreStatistics {
 
-    private def printUsage() {
+    private def printUsage(): Unit = {
       val usage: String = "BigPetStore Analytics Module." +
       "\n" +
       "Usage: spark-submit ... inputDir outputFile\n " +
@@ -119,26 +114,24 @@ GROUP BY productId, zipcode""")
   def runQueries(r:(RDD[Location], RDD[Store], RDD[Customer], RDD[Product],
     RDD[Transaction]), sc: SparkContext): Statistics = {
 
-    val sqlContext = new org.apache.spark.sql.SQLContext(sc)
-    import sqlContext._
-    import sqlContext.implicits._
+    val spark = SparkSession.builder().config(sc.getConf).getOrCreate()
 
     // Transform the Non-SparkSQL Calendar into a SparkSQL-friendly field.
     val mappableTransactions:RDD[TransactionSQL] =
       r._5.map { trans => trans.toSQL() }
 
-    r._1.toDF().registerTempTable("Locations")
-    r._2.toDF().registerTempTable("Stores")
-    r._3.toDF().registerTempTable("Customers")
-    r._4.toDF().registerTempTable("Product")
-    mappableTransactions.toDF().registerTempTable("Transactions")
+    spark.createDataFrame(r._1).toDF().createOrReplaceTempView("Locations")
+    spark.createDataFrame(r._2).createOrReplaceTempView("Stores")
+    spark.createDataFrame(r._3).createOrReplaceTempView("Customers")
+    spark.createDataFrame(r._4).createOrReplaceTempView("Product")
+    
spark.createDataFrame(mappableTransactions).createOrReplaceTempView("Transactions")
 
 
-    val txByMonth = queryTxByMonth(sqlContext)
-    val txByProduct = queryTxByProduct(sqlContext)
-    val txByProductZip = queryTxByProductZip(sqlContext)
+    val txByMonth = queryTxByMonth(spark.sqlContext)
+    val txByProduct = queryTxByProduct(spark.sqlContext)
+    val txByProductZip = queryTxByProductZip(spark.sqlContext)
 
-    return Statistics(
+    Statistics(
       txByMonth.map { s => s.count }.reduce(_+_),  // Total number of 
transactions
       txByMonth,
       txByProduct,
@@ -150,7 +143,7 @@ GROUP BY productId, zipcode""")
     * We keep a "run" method which can be called easily from tests and also is 
used by main.
     */
     def run(txInputDir:String, statsOutputFile:String,
-      sc:SparkContext) {
+      sc:SparkContext): Unit = {
 
       System.out.println("Running w/ input = " + txInputDir)
 
@@ -164,7 +157,7 @@ GROUP BY productId, zipcode""")
       System.out.println("Output JSON Stats stored : " + statsOutputFile)
     }
 
-  def main(args: Array[String]) {
+  def main(args: Array[String]): Unit = {
     // Get or else : On failure (else) we exit.
     val (inputPath,outputPath) = parseArgs(args)
 
diff --git 
a/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/analytics/RecommendProducts.scala
 
b/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/analytics/RecommendProducts.scala
index 5ac0fbaf5..fb9b8a7a1 100644
--- 
a/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/analytics/RecommendProducts.scala
+++ 
b/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/analytics/RecommendProducts.scala
@@ -19,7 +19,6 @@ package org.apache.bigtop.bigpetstore.spark.analytics
 
 import org.apache.bigtop.bigpetstore.spark.datamodel._
 import org.apache.spark.{SparkContext, SparkConf}
-import org.apache.spark.SparkContext._
 import org.apache.spark.rdd._
 import org.apache.spark.mllib.recommendation._
 
@@ -29,7 +28,7 @@ case class PRParameters(inputDir: String, outputFile: String)
 
 object RecommendProducts {
 
-  private def printUsage() {
+  private def printUsage(): Unit = {
     val usage = "BigPetStore Product Recommendation Module\n" +
     "\n" +
     "Usage: transformed_data recommendations\n" +
@@ -51,7 +50,7 @@ object RecommendProducts {
 
   def prepareRatings(tx: RDD[Transaction]): RDD[Rating] = {
     val productPairs = tx.map { t => ((t.customerId, t.productId), 1) }
-    val pairCounts = productPairs.reduceByKey { case (v1, v2) => v1 + v2 }
+    val pairCounts = productPairs.reduceByKey { (v1, v2) => v1 + v2 }
     val ratings = pairCounts.map { p => Rating(p._1._1.toInt, p._1._2.toInt, 
p._2) }
 
     ratings
@@ -83,7 +82,7 @@ object RecommendProducts {
     */
   def run(txInputDir: String, recOutputFile: String, sc: SparkContext,
   nIterations: Int = 20, alpha: Double = 40.0, rank:Int = 10, lambda: Double = 
1.0,
-  nRecommendations: Int = 5) {
+  nRecommendations: Int = 5): Unit = {
 
     println("input : " + txInputDir)
     println(sc)
@@ -104,7 +103,7 @@ object RecommendProducts {
     IOUtils.saveLocalAsJSON(new File(recOutputFile), prodRec)
   }
 
-  def main(args: Array[String]) {
+  def main(args: Array[String]): Unit = {
     val params: PRParameters = parseArgsOrDie(args)
 
     val conf = new SparkConf().setAppName("BPS Product Recommendations")
diff --git 
a/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/datamodel/DataModel.scala
 
b/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/datamodel/DataModel.scala
index 4ec30e2d3..71edd7f5b 100644
--- 
a/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/datamodel/DataModel.scala
+++ 
b/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/datamodel/DataModel.scala
@@ -20,10 +20,7 @@ package org.apache.bigtop.bigpetstore.spark.datamodel
 import java.sql.Timestamp
 import java.util.Calendar
 
-import org.apache.spark.sql
 import org.joda.time.DateTime
-import org.json4s.CustomSerializer
-import org.json4s.JsonAST.{JString, JField, JInt, JObject}
 
 /**
  * Statistics phase.  Represents JSON for a front end.
@@ -58,8 +55,7 @@ case class Transaction(customerId: Long, transactionId: Long, 
storeId: Long, dat
    */
   def toSQL(): TransactionSQL = {
     val dt = new DateTime(dateTime)
-    val ts = new Timestamp(dt.getMillis)
-    return TransactionSQL(customerId,transactionId,storeId,
+    TransactionSQL(customerId,transactionId,storeId,
       new Timestamp(
         new DateTime(dateTime).getMillis),
         productId,
diff --git 
a/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/datamodel/IOUtils.scala
 
b/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/datamodel/IOUtils.scala
index 1485ac318..d28756dc6 100644
--- 
a/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/datamodel/IOUtils.scala
+++ 
b/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/datamodel/IOUtils.scala
@@ -18,22 +18,14 @@
 package org.apache.bigtop.bigpetstore.spark.datamodel
 
 import java.io.File
-import java.util.Date
-import java.nio.file.{Path, Paths, Files}
+import java.nio.file.Files
 import java.nio.charset.StandardCharsets
 
-import org.apache.spark.{SparkContext, SparkConf}
-import org.apache.spark.SparkContext._
+import org.apache.spark.SparkContext
 import org.apache.spark.rdd._
 
-import org.apache.bigtop.bigpetstore.spark.datamodel._
-import org.json4s.JsonDSL._
-import org.json4s.JsonDSL.WithDouble._
-import org.json4s.JsonDSL.WithBigDecimal._
 import org.json4s.jackson.Serialization
 import org.json4s._
-import org.json4s.JsonDSL._
-import org.json4s.jackson.JsonMethods._
 import org.json4s.jackson.Serialization.{read, write}
 
 /**
@@ -60,7 +52,7 @@ object IOUtils {
     */
   def save(outputDir: String, locationRDD: RDD[Location],
     storeRDD: RDD[Store], customerRDD: RDD[Customer],
-    productRDD: RDD[Product], transactionRDD: RDD[Transaction]) {
+    productRDD: RDD[Product], transactionRDD: RDD[Transaction]): Unit = {
 
     locationRDD.saveAsObjectFile(outputDir + "/" + LOCATION_DIR)
     storeRDD.saveAsObjectFile(outputDir + "/" + STORE_DIR)
@@ -69,7 +61,7 @@ object IOUtils {
     transactionRDD.saveAsObjectFile(outputDir + "/" + TRANSACTION_DIR)
   }
 
-  def saveLocalAsJSON(outputDir: File, statistics: Statistics) {
+  def saveLocalAsJSON(outputDir: File, statistics: Statistics): Unit = {
     //load the write/read methods.
     implicit val formats = Serialization.formats(NoTypeHints)
     val json:String = write(statistics)
@@ -81,10 +73,10 @@ object IOUtils {
     implicit val formats = Serialization.formats(NoTypeHints)
     //Read file as String, and serialize it into Stats object.
     //See http://json4s.org/ examples.
-    
read[Statistics](scala.io.Source.fromFile(jsonFile).getLines.reduceLeft(_+_))
+    
read[Statistics](scala.io.Source.fromFile(jsonFile).getLines().reduceLeft(_+_))
   }
 
-  def saveLocalAsJSON(outputDir: File, recommendations:ProductRecommendations) 
{
+  def saveLocalAsJSON(outputDir: File, 
recommendations:ProductRecommendations): Unit = {
     //load the write/read methods.
     implicit val formats = Serialization.formats(NoTypeHints)
     val json:String = write(recommendations)
@@ -96,7 +88,7 @@ object IOUtils {
     implicit val formats = Serialization.formats(NoTypeHints)
     //Read file as String, and serialize it into Stats object.
     //See http://json4s.org/ examples.
-    
read[ProductRecommendations](scala.io.Source.fromFile(jsonFile).getLines.reduceLeft(_+_))
+    
read[ProductRecommendations](scala.io.Source.fromFile(jsonFile).getLines().reduceLeft(_+_))
   }
 
 
diff --git 
a/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/etl/ETL.scala
 
b/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/etl/ETL.scala
index 66d2f5247..3f046a66d 100644
--- 
a/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/etl/ETL.scala
+++ 
b/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/etl/ETL.scala
@@ -20,11 +20,8 @@ package org.apache.bigtop.bigpetstore.spark.etl
 import org.apache.bigtop.bigpetstore.spark.datamodel._
 
 import org.apache.spark.{SparkContext, SparkConf}
-import org.apache.spark.SparkContext._
 import org.apache.spark.rdd._
 
-import java.io.File
-import java.text.DateFormat
 import java.text.SimpleDateFormat
 import java.util._
 
@@ -37,7 +34,7 @@ object SparkETL {
 
   private val NPARAMS = 2
 
-  private def printUsage() {
+  private def printUsage(): Unit = {
     val usage: String = "BigPetStore Spark ETL\n" +
       "\n" +
       "Usage: spark-submit ... inputDir outputDir\n" +
@@ -92,7 +89,7 @@ object SparkETL {
       val txId = cols(10).toInt
       val df = new SimpleDateFormat("EEE MMM dd kk:mm:ss z yyyy", Locale.US)
       val txDate = df.parse(cols(11))
-      val txCal = Calendar.getInstance(Locale.US)
+      val txCal = 
Calendar.getInstance(TimeZone.getTimeZone("America/New_York"), Locale.US)
       txCal.setTime(txDate)
       txCal.set(Calendar.MILLISECOND, 0)
       val txProduct = cols(12)
@@ -196,15 +193,11 @@ object SparkETL {
     IOUtils.save(parameters.outputDir, locationRDD, storeRDD,
       customerRDD, productRDD, transactionRDD)
 
-    return (locationRDD.count(),
-        storeRDD.count(),
-        customerRDD.count(),
-        productRDD.count(),
-        transactionRDD.count()
-        );
+    (locationRDD.count(), storeRDD.count(), customerRDD.count(),
+      productRDD.count(), transactionRDD.count())
   }
 
-  def main(args: Array[String]) {
+  def main(args: Array[String]): Unit = {
     val parameters = parseArgs(args)
 
     println("Creating SparkConf")
diff --git 
a/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/generator/SparkDriver.scala
 
b/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/generator/SparkDriver.scala
index f27a15454..6169a0039 100644
--- 
a/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/generator/SparkDriver.scala
+++ 
b/bigtop-bigpetstore/bigpetstore-spark/src/main/scala/org/apache/bigpetstore/spark/generator/SparkDriver.scala
@@ -17,19 +17,17 @@
 
 package org.apache.bigtop.bigpetstore.spark.generator
 
-import com.github.rnowling.bps.datagenerator.datamodels.inputs.ZipcodeRecord
 import com.github.rnowling.bps.datagenerator.datamodels._
-import 
com.github.rnowling.bps.datagenerator.{DataLoader,StoreGenerator,CustomerGenerator
 => CustGen, PurchasingProfileGenerator,TransactionGenerator}
+import com.github.rnowling.bps.datagenerator.{DataLoader, 
PurchasingProfileGenerator, StoreGenerator, TransactionGenerator, 
CustomerGenerator => CustGen}
 import com.github.rnowling.bps.datagenerator.framework.SeedFactory
-import scala.collection.JavaConversions._
-import org.apache.spark.{SparkContext, SparkConf}
-import org.apache.spark.SparkContext._
+
+import scala.jdk.CollectionConverters._
+import org.apache.spark.{SparkConf, SparkContext}
 import org.apache.spark.rdd._
 
 import java.util.ArrayList
-import scala.util.Random
-import java.io.File
 import java.util.Date
+import scala.util.Random
 
 /**
  * This driver uses the data generator API to generate
@@ -50,7 +48,7 @@ object SparkDriver {
   private val NPARAMS = 5
   private val BURNIN_TIME = 7.0 // days
 
-  private def printUsage() {
+  private def printUsage(): Unit = {
     val usage: String =
       "BigPetStore Data Generator.\n" +
       "Usage: spark-submit ... outputDir nStores nCustomers simulationLength 
[seed]\n" +
@@ -62,7 +60,7 @@ object SparkDriver {
     System.err.println(usage)
   }
 
-  def parseArgs(args: Array[String]) {
+  def parseArgs(args: Array[String]): Unit = {
     if(args.length != NPARAMS && args.length != (NPARAMS - 1)) {
       printUsage()
       System.exit(1)
@@ -109,7 +107,7 @@ object SparkDriver {
       }
     }
     else {
-      seed = (new Random()).nextLong
+      seed = (new Random()).nextLong()
     }
   }
 
@@ -199,7 +197,7 @@ object SparkDriver {
   }
 
   def lineItem(t: Transaction, date:Date, p:Product): String = {
-      t.getStore.getId + "," +
+      t.getStore.getId.toString + "," +
       t.getStore.getLocation.getZipcode + "," +
       t.getStore.getLocation.getCity + "," +
       t.getStore.getLocation.getState + "," +
@@ -211,7 +209,7 @@ object SparkDriver {
       t.getId + "," +
       date + "," + p
   }
-  def writeData(transactionRDD : RDD[Transaction]) {
+  def writeData(transactionRDD : RDD[Transaction]): Unit = {
     val initialDate : Long = new Date().getTime()
 
     val transactionStringsRDD = transactionRDD.flatMap {
@@ -225,7 +223,7 @@ object SparkDriver {
         * So we ultimately define an RDD of strings, where each string 
represents
         * an instance where of a item purchase.
         * ********************************************************/
-        val records = products.map{
+        val records = products.asScala.map {
           product =>
             val storeLocation = transaction.getStore().getLocation()
             // days -> milliseconds = days * 24 h / day * 60 min / hr * 60 sec 
/ min * 1000 ms / sec
@@ -240,7 +238,7 @@ object SparkDriver {
     transactionStringsRDD.saveAsTextFile(outputDir + "/transactions")
   }
 
-  def main(args: Array[String]) {
+  def main(args: Array[String]): Unit = {
     parseArgs(args)
     val conf = new SparkConf().setAppName("BPS Data Generator")
     val sc = new SparkContext(conf)
diff --git 
a/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/TestFullPipeline.scala
 
b/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/TestFullPipeline.scala
index bc7e27f7f..92d5fe1ee 100644
--- 
a/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/TestFullPipeline.scala
+++ 
b/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/TestFullPipeline.scala
@@ -22,32 +22,25 @@ import 
org.apache.bigtop.bigpetstore.spark.analytics.RecommendProducts
 import org.apache.bigtop.bigpetstore.spark.datamodel.{Statistics, IOUtils}
 import org.apache.bigtop.bigpetstore.spark.etl.ETLParameters
 import org.apache.bigtop.bigpetstore.spark.etl.SparkETL
-import org.apache.bigtop.bigpetstore.spark.etl.{ETLParameters, SparkETL}
 import org.apache.bigtop.bigpetstore.spark.generator.SparkDriver
 import org.apache.spark.SparkConf
 import org.apache.spark.SparkContext
 import org.junit.runner.RunWith
-import org.scalatest.{BeforeAndAfterAll, FunSuite}
-import org.scalatest.junit.JUnitRunner
-
-import Array._
+import org.scalatest.BeforeAndAfterAll
+import org.scalatest.funsuite.AnyFunSuite
+import org.scalatestplus.junit.JUnitRunner
 
 import java.io.File
 import java.nio.file.Files
 
-import org.apache.spark.{SparkContext, SparkConf}
-import org.scalatest.junit.JUnitRunner
-import org.junit.runner.RunWith
-
-
 // hack for running tests with Gradle
 @RunWith(classOf[JUnitRunner])
-class TestFullPipeline extends FunSuite with BeforeAndAfterAll {
+class TestFullPipeline extends AnyFunSuite with BeforeAndAfterAll {
 
   val conf = new SparkConf().setAppName("BPS Data Generator Test 
Suite").setMaster("local[2]")
   val sc = new SparkContext(conf)
 
-  override def afterAll() {
+  override def afterAll(): Unit = {
     sc.stop()
   }
 
@@ -97,7 +90,6 @@ class TestFullPipeline extends FunSuite with 
BeforeAndAfterAll {
       recommJson.getAbsolutePath,
       sc, nIterations=5)
 
-
     sc.stop()
   }
 }
diff --git 
a/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/analytics/AnalyticsSuite.scala
 
b/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/analytics/AnalyticsSuite.scala
index d9ed39076..66e245ea0 100644
--- 
a/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/analytics/AnalyticsSuite.scala
+++ 
b/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/analytics/AnalyticsSuite.scala
@@ -17,24 +17,16 @@ package org.apache.bigpetstore.spark.analytics
  *  limitations under the License.
  */
 
-import com.google.common.collect.ImmutableMap
 import org.apache.bigtop.bigpetstore.spark.analytics.PetStoreStatistics
 import org.apache.bigtop.bigpetstore.spark.datamodel.Product
 import org.junit.runner.RunWith
-import org.scalatest.{BeforeAndAfterAll, FunSuite}
-import org.scalatest.junit.JUnitRunner
-
-import Array._
-
-import java.io.File
-import java.nio.file.Files
-import java.util.Calendar
-import java.util.Locale
-
+import org.scalatest.BeforeAndAfterAll
+import org.scalatest.funsuite.AnyFunSuite
+import org.scalatestplus.junit.JUnitRunner
 
 // hack for running tests with Gradle
 @RunWith(classOf[JUnitRunner])
-class AnalyticsSuite extends FunSuite with BeforeAndAfterAll {
+class AnalyticsSuite extends AnyFunSuite with BeforeAndAfterAll {
 
   test("product mapper") {
     val p = Product(1L, "cat1", Map(("a","a1"), ("b","b1")))
diff --git 
a/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/datamodel/IOUtilsSuite.scala
 
b/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/datamodel/IOUtilsSuite.scala
index 49ff1a4e6..126ef16c8 100644
--- 
a/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/datamodel/IOUtilsSuite.scala
+++ 
b/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/datamodel/IOUtilsSuite.scala
@@ -17,29 +17,25 @@
 
 package org.apache.bigtop.bigpetstore.spark.datamodel
 
-import Array._
-
-import java.io.File
 import java.nio.file.Files
 import java.util.Calendar
 import java.util.Locale
 
 import org.apache.spark.{SparkContext, SparkConf}
 
-import org.scalatest.{BeforeAndAfterAll, FunSuite}
-import org.scalatest.junit.JUnitRunner
+import org.scalatest.BeforeAndAfterAll
+import org.scalatest.funsuite.AnyFunSuite
+import org.scalatestplus.junit.JUnitRunner
 import org.junit.runner.RunWith
 
-import org.apache.bigtop.bigpetstore.spark.datamodel._
-
 // hack for running tests with Gradle
 @RunWith(classOf[JUnitRunner])
-class IOUtilsSuite extends FunSuite with BeforeAndAfterAll {
+class IOUtilsSuite extends AnyFunSuite with BeforeAndAfterAll {
 
   val conf = new SparkConf().setAppName("BPS Data Generator Test 
Suite").setMaster("local[2]")
   val sc = new SparkContext(conf)
 
-  override def afterAll() {
+  override def afterAll(): Unit = {
     sc.stop()
   }
 
@@ -92,6 +88,5 @@ class IOUtilsSuite extends FunSuite with BeforeAndAfterAll {
     assert(customerRDD.collect().toSet === readCustomerRDD.collect().toSet)
     assert(productRDD.collect().toSet === readProductRDD.collect().toSet)
     assert(transactionRDD.collect().toSet === 
readTransactionRDD.collect().toSet)
-
   }
 }
diff --git 
a/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/etl/ETLSuite.scala
 
b/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/etl/ETLSuite.scala
index 510e5f33b..70a3b354f 100644
--- 
a/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/etl/ETLSuite.scala
+++ 
b/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/etl/ETLSuite.scala
@@ -17,17 +17,14 @@
 
 package org.apache.bigtop.bigpetstore.spark.etl
 
-import org.apache.spark.rdd.RDD
-
-import Array._
-
 import java.util.Calendar
 import java.util.Locale
 import java.util.TimeZone
 
 import org.apache.spark.{SparkContext, SparkConf}
-import org.scalatest._
-import org.scalatest.junit.JUnitRunner
+import org.scalatest.BeforeAndAfterAll
+import org.scalatest.funsuite.AnyFunSuite
+import org.scalatestplus.junit.JUnitRunner
 import org.junit.runner.RunWith
 
 import org.apache.bigtop.bigpetstore.spark.datamodel._
@@ -44,7 +41,7 @@ import org.apache.bigtop.bigpetstore.spark.datamodel._
  * RunWith annotation is just a hack for running tests with Gradle
  */
 @RunWith(classOf[JUnitRunner])
-class ETLSuite extends FunSuite with BeforeAndAfterAll {
+class ETLSuite extends AnyFunSuite with BeforeAndAfterAll {
 
   /**
    * TODO : We are using Option monads as a replacement for nulls.
@@ -75,7 +72,7 @@ class ETLSuite extends FunSuite with BeforeAndAfterAll {
     "1,98110,Bainbridge 
Islan,WA,999,Cesareo,Lamplough,20152,Chantilly,VA,31,Mon Nov 02 17:51:37 EST 
2015,category=poop bags;brand=Dog 
Days;color=Blue;size=60.0;per_unit_cost=0.21;",
     "6,66067,Ottawa,KS,999,Cesareo,Lamplough,20152,Chantilly,VA,30,Mon Oct 12 
04:29:46 EDT 2015,category=dry cat food;brand=Feisty Feline;flavor=Chicken & 
Rice;size=14.0;per_unit_cost=2.14;")
 
-  override def beforeAll() {
+  override def beforeAll(): Unit = {
 
     val cal1 = Calendar.getInstance(TimeZone.getTimeZone("America/New_York"), 
Locale.US)
     val cal2 = Calendar.getInstance(TimeZone.getTimeZone("America/New_York"), 
Locale.US)
@@ -104,7 +101,7 @@ class ETLSuite extends FunSuite with BeforeAndAfterAll {
   }
 
 
-  override def afterAll() {
+  override def afterAll(): Unit = {
     sc.stop()
   }
 
@@ -113,7 +110,7 @@ class ETLSuite extends FunSuite with BeforeAndAfterAll {
     val expectedRecords = rawRecords.get
 
     //Goal: Confirm that these RDD's are identical to the expected ones.
-    val rdd = SparkETL.parseRawData(rawRDD).collect
+    val rdd = SparkETL.parseRawData(rawRDD).collect()
 
     /**
      * Assumption: Order of RDD elements will be same as the mock records.
@@ -136,12 +133,12 @@ class ETLSuite extends FunSuite with BeforeAndAfterAll {
       assert(rawRecord._5.storeId === expectedRecord._5.storeId)
 
       //BIGTOP-1586 : We want granular assertions, and we don't care to 
compare millisecond timestamps.
-      assert(rawRecord._5.dateTime.getTime.getYear === 
expectedRecord._5.dateTime.getTime.getYear)
-      assert(rawRecord._5.dateTime.getTime.getMonth === 
expectedRecord._5.dateTime.getTime.getMonth)
-      assert(rawRecord._5.dateTime.getTime.getDay === 
expectedRecord._5.dateTime.getTime.getDay)
-      assert(rawRecord._5.dateTime.getTime.getHours === 
expectedRecord._5.dateTime.getTime.getHours)
-      assert(rawRecord._5.dateTime.getTime.getMinutes === 
expectedRecord._5.dateTime.getTime.getMinutes)
-      assert(rawRecord._5.dateTime.getTime.getSeconds=== 
expectedRecord._5.dateTime.getTime.getSeconds)
+      assert(rawRecord._5.dateTime.get(Calendar.YEAR) === 
expectedRecord._5.dateTime.get(Calendar.YEAR))
+      assert(rawRecord._5.dateTime.get(Calendar.MONTH) === 
expectedRecord._5.dateTime.get(Calendar.MONTH))
+      assert(rawRecord._5.dateTime.get(Calendar.DAY_OF_MONTH) === 
expectedRecord._5.dateTime.get(Calendar.DAY_OF_MONTH))
+      assert(rawRecord._5.dateTime.get(Calendar.HOUR_OF_DAY) === 
expectedRecord._5.dateTime.get(Calendar.HOUR_OF_DAY))
+      assert(rawRecord._5.dateTime.get(Calendar.MINUTE) === 
expectedRecord._5.dateTime.get(Calendar.MINUTE))
+      assert(rawRecord._5.dateTime.get(Calendar.SECOND) === 
expectedRecord._5.dateTime.get(Calendar.SECOND))
     }
 
   }
diff --git 
a/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/generator/SparkDriverSuite.scala
 
b/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/generator/SparkDriverSuite.scala
index a3c6e1c02..6ab00c17f 100644
--- 
a/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/generator/SparkDriverSuite.scala
+++ 
b/bigtop-bigpetstore/bigpetstore-spark/src/test/scala/org/apache/bigpetstore/spark/generator/SparkDriverSuite.scala
@@ -17,28 +17,25 @@
 
 package org.apache.bigtop.bigpetstore.spark.generator
 
-import org.apache.bigtop.bigpetstore.spark.etl.{ETLParameters, SparkETL}
-
-import Array._
-
 import java.io.File
 import java.nio.file.Files
 
 import org.apache.spark.{SparkContext, SparkConf}
 
-import org.scalatest.{BeforeAndAfterAll, FunSuite}
-import org.scalatest.junit.JUnitRunner
+import org.scalatest.BeforeAndAfterAll
+import org.scalatest.funsuite.AnyFunSuite
+import org.scalatestplus.junit.JUnitRunner
 import org.junit.runner.RunWith
 
 
 // hack for running tests with Gradle
 @RunWith(classOf[JUnitRunner])
-class SparkDriverSuite extends FunSuite  with BeforeAndAfterAll {
+class SparkDriverSuite extends AnyFunSuite  with BeforeAndAfterAll {
 
   val conf = new SparkConf().setAppName("BPS Data Generator Test 
Suite").setMaster("local[2]")
   val sc = new SparkContext(conf);
 
-  override def afterAll() {
+  override def afterAll(): Unit = {
       sc.stop();
   }
 


Reply via email to