From 601c217cc402dc834da8091192525712e424fd33 Mon Sep 17 00:00:00 2001 From: Daniel Kavan Date: Thu, 3 Dec 2020 14:29:49 +0100 Subject: [PATCH 1/6] #48 no storer write fix - hadoopfs default storer is used - hdfs test enabled for build, while s3 ignored - readme update --- README.md | 28 ++++++++++++++++--- .../core/SparkQueryExecutionListener.scala | 4 ++- examples/pom.xml | 2 +- ...mpleMeasurementsS3RunnerExampleSpec.scala} | 4 ++- 4 files changed, 31 insertions(+), 7 deletions(-) rename examples/src/test/scala/za/co/absa/atum/examples/{SampleMeasurementsS3RunnerSpec.scala => SampleMeasurementsS3RunnerExampleSpec.scala} (89%) diff --git a/README.md b/README.md index 779dae08..82369369 100644 --- a/README.md +++ b/README.md @@ -160,7 +160,7 @@ object ExampleSparkJob { import spark.implicits._ // implicit FS is needed for enableControlMeasuresTracking, setCheckpoint calls, e.g. standard HDFS here: - implicit val localHdfs = FileSystem.get(new Configuration) + implicit val localHdfs = FileSystem.get(spark.sparkContext.hadoopConfiguration) // Initializing library to hook up to Apache Spark spark.enableControlMeasuresTracking(sourceInfoFile = "data/input/_INFO") @@ -188,8 +188,28 @@ in 'data/input/_INFO'. Two checkpoints are created. Any business logic can be in and saving it to Parquet format. ### Storing Measurements in AWS S3 -Starting with version 3.0.0, persistence support for AWS S3 has been added. -AWS S3 can be both used for loading the measurement data from as well as saving the measurements back to. + +#### AWS S3 via Hadoop FS API +Since version 3.1.0, persistence support for AWS S3 via Hadoop FS API is available. The usage is the same as with +regular HDFS with the exception of providing a different file system, e.g.: +```scala +import java.net.URI +import org.apache.hadoop.fs.FileSystem +import org.apache.spark.sql.SparkSession + +val spark = SparkSession + .builder() + .appName("Example Spark Job") + .getOrCreate() + +val s3Uri = new URI("s3://my-awesome-bucket") +implicit val fs = FileSystem.get(s3Uri, spark.sparkContext.hadoopConfiguration) + +``` +The rest of the usage is the same in the example listed above. + +#### AWS S3 via AWS SDK for S3 +Starting with version 3.0.0, there is also persistence support for AWS S3 via AWS SDK S3. The following example demonstrates the setup: ```scala @@ -230,7 +250,7 @@ object S3Example { } ``` -The rest of the processing logic and programatic approach to the library remains unchanged. +The rest of the processing logic and programmatic approach to the library remains unchanged. ## Atum library routines diff --git a/atum/src/main/scala/za/co/absa/atum/core/SparkQueryExecutionListener.scala b/atum/src/main/scala/za/co/absa/atum/core/SparkQueryExecutionListener.scala index 84bf0c6e..4a0eeb5b 100644 --- a/atum/src/main/scala/za/co/absa/atum/core/SparkQueryExecutionListener.scala +++ b/atum/src/main/scala/za/co/absa/atum/core/SparkQueryExecutionListener.scala @@ -44,7 +44,9 @@ class SparkQueryExecutionListener(cf: ControlFrameworkState) extends QueryExecut writeInfoFileForQuery(qe)(hadoopStorer.outputFs) case _ => - Atum.log.info("No usable storer is set, therefore no data will be written the automatically with DF-save to an _INFO file.") + Atum.log.info("No storer is set, using default HadoopFs-based bound with DF-save to an inferred _INFO file path.") + val defaultFs = FileSystem.get(qe.sparkSession.sparkContext.hadoopConfiguration) + writeInfoFileForQuery(qe)(defaultFs) } // Notify listeners diff --git a/examples/pom.xml b/examples/pom.xml index 79010e0e..966c961d 100644 --- a/examples/pom.xml +++ b/examples/pom.xml @@ -75,7 +75,7 @@ scalatest-maven-plugin ${scalatest.maven.version} - true + false diff --git a/examples/src/test/scala/za/co/absa/atum/examples/SampleMeasurementsS3RunnerSpec.scala b/examples/src/test/scala/za/co/absa/atum/examples/SampleMeasurementsS3RunnerExampleSpec.scala similarity index 89% rename from examples/src/test/scala/za/co/absa/atum/examples/SampleMeasurementsS3RunnerSpec.scala rename to examples/src/test/scala/za/co/absa/atum/examples/SampleMeasurementsS3RunnerExampleSpec.scala index 724a861c..2e3ec071 100644 --- a/examples/src/test/scala/za/co/absa/atum/examples/SampleMeasurementsS3RunnerSpec.scala +++ b/examples/src/test/scala/za/co/absa/atum/examples/SampleMeasurementsS3RunnerExampleSpec.scala @@ -15,10 +15,12 @@ package za.co.absa.atum.examples +import org.scalatest.Ignore import org.scalatest.funsuite.AnyFunSuite import za.co.absa.atum.utils._ -class SampleMeasurementsS3RunnerSpec extends AnyFunSuite +@Ignore +class SampleMeasurementsS3RunnerExampleSpec extends AnyFunSuite with SparkJobRunnerMethods with SparkLocalMaster { From 5c8057d812214c9c7b39149ca25c49a1b71ac911 Mon Sep 17 00:00:00 2001 From: Daniel Kavan Date: Thu, 3 Dec 2020 15:27:01 +0100 Subject: [PATCH 2/6] #48 projected ability for s3-over-hadoopFs pending measurements saving TODO test/testadd if it works with s3-over-hadoopFs like this --- .../core/SparkQueryExecutionListener.scala | 32 ++++++++++--------- .../absa/atum/utils/ExecutionPlanUtils.scala | 10 ++++-- .../za/co/absa/atum/utils/InfoFile.scala | 2 +- 3 files changed, 26 insertions(+), 18 deletions(-) diff --git a/atum/src/main/scala/za/co/absa/atum/core/SparkQueryExecutionListener.scala b/atum/src/main/scala/za/co/absa/atum/core/SparkQueryExecutionListener.scala index 4a0eeb5b..9bd3d7e3 100644 --- a/atum/src/main/scala/za/co/absa/atum/core/SparkQueryExecutionListener.scala +++ b/atum/src/main/scala/za/co/absa/atum/core/SparkQueryExecutionListener.scala @@ -17,14 +17,14 @@ package za.co.absa.atum.core import java.io.{PrintWriter, StringWriter} -import org.apache.hadoop.fs.FileSystem +import org.apache.hadoop.fs.Path import org.apache.spark.sql.execution.QueryExecution import org.apache.spark.sql.util.QueryExecutionListener import software.amazon.awssdk.auth.credentials.AwsCredentialsProvider import software.amazon.awssdk.regions.Region -import za.co.absa.atum.persistence.{HadoopFsControlMeasuresStorer, S3ControlMeasuresStorer, S3KmsSettings} +import za.co.absa.atum.persistence.{S3ControlMeasuresStorer, S3KmsSettings} import za.co.absa.atum.utils.ExecutionPlanUtils._ -import za.co.absa.atum.utils.S3Utils +import za.co.absa.atum.utils.{InfoFile, S3Utils} /** * The class is responsible for listening to DataSet save events and outputting corresponding control measurements. @@ -39,14 +39,9 @@ class SparkQueryExecutionListener(cf: ControlFrameworkState) extends QueryExecut Atum.log.debug(s"SparkQueryExecutionListener.onSuccess for S3ControlMeasuresStorer: writing to ${s3storer.outputLocation.s3String}") writeInfoFileForQueryForSdkS3(qe, s3storer.outputLocation.region, s3storer.kmsSettings)(s3storer.credentialsProvider) - case Some(hadoopStorer: HadoopFsControlMeasuresStorer) => - Atum.log.debug(s"SparkQueryExecutionListener.onSuccess: writing to Hadoop FS") - writeInfoFileForQuery(qe)(hadoopStorer.outputFs) - case _ => - Atum.log.info("No storer is set, using default HadoopFs-based bound with DF-save to an inferred _INFO file path.") - val defaultFs = FileSystem.get(qe.sparkSession.sparkContext.hadoopConfiguration) - writeInfoFileForQuery(qe)(defaultFs) + Atum.log.debug(s"SparkQueryExecutionListener.onSuccess: writing to Hadoop FS") + writeInfoFileForQuery(qe) } // Notify listeners @@ -66,14 +61,21 @@ class SparkQueryExecutionListener(cf: ControlFrameworkState) extends QueryExecut } /** Write _INFO file with control measurements to the output directory based on the query plan */ - private def writeInfoFileForQuery(qe: QueryExecution)(implicit outputFs: FileSystem): Unit = { - val infoFilePath = inferOutputInfoFileName(qe, cf.outputInfoFileName) + private[core] def writeInfoFileForQuery(qe: QueryExecution)(): Unit = { + val infoFileDir: Option[String] = inferOutputInfoFileDir(qe) + + implicit val hadoopConf = qe.sparkSession.sparkContext.hadoopConfiguration + val fsWithDir = infoFileDir + .map(InfoFile) + .flatMap(_.toOptFsPath) // path + FS based on HDFS or S3 over hadoopFS // Write _INFO file to the output directory - infoFilePath.foreach(path => { + fsWithDir.foreach { case (fs, dir) => { + val path = new Path(dir, cf.outputInfoFileName) + Atum.log.info(s"Inferred _INFO Path = ${path.toUri.toString}") - cf.storeCurrentInfoFile(path) - }) + cf.storeCurrentInfoFile(path)(fs) + }} // Write _INFO file to a registered storer if (cf.accumulator.isStorerLoaded) { diff --git a/atum/src/main/scala/za/co/absa/atum/utils/ExecutionPlanUtils.scala b/atum/src/main/scala/za/co/absa/atum/utils/ExecutionPlanUtils.scala index adf46e2c..a4d95b38 100644 --- a/atum/src/main/scala/za/co/absa/atum/utils/ExecutionPlanUtils.scala +++ b/atum/src/main/scala/za/co/absa/atum/utils/ExecutionPlanUtils.scala @@ -97,11 +97,17 @@ object ExecutionPlanUtils { * @return The inferred output control measurements file path of the source dataset */ def inferOutputInfoFileName(qe: QueryExecution, infoFileName: String = Constants.DefaultInfoFileName): Option[Path] = { + inferOutputInfoFileDir(qe).map { dir => + new Path(dir, infoFileName) + } + } + + private[atum] def inferOutputInfoFileDir(qe: QueryExecution): Option[String] = { qe.analyzed match { case s: SaveIntoDataSourceCommand => - Some(new Path(s.options("path"), infoFileName)) + Some(s.options("path")) case h: InsertIntoHadoopFsRelationCommand => - Some(new Path(h.outputPath, infoFileName)) + Some(h.outputPath.toString) case a => log.warn(s"Logical plan: ${qe.logical.treeString}") log.warn(s"Analyzed plan: ${qe.analyzed.treeString}") diff --git a/atum/src/main/scala/za/co/absa/atum/utils/InfoFile.scala b/atum/src/main/scala/za/co/absa/atum/utils/InfoFile.scala index ad80550d..5f2a9bf9 100644 --- a/atum/src/main/scala/za/co/absa/atum/utils/InfoFile.scala +++ b/atum/src/main/scala/za/co/absa/atum/utils/InfoFile.scala @@ -12,7 +12,7 @@ private[atum] case class InfoFile(infoFile: String) { private val validatedInfoFile: Option[String] = if (infoFile.isEmpty) None else Some(infoFile) - private def toOptFsPath(implicit hadoopConfiguration: Configuration): Option[(FileSystem, Path)] = { + def toOptFsPath(implicit hadoopConfiguration: Configuration): Option[(FileSystem, Path)] = { validatedInfoFile.map { definedInfoFile => definedInfoFile.toS3Location match { From 6872dfad1eee77f03c45597a36df5f151bae1806 Mon Sep 17 00:00:00 2001 From: Daniel Kavan Date: Thu, 3 Dec 2020 16:49:41 +0100 Subject: [PATCH 3/6] #48 implicit saving test adding 1 (loader non-"", storer = "") --- .../absa/atum/HdfsInfoIntegrationSuite.scala | 65 +++++++++++++++++++ .../za/co/absa/atum/LocalFsTestUtils.scala | 43 ++++++++++++ 2 files changed, 108 insertions(+) create mode 100644 examples/src/test/scala/za/co/absa/atum/HdfsInfoIntegrationSuite.scala create mode 100644 examples/src/test/scala/za/co/absa/atum/LocalFsTestUtils.scala diff --git a/examples/src/test/scala/za/co/absa/atum/HdfsInfoIntegrationSuite.scala b/examples/src/test/scala/za/co/absa/atum/HdfsInfoIntegrationSuite.scala new file mode 100644 index 00000000..6ba4c8f3 --- /dev/null +++ b/examples/src/test/scala/za/co/absa/atum/HdfsInfoIntegrationSuite.scala @@ -0,0 +1,65 @@ +package za.co.absa.atum + +import org.apache.hadoop.fs.FileSystem +import org.apache.log4j.LogManager +import org.apache.spark.sql.{DataFrame, SaveMode} +import org.scalatest.flatspec.AnyFlatSpec +import org.scalatest.matchers.should.Matchers +import za.co.absa.atum.model.{Checkpoint, Measurement} +import za.co.absa.atum.persistence.ControlMeasuresParser +import za.co.absa.atum.utils.SparkTestBase + +class HdfsInfoIntegrationSuite extends AnyFlatSpec with SparkTestBase with Matchers { + + private val log = LogManager.getLogger(this.getClass) + + + private val inputCsv = "data/input/wikidata.csv" + private def readSparkInputCsv(inputCsvPath: String): DataFrame = spark.read + .option("header", "true") + .option("inferSchema", "true") + .csv(inputCsvPath) + + private def writeSparkData(df: DataFrame, outputPath: String): Unit = + df.write.mode(SaveMode.Overwrite) + .parquet(outputPath) + + + "_INFO" should "be written implicitly on spark.write" in { + val tempDir = LocalFsTestUtils.createLocalTemporaryDirectory("hdfsTestOutput") // todo beforeAll+cleanup afterAll? + + import spark.implicits._ + import za.co.absa.atum.AtumImplicits._ + + val hadoopConfiguration = spark.sparkContext.hadoopConfiguration + implicit val fs = FileSystem.get(hadoopConfiguration) + + // Initializing library to hook up to Apache Spark + spark.enableControlMeasuresTracking(sourceInfoFile = "data/input/wikidata.csv.info") // todo version with None, None, too? + .setControlMeasuresWorkflow("Job 1") + + val df1 = readSparkInputCsv(inputCsv) + df1.setCheckpoint("Checkpoint0") + val filteredDf1 = df1.filter($"total_response_size" > 1000) + filteredDf1.setCheckpoint("Checkpoint1") // stateful, do not need return value + + val outputPath = s"$tempDir/hdfsOutput/implicitTest1" + writeSparkData(filteredDf1, outputPath) + + spark.disableControlMeasuresTracking() + + log.info(s"Checking $outputPath/_INFO to contain expected values") + + val infoContentJson = LocalFsTestUtils.readFileAsString(s"$outputPath/_INFO") + val infoControlMeasures = ControlMeasuresParser.fromJson(infoContentJson) + + infoControlMeasures.checkpoints.map(_.name) shouldBe Seq("Source", "Raw", "Checkpoint0", "Checkpoint1") + val checkpoint0 = infoControlMeasures.checkpoints.collectFirst{ case c: Checkpoint if c.name == "Checkpoint0" => c }.get // todo generalize + checkpoint0.controls should contain (Measurement("recordCount", "count", "*", "5000")) + + val checkpoint1 = infoControlMeasures.checkpoints.collectFirst{ case c: Checkpoint if c.name == "Checkpoint1" => c }.get + checkpoint1.controls should contain (Measurement("recordCount", "count", "*", "4964")) + + LocalFsTestUtils.safeDeleteTestDir(tempDir) + } +} diff --git a/examples/src/test/scala/za/co/absa/atum/LocalFsTestUtils.scala b/examples/src/test/scala/za/co/absa/atum/LocalFsTestUtils.scala new file mode 100644 index 00000000..d2fc9f66 --- /dev/null +++ b/examples/src/test/scala/za/co/absa/atum/LocalFsTestUtils.scala @@ -0,0 +1,43 @@ +package za.co.absa.atum + +import java.io.File +import java.nio.file.Files + +import org.apache.commons.io.FileUtils +import org.apache.log4j.LogManager + +import scala.io.Source +import scala.util.control.NonFatal + +object LocalFsTestUtils { + private val log = LogManager.getLogger(this.getClass) + + /** + * Creates a temporary directory in the local filesystem. + * + * @param prefix A prefix to use for the temporary directory. + * @return A path to a temporary directory. + */ + def createLocalTemporaryDirectory(prefix: String): String = { + val tmpPath = Files.createTempDirectory(prefix) + tmpPath.toAbsolutePath.toString + } + + def safeDeleteTestDir(path: String): Unit = { + try { + FileUtils.deleteDirectory(new File(path)) + } catch { + case NonFatal(e) => log.warn(s"Unable to delete a test directory $path") + } + } + + def readFileAsString(filename: String, lineSeparator: String = "\n"): String = { + val sourceFile = Source.fromFile(filename) + try { + sourceFile.getLines().mkString(lineSeparator) + } finally { + sourceFile.close() + } + } + +} From c140413ea317e720ec79b91ba43eeaf45fa0e468 Mon Sep 17 00:00:00 2001 From: Daniel Kavan Date: Thu, 10 Dec 2020 11:13:05 +0100 Subject: [PATCH 4/6] #48 explicit saving test adding (loader non-"", storer = defined) --- .../absa/atum/HdfsInfoIntegrationSuite.scala | 69 +++++++++++-------- 1 file changed, 40 insertions(+), 29 deletions(-) diff --git a/examples/src/test/scala/za/co/absa/atum/HdfsInfoIntegrationSuite.scala b/examples/src/test/scala/za/co/absa/atum/HdfsInfoIntegrationSuite.scala index 6ba4c8f3..b98f5a8a 100644 --- a/examples/src/test/scala/za/co/absa/atum/HdfsInfoIntegrationSuite.scala +++ b/examples/src/test/scala/za/co/absa/atum/HdfsInfoIntegrationSuite.scala @@ -3,16 +3,21 @@ package za.co.absa.atum import org.apache.hadoop.fs.FileSystem import org.apache.log4j.LogManager import org.apache.spark.sql.{DataFrame, SaveMode} +import org.scalatest.BeforeAndAfterAll import org.scalatest.flatspec.AnyFlatSpec import org.scalatest.matchers.should.Matchers import za.co.absa.atum.model.{Checkpoint, Measurement} import za.co.absa.atum.persistence.ControlMeasuresParser import za.co.absa.atum.utils.SparkTestBase -class HdfsInfoIntegrationSuite extends AnyFlatSpec with SparkTestBase with Matchers { +class HdfsInfoIntegrationSuite extends AnyFlatSpec with SparkTestBase with Matchers with BeforeAndAfterAll { private val log = LogManager.getLogger(this.getClass) + val tempDir: String = LocalFsTestUtils.createLocalTemporaryDirectory("hdfsTestOutput") + override def afterAll = { + LocalFsTestUtils.safeDeleteTestDir(tempDir) + } private val inputCsv = "data/input/wikidata.csv" private def readSparkInputCsv(inputCsvPath: String): DataFrame = spark.read @@ -24,42 +29,48 @@ class HdfsInfoIntegrationSuite extends AnyFlatSpec with SparkTestBase with Match df.write.mode(SaveMode.Overwrite) .parquet(outputPath) + { + val outputPath = s"$tempDir/outputCheck1" + // implicit variant only writes to derived outputPath, explicit writes to both implicit derived path and the explicit one, too. + Seq( + ("implicit output _INFO path only", "", Seq(s"$outputPath/_INFO")), + ("implicit & explicit output _INFO path", s"$outputPath/extra/_INFO2", Seq(s"$outputPath/_INFO", s"$outputPath/extra/_INFO2")) + ).foreach { case (testCaseName, destinationInfoFilePath, expectedPaths) => - "_INFO" should "be written implicitly on spark.write" in { - val tempDir = LocalFsTestUtils.createLocalTemporaryDirectory("hdfsTestOutput") // todo beforeAll+cleanup afterAll? - - import spark.implicits._ - import za.co.absa.atum.AtumImplicits._ - - val hadoopConfiguration = spark.sparkContext.hadoopConfiguration - implicit val fs = FileSystem.get(hadoopConfiguration) + "_INFO" should s"be written on spark.write ($testCaseName)" in { + import spark.implicits._ + import za.co.absa.atum.AtumImplicits._ - // Initializing library to hook up to Apache Spark - spark.enableControlMeasuresTracking(sourceInfoFile = "data/input/wikidata.csv.info") // todo version with None, None, too? - .setControlMeasuresWorkflow("Job 1") + val hadoopConfiguration = spark.sparkContext.hadoopConfiguration + implicit val fs = FileSystem.get(hadoopConfiguration) - val df1 = readSparkInputCsv(inputCsv) - df1.setCheckpoint("Checkpoint0") - val filteredDf1 = df1.filter($"total_response_size" > 1000) - filteredDf1.setCheckpoint("Checkpoint1") // stateful, do not need return value + // Initializing library to hook up to Apache Spark + spark.enableControlMeasuresTracking(sourceInfoFile = "data/input/wikidata.csv.info", destinationInfoFile = destinationInfoFilePath) + .setControlMeasuresWorkflow("Job 1") - val outputPath = s"$tempDir/hdfsOutput/implicitTest1" - writeSparkData(filteredDf1, outputPath) + val df1 = readSparkInputCsv(inputCsv) + df1.setCheckpoint("Checkpoint0") + val filteredDf1 = df1.filter($"total_response_size" > 1000) + filteredDf1.setCheckpoint("Checkpoint1") // stateful, do not need return value + writeSparkData(filteredDf1, outputPath) // implicit output _INFO file path is derived from this path passed to spark.write - spark.disableControlMeasuresTracking() + spark.disableControlMeasuresTracking() - log.info(s"Checking $outputPath/_INFO to contain expected values") + expectedPaths.foreach { expectedPath => + log.info(s"Checking $expectedPath to contain expected values") - val infoContentJson = LocalFsTestUtils.readFileAsString(s"$outputPath/_INFO") - val infoControlMeasures = ControlMeasuresParser.fromJson(infoContentJson) + val infoContentJson = LocalFsTestUtils.readFileAsString(expectedPath) + val infoControlMeasures = ControlMeasuresParser.fromJson(infoContentJson) - infoControlMeasures.checkpoints.map(_.name) shouldBe Seq("Source", "Raw", "Checkpoint0", "Checkpoint1") - val checkpoint0 = infoControlMeasures.checkpoints.collectFirst{ case c: Checkpoint if c.name == "Checkpoint0" => c }.get // todo generalize - checkpoint0.controls should contain (Measurement("recordCount", "count", "*", "5000")) + infoControlMeasures.checkpoints.map(_.name) shouldBe Seq("Source", "Raw", "Checkpoint0", "Checkpoint1") + val checkpoint0 = infoControlMeasures.checkpoints.collectFirst { case c: Checkpoint if c.name == "Checkpoint0" => c }.get + checkpoint0.controls should contain(Measurement("recordCount", "count", "*", "5000")) - val checkpoint1 = infoControlMeasures.checkpoints.collectFirst{ case c: Checkpoint if c.name == "Checkpoint1" => c }.get - checkpoint1.controls should contain (Measurement("recordCount", "count", "*", "4964")) - - LocalFsTestUtils.safeDeleteTestDir(tempDir) + val checkpoint1 = infoControlMeasures.checkpoints.collectFirst { case c: Checkpoint if c.name == "Checkpoint1" => c }.get + checkpoint1.controls should contain(Measurement("recordCount", "count", "*", "4964")) + } + } + } } + } From a051dfa06103778ebc0b4a604ac6b12ebfb043e7 Mon Sep 17 00:00:00 2001 From: Daniel Kavan Date: Mon, 4 Jan 2021 11:24:00 +0100 Subject: [PATCH 5/6] #48 PR touchups (explicit types, comments, etc.) --- .../za/co/absa/atum/core/SparkQueryExecutionListener.scala | 3 ++- .../scala/za/co/absa/atum/utils/ExecutionPlanUtils.scala | 7 ++++++- .../hdfs/ControlMeasuresHdfsStorerJsonSpec.scala | 7 ++++--- .../za/co/absa/atum/examples/SampleMeasurements1.scala | 2 +- .../za/co/absa/atum/examples/SampleMeasurements2.scala | 2 +- .../co/absa/atum/examples/SampleSdkS3Measurements1.scala | 2 +- .../co/absa/atum/examples/SampleSdkS3Measurements2.scala | 2 +- .../scala/za/co/absa/atum/HdfsInfoIntegrationSuite.scala | 4 ++-- .../src/test/scala/za/co/absa/atum/LocalFsTestUtils.scala | 2 +- 9 files changed, 19 insertions(+), 12 deletions(-) diff --git a/atum/src/main/scala/za/co/absa/atum/core/SparkQueryExecutionListener.scala b/atum/src/main/scala/za/co/absa/atum/core/SparkQueryExecutionListener.scala index 9bd3d7e3..4f38d1ee 100644 --- a/atum/src/main/scala/za/co/absa/atum/core/SparkQueryExecutionListener.scala +++ b/atum/src/main/scala/za/co/absa/atum/core/SparkQueryExecutionListener.scala @@ -17,6 +17,7 @@ package za.co.absa.atum.core import java.io.{PrintWriter, StringWriter} +import org.apache.hadoop.conf.Configuration import org.apache.hadoop.fs.Path import org.apache.spark.sql.execution.QueryExecution import org.apache.spark.sql.util.QueryExecutionListener @@ -64,7 +65,7 @@ class SparkQueryExecutionListener(cf: ControlFrameworkState) extends QueryExecut private[core] def writeInfoFileForQuery(qe: QueryExecution)(): Unit = { val infoFileDir: Option[String] = inferOutputInfoFileDir(qe) - implicit val hadoopConf = qe.sparkSession.sparkContext.hadoopConfiguration + implicit val hadoopConf: Configuration = qe.sparkSession.sparkContext.hadoopConfiguration val fsWithDir = infoFileDir .map(InfoFile) .flatMap(_.toOptFsPath) // path + FS based on HDFS or S3 over hadoopFS diff --git a/atum/src/main/scala/za/co/absa/atum/utils/ExecutionPlanUtils.scala b/atum/src/main/scala/za/co/absa/atum/utils/ExecutionPlanUtils.scala index a4d95b38..ed1007c9 100644 --- a/atum/src/main/scala/za/co/absa/atum/utils/ExecutionPlanUtils.scala +++ b/atum/src/main/scala/za/co/absa/atum/utils/ExecutionPlanUtils.scala @@ -102,7 +102,12 @@ object ExecutionPlanUtils { } } - private[atum] def inferOutputInfoFileDir(qe: QueryExecution): Option[String] = { + /** + * Based on the `qe` supplied, output _INFO file path is inference is attempted + * @param qe QueryExecution - path inference basis + * @return optional inferred _INFO file path + */ + def inferOutputInfoFileDir(qe: QueryExecution): Option[String] = { qe.analyzed match { case s: SaveIntoDataSourceCommand => Some(s.options("path")) diff --git a/atum/src/test/scala/za/co/absa/atum/persistence/hdfs/ControlMeasuresHdfsStorerJsonSpec.scala b/atum/src/test/scala/za/co/absa/atum/persistence/hdfs/ControlMeasuresHdfsStorerJsonSpec.scala index 7c3b0579..9d29a790 100644 --- a/atum/src/test/scala/za/co/absa/atum/persistence/hdfs/ControlMeasuresHdfsStorerJsonSpec.scala +++ b/atum/src/test/scala/za/co/absa/atum/persistence/hdfs/ControlMeasuresHdfsStorerJsonSpec.scala @@ -4,16 +4,17 @@ import org.apache.hadoop.conf.Configuration import org.apache.hadoop.fs.{FileSystem, Path} import org.scalatest.flatspec.AnyFlatSpec import org.scalatest.matchers.should.Matchers +import za.co.absa.atum.model.ControlMeasure import za.co.absa.atum.persistence.TestResources import za.co.absa.atum.utils.{FileUtils, HdfsFileUtils} class ControlMeasuresHdfsStorerJsonSpec extends AnyFlatSpec with Matchers { val expectedFilePath: String = TestResources.InputInfo.localPath - val inputControlMeasure = TestResources.InputInfo.controlMeasure + val inputControlMeasure: ControlMeasure = TestResources.InputInfo.controlMeasure - val hadoopConfiguration = new Configuration() - implicit val fs = FileSystem.get(hadoopConfiguration) + val hadoopConfiguration: Configuration = new Configuration() + implicit val fs: FileSystem = FileSystem.get(hadoopConfiguration) "ControlMeasuresHdfsStorerJsonFile" should "store json file to HDFS" in { diff --git a/examples/src/main/scala/za/co/absa/atum/examples/SampleMeasurements1.scala b/examples/src/main/scala/za/co/absa/atum/examples/SampleMeasurements1.scala index b9e4b11d..f71fa4c2 100644 --- a/examples/src/main/scala/za/co/absa/atum/examples/SampleMeasurements1.scala +++ b/examples/src/main/scala/za/co/absa/atum/examples/SampleMeasurements1.scala @@ -29,7 +29,7 @@ object SampleMeasurements1 { import spark.implicits._ val hadoopConfiguration = spark.sparkContext.hadoopConfiguration - implicit val fs = FileSystem.get(hadoopConfiguration) + implicit val fs: FileSystem = FileSystem.get(hadoopConfiguration) // Initializing library to hook up to Apache Spark spark.enableControlMeasuresTracking(sourceInfoFile = "data/input/wikidata.csv.info") diff --git a/examples/src/main/scala/za/co/absa/atum/examples/SampleMeasurements2.scala b/examples/src/main/scala/za/co/absa/atum/examples/SampleMeasurements2.scala index a2e4238e..8ea0f7a1 100644 --- a/examples/src/main/scala/za/co/absa/atum/examples/SampleMeasurements2.scala +++ b/examples/src/main/scala/za/co/absa/atum/examples/SampleMeasurements2.scala @@ -30,7 +30,7 @@ object SampleMeasurements2 { import spark.implicits._ val hadoopConfiguration = spark.sparkContext.hadoopConfiguration - implicit val fs = FileSystem.get(hadoopConfiguration) + implicit val fs: FileSystem = FileSystem.get(hadoopConfiguration) // Initializing library to hook up to Apache Spark // No need to specify datasetName and datasetVersion as it is stage 2 and it will be determined automatically diff --git a/examples/src/main/scala/za/co/absa/atum/examples/SampleSdkS3Measurements1.scala b/examples/src/main/scala/za/co/absa/atum/examples/SampleSdkS3Measurements1.scala index 78cd44a3..1d15eab9 100644 --- a/examples/src/main/scala/za/co/absa/atum/examples/SampleSdkS3Measurements1.scala +++ b/examples/src/main/scala/za/co/absa/atum/examples/SampleSdkS3Measurements1.scala @@ -32,7 +32,7 @@ object SampleSdkS3Measurements1 { import spark.implicits._ val hadoopConfiguration = spark.sparkContext.hadoopConfiguration - implicit val fs = FileSystem.get(hadoopConfiguration) + implicit val fs: FileSystem = FileSystem.get(hadoopConfiguration) // This sample example relies on local credentials profile named "saml" with access to the s3 location defined below implicit val samlCredentialsProvider = S3Utils.getLocalProfileCredentialsProvider("saml") diff --git a/examples/src/main/scala/za/co/absa/atum/examples/SampleSdkS3Measurements2.scala b/examples/src/main/scala/za/co/absa/atum/examples/SampleSdkS3Measurements2.scala index 3dc619ef..27726f3f 100644 --- a/examples/src/main/scala/za/co/absa/atum/examples/SampleSdkS3Measurements2.scala +++ b/examples/src/main/scala/za/co/absa/atum/examples/SampleSdkS3Measurements2.scala @@ -34,7 +34,7 @@ object SampleSdkS3Measurements2 { import spark.implicits._ val hadoopConfiguration = spark.sparkContext.hadoopConfiguration - implicit val fs = FileSystem.get(hadoopConfiguration) + implicit val fs: FileSystem = FileSystem.get(hadoopConfiguration) // This sample example relies on local credentials profile named "saml" with access to the s3 location defined below // AND by having explicitly defined KMS Key ID diff --git a/examples/src/test/scala/za/co/absa/atum/HdfsInfoIntegrationSuite.scala b/examples/src/test/scala/za/co/absa/atum/HdfsInfoIntegrationSuite.scala index b98f5a8a..ad3aac5b 100644 --- a/examples/src/test/scala/za/co/absa/atum/HdfsInfoIntegrationSuite.scala +++ b/examples/src/test/scala/za/co/absa/atum/HdfsInfoIntegrationSuite.scala @@ -15,7 +15,7 @@ class HdfsInfoIntegrationSuite extends AnyFlatSpec with SparkTestBase with Match private val log = LogManager.getLogger(this.getClass) val tempDir: String = LocalFsTestUtils.createLocalTemporaryDirectory("hdfsTestOutput") - override def afterAll = { + override def afterAll: Unit = { LocalFsTestUtils.safeDeleteTestDir(tempDir) } @@ -42,7 +42,7 @@ class HdfsInfoIntegrationSuite extends AnyFlatSpec with SparkTestBase with Match import za.co.absa.atum.AtumImplicits._ val hadoopConfiguration = spark.sparkContext.hadoopConfiguration - implicit val fs = FileSystem.get(hadoopConfiguration) + implicit val fs: FileSystem = FileSystem.get(hadoopConfiguration) // Initializing library to hook up to Apache Spark spark.enableControlMeasuresTracking(sourceInfoFile = "data/input/wikidata.csv.info", destinationInfoFile = destinationInfoFilePath) diff --git a/examples/src/test/scala/za/co/absa/atum/LocalFsTestUtils.scala b/examples/src/test/scala/za/co/absa/atum/LocalFsTestUtils.scala index d2fc9f66..5f40f107 100644 --- a/examples/src/test/scala/za/co/absa/atum/LocalFsTestUtils.scala +++ b/examples/src/test/scala/za/co/absa/atum/LocalFsTestUtils.scala @@ -27,7 +27,7 @@ object LocalFsTestUtils { try { FileUtils.deleteDirectory(new File(path)) } catch { - case NonFatal(e) => log.warn(s"Unable to delete a test directory $path") + case NonFatal(_) => log.warn(s"Unable to delete a test directory $path") } } From 61adf96076eb9238a59612d5e6f8c5d78812d0a5 Mon Sep 17 00:00:00 2001 From: Daniel Kavan Date: Tue, 5 Jan 2021 09:44:40 +0100 Subject: [PATCH 6/6] #48 mergefix --- build-all.sh | 0 1 file changed, 0 insertions(+), 0 deletions(-) mode change 100755 => 100644 build-all.sh diff --git a/build-all.sh b/build-all.sh old mode 100755 new mode 100644