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();
}