diff --git a/.github/workflows/jacoco.yml b/.github/workflows/jacoco.yml
index aae932eeb..b19e4457d 100644
--- a/.github/workflows/jacoco.yml
+++ b/.github/workflows/jacoco.yml
@@ -28,8 +28,8 @@ concurrency:
cancel-in-progress: true
env:
- SCALA_VERSION: 2.12.20
- SPARK_VERSION: 3.4.4
+ SCALA_VERSION: 2.12.21
+ SPARK_VERSION: 3.5.5
# hint: "group thresholds" are in format: 'overall*changed-files-average*per-changed-file'
REPORT_GROUPS: |
@@ -89,7 +89,7 @@ jobs:
uses: actions/setup-java@de7274f081f381c8f8158605e0321c36c376e2e6 # v6.0.1
with:
distribution: temurin
- java-version: 8
+ java-version: 17
cache: sbt
- name: Build and run tests with coverage
diff --git a/.github/workflows/scala.yml b/.github/workflows/scala.yml
index 922e7cbca..2a10e5023 100644
--- a/.github/workflows/scala.yml
+++ b/.github/workflows/scala.yml
@@ -22,17 +22,13 @@ jobs:
strategy:
fail-fast: false
matrix:
- scala: [2.11.12, 2.12.20, 2.13.16]
- spark: [2.4.8, 3.4.4, 3.5.5]
+ scala: [2.12.21, 2.13.18]
+ spark: [3.5.5, 4.1.2]
exclude:
- - scala: 2.11.12
- spark: 3.4.4
- - scala: 2.11.12
+ - scala: 2.13.18
spark: 3.5.5
- - scala: 2.12.20
- spark: 2.4.8
- - scala: 2.13.16
- spark: 2.4.8
+ - scala: 2.12.21
+ spark: 4.1.2
name: Test Spark ${{matrix.spark}} on Scala ${{matrix.scala}}
steps:
- name: Checkout code
@@ -43,7 +39,7 @@ jobs:
uses: actions/setup-java@v4.2.1
with:
distribution: temurin
- java-version: 8
+ java-version: 17
cache: sbt
- name: Build and run unit tests
working-directory: ./pramen
diff --git a/README.md b/README.md
index abb0b7de7..10512acff 100644
--- a/README.md
+++ b/README.md
@@ -85,7 +85,7 @@ In addition to basic error notification, typical operational warnings are genera
```sh
git clone https://github.com/AbsaOSS/pramen
cd pramen
- sbt -DSPARK_VERSION="3.3.4" ++2.12.18 assembly
+ sbt -DSPARK_VERSION="3.5.5" ++2.12.21 assembly
```
(You need JDK 1.8 installed to run this)
@@ -203,15 +203,15 @@ Pramen for Python transformers is available in PyPi: [ section of the project.
-- Or by building Pramen from source and creating an uber JAR file that contains all dependencies required to run the pipeline on a Spark cluster (see below).
+- Or by building Pramen from source and creating an uber JAR file that contains all dependencies required to run the
+ pipeline on a Spark cluster (see below).
### Building a Pramen runner JAR from sources
Creating an uber jar for Pramen is very easy. Just clone the repository and run one of the following commands:
```sh
-sbt ++2.11.12 assembly
-sbt ++2.12.20 assembly
-sbt ++2.13.16 assembly
+sbt ++2.12.21 assembly
+sbt ++2.13.18 assembly
```
You can collect the uber jar of Pramen either at
@@ -222,14 +222,13 @@ Since `1.7.0` Pramen runner bundle does not include Delta Lake format classes si
Spark distributions. This makes the runner independent of Spark version. But if you want to include Delta Lake files
in your bundle, use one of example commands specifying your Spark version:
```sh
-sbt -DSPARK_VERSION="2.4.8" -Dassembly.features="includeDelta" ++2.11.12 assembly
-sbt -DSPARK_VERSION="3.3.4" -Dassembly.features="includeDelta" ++2.12.20 assembly
-sbt -DSPARK_VERSION="3.5.5" -Dassembly.features="includeDelta" ++2.13.16 assembly
+sbt -DSPARK_VERSION="3.5.5" -Dassembly.features="includeDelta" ++2.12.21 assembly
+sbt -DSPARK_VERSION="4.1.2" -Dassembly.features="includeDelta" ++2.13.18 assembly
```
Then, run `spark-shell` or `spark-submit` adding the fat jar as the option.
```sh
-$ spark-shell --jars pramen-runner_2.12-1.7.5-SNAPSHOT.jar
+$ spark-shell --jars pramen-runner_2.12-1.15.1-SNAPSHOT.jar
```
# Creating a data pipeline
diff --git a/pramen/api/pom.xml b/pramen/api/pom.xml
index 153d9835e..9ce683d61 100644
--- a/pramen/api/pom.xml
+++ b/pramen/api/pom.xml
@@ -27,7 +27,7 @@
za.co.absa.pramen
pramen
- 1.14.9-SNAPSHOT
+ 1.15.0-SNAPSHOT
diff --git a/pramen/build.sbt b/pramen/build.sbt
index a450516a7..dbfa33965 100644
--- a/pramen/build.sbt
+++ b/pramen/build.sbt
@@ -18,17 +18,16 @@ import Dependencies._
import Versions._
import BuildInfoTemplateSettings._
-val scala211 = "2.11.12"
-val scala212 = "2.12.20"
-val scala213 = "2.13.16"
+val scala212 = "2.12.21"
+val scala213 = "2.13.18"
ThisBuild / organization := "za.co.absa.pramen"
ThisBuild / scalaVersion := scala212
-ThisBuild / crossScalaVersions := Seq(scala211, scala212, scala213)
+ThisBuild / crossScalaVersions := Seq(scala212, scala213)
ThisBuild / scalacOptions := Seq("-unchecked", "-deprecation", "-target:jvm-1.8")
-ThisBuild / javacOptions := Seq("-source", "1.8", "-target", "1.8")
+ThisBuild / javacOptions := Seq("-source", "1.8", "-target", "1.8", "--release", "8")
ThisBuild / versionScheme := Some("early-semver")
@@ -58,6 +57,32 @@ def runnerSparkVersionSuffix(moduleName: String, scalaVersion: String, includeDe
} else ""
}
+val projectJavaOptions = Seq(
+ "-XX:+IgnoreUnrecognizedVMOptions",
+ "-Xmx2048m",
+ "--add-modules=jdk.incubator.vector",
+ "--add-opens=java.base/java.lang=ALL-UNNAMED",
+ "--add-opens=java.base/java.lang.invoke=ALL-UNNAMED",
+ "--add-opens=java.base/java.lang.reflect=ALL-UNNAMED",
+ "--add-opens=java.base/java.io=ALL-UNNAMED",
+ "--add-opens=java.base/java.net=ALL-UNNAMED",
+ "--add-opens=java.base/java.nio=ALL-UNNAMED",
+ "--add-opens=java.base/java.util=ALL-UNNAMED",
+ "--add-opens=java.base/java.util.concurrent=ALL-UNNAMED",
+ "--add-opens=java.base/java.util.concurrent.atomic=ALL-UNNAMED",
+ "--add-opens=java.base/jdk.internal.ref=ALL-UNNAMED",
+ "--add-opens=java.base/sun.nio.ch=ALL-UNNAMED",
+ "--add-opens=java.base/sun.nio.cs=ALL-UNNAMED",
+ "--add-opens=java.base/sun.security.action=ALL-UNNAMED",
+ "--add-opens=java.base/sun.util.calendar=ALL-UNNAMED",
+ "--add-opens=java.base/sun.net.www.protocol.jar=ALL-UNNAMED",
+ "-Djdk.reflect.useDirectMethodHandle=false",
+ "-Dio.netty.tryReflectionSetAccessible=true",
+ "-Dio.netty.allocator.type=pooled",
+ "-Dio.netty.handler.ssl.defaultEndpointVerificationAlgorithm=NONE",
+ "--enable-native-access=ALL-UNNAMED"
+)
+
val assemblyFeatures = settingKey[Seq[String]]("Define assembly scope")
lazy val UnitTest = config("unit") extend Test
@@ -67,7 +92,7 @@ lazy val pramen = (project in file("."))
.disablePlugins(sbtassembly.AssemblyPlugin)
.settings(
name := "pramen",
- crossScalaVersions := List(scala211, scala212, scala213),
+ crossScalaVersions := List(scala212, scala213),
// No need to publish the aggregation [empty] artifact
publishArtifact := false,
@@ -85,7 +110,7 @@ lazy val api = (project in file("api"))
.settings( inConfig(IntegrationTest)(Defaults.testTasks) : _*)
.settings(
name := "pramen-api",
- crossScalaVersions := List(scala211, scala212, scala213),
+ crossScalaVersions := List(scala212, scala213),
printSparkVersion := {
val log = streams.value.log
log.info(s"Building with Spark ${sparkVersion(scalaVersion.value)}, Scala ${scalaVersion.value}")
@@ -106,7 +131,7 @@ lazy val core = (project in file("core"))
.settings( inConfig(IntegrationTest)(Defaults.testTasks) : _*)
.settings(
name := "pramen-core",
- crossScalaVersions := List(scala211, scala212, scala213),
+ crossScalaVersions := List(scala212, scala213),
printSparkVersion := {
val log = streams.value.log
log.info(s"Building with Spark ${sparkVersion(scalaVersion.value)}, Scala ${scalaVersion.value}")
@@ -116,20 +141,19 @@ lazy val core = (project in file("core"))
Compile / unmanagedSourceDirectories += {
val sourceDir = (Compile / sourceDirectory).value
CrossVersion.partialVersion(scalaVersion.value) match {
- case Some((2, n)) if n == 11 => sourceDir / "scala_2.11"
case Some((2, n)) if n == 12 => sourceDir / "scala_2.12"
case Some((2, n)) if n == 13 => sourceDir / "scala_2.13"
case _ => throw new RuntimeException("Unsupported Scala version")
}
},
assemblyFeatures := sys.props.getOrElse("assembly.features", "").split(',').toSeq,
- libraryDependencies ++= CoreDependencies(scalaVersion.value, assemblyFeatures.value.contains("includeDelta")) ++
- getSparkVersionRelatedDeps(sparkVersion(scalaVersion.value)) :+
+ libraryDependencies ++= CoreDependencies(scalaVersion.value, assemblyFeatures.value.contains("includeDelta")) :+
getScalaDependency(scalaVersion.value),
(Test / testOptions) := Seq(Tests.Filter(allFilter)),
(UnitTest / testOptions) := Seq(Tests.Filter(unitFilter)),
(IntegrationTest / testOptions) := Seq(Tests.Filter(itFilter)),
Test / fork := true,
+ Test / javaOptions ++= projectJavaOptions,
populateBuildInfoTemplate,
jacocoReportName := "pramen:core Jacoco Report",
jacocoReportFormats := Set("html", "xml"),
@@ -146,7 +170,7 @@ lazy val extras = (project in file("extras"))
.settings( inConfig(IntegrationTest)(Defaults.testTasks) : _*)
.settings(
name := "pramen-extras",
- crossScalaVersions := List(scala211, scala212, scala213),
+ crossScalaVersions := List(scala212, scala213),
printSparkVersion := {
val log = streams.value.log
log.info(s"Building with Spark ${sparkVersion(scalaVersion.value)}, Scala ${scalaVersion.value}")
@@ -154,14 +178,14 @@ lazy val extras = (project in file("extras"))
},
(Compile / compile) := ((Compile / compile) dependsOn printSparkVersion).value,
assemblyFeatures := sys.props.getOrElse("assembly.features", "").split(',').toSeq,
- libraryDependencies ++= ExtrasJobsDependencies(scalaVersion.value) ++
- getSparkVersionRelatedDeps(sparkVersion(scalaVersion.value)) :+
+ libraryDependencies ++= ExtrasJobsDependencies(scalaVersion.value) :+
getScalaDependency(scalaVersion.value),
resolvers += "confluent" at "https://packages.confluent.io/maven/",
(Test / testOptions) := Seq(Tests.Filter(allFilter)),
(UnitTest / testOptions) := Seq(Tests.Filter(unitFilter)),
(IntegrationTest / testOptions) := Seq(Tests.Filter(itFilter)),
Test / fork := true,
+ Test / javaOptions ++= projectJavaOptions,
jacocoReportName := "pramen-extras Jacoco Report",
jacocoReportFormats := Set("html", "xml"),
assemblySettingsExtras
diff --git a/pramen/core/pom.xml b/pramen/core/pom.xml
index 9df768c93..de56a8778 100644
--- a/pramen/core/pom.xml
+++ b/pramen/core/pom.xml
@@ -27,7 +27,7 @@
za.co.absa.pramen
pramen
- 1.14.9-SNAPSHOT
+ 1.15.0-SNAPSHOT
@@ -128,7 +128,7 @@
javax.mail
-
+
com.lihaoyi
requests_${scala.compat.version}
diff --git a/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/SparkCompatUtils.scala b/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/SparkCompatUtils.scala
new file mode 100644
index 000000000..312210c38
--- /dev/null
+++ b/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/SparkCompatUtils.scala
@@ -0,0 +1,76 @@
+/*
+ * Copyright 2022 ABSA Group Limited
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package za.co.absa.pramen.core.utils
+
+import org.apache.spark.sql.Column
+import org.apache.spark.sql.catalyst.expressions.Expression
+
+/**
+ * Compatibility layer for Column <-> Expression conversions across Spark versions.
+ *
+ * In Spark 3.x, Column has a direct `.expr` property and a constructor taking an Expression.
+ * In Spark 4.x, Column wraps a ColumnNode AST and conversions go through
+ * `org.apache.spark.sql.classic.ColumnConversions` and `ExpressionUtils`.
+ *
+ * This object uses reflection to support both versions from a single codebase,
+ * following the same pattern used in AbrisAvroDeserializer.
+ */
+object SparkCompatUtils {
+ /** Convert a Column to a Catalyst Expression. Replaces `column.expr`. */
+ lazy val col2expr: Column => Expression = {
+ // Try Spark 3.x first: Column.expr is a direct method
+ val spark3Method = try {
+ Some(classOf[Column].getMethod("expr"))
+ } catch {
+ case _: NoSuchMethodException => None
+ }
+
+ spark3Method match {
+ case Some(method) =>
+ (column: Column) => method.invoke(column).asInstanceOf[Expression]
+
+ case None =>
+ // Spark 4.x: use org.apache.spark.sql.classic.ColumnConversions.expression(column)
+ val clazz = Class.forName("org.apache.spark.sql.classic.ColumnConversions$")
+ val instance = clazz.getField("MODULE$").get(null)
+ val method = clazz.getMethod("expression", classOf[Column])
+ (column: Column) => method.invoke(instance, column).asInstanceOf[Expression]
+ }
+ }
+
+ /** Convert a Catalyst Expression to a Column. Replaces `new Column(expr)`. */
+ lazy val expr2col: Expression => Column = {
+ // Try Spark 3.x first: new Column(Expression)
+ val spark3Ctor = try {
+ Some(classOf[Column].getConstructor(classOf[Expression]))
+ } catch {
+ case _: NoSuchMethodException => None
+ }
+
+ spark3Ctor match {
+ case Some(ctor) =>
+ (expr: Expression) => ctor.newInstance(expr)
+
+ case None =>
+ // Spark 4.x: use ExpressionUtils.column(expr)
+ val clazz = Class.forName("org.apache.spark.sql.classic.ExpressionUtils$")
+ val instance = clazz.getField("MODULE$").get(null)
+ val method = clazz.getMethod("column", classOf[Expression])
+ (expr: Expression) => method.invoke(instance, expr).asInstanceOf[Column]
+ }
+ }
+}
diff --git a/pramen/core/src/main/scala_2.11/za/co/absa/pramen/core/bookkeeper/BookkeeperDeltaTable.scala b/pramen/core/src/main/scala_2.11/za/co/absa/pramen/core/bookkeeper/BookkeeperDeltaTable.scala
deleted file mode 100644
index ab1628609..000000000
--- a/pramen/core/src/main/scala_2.11/za/co/absa/pramen/core/bookkeeper/BookkeeperDeltaTable.scala
+++ /dev/null
@@ -1,142 +0,0 @@
-/*
- * Copyright 2022 ABSA Group Limited
- *
- * Licensed under the Apache License, Version 2.0 (the "License");
- * you may not use this file except in compliance with the License.
- * You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package za.co.absa.pramen.core.bookkeeper
-
-import org.apache.spark.sql.functions.col
-import org.apache.spark.sql.{AnalysisException, Column, Dataset, SaveMode, SparkSession}
-import za.co.absa.pramen.core.bookkeeper.model.TableSchemaJson
-import za.co.absa.pramen.core.model.{DataChunk, TableSchema}
-
-import java.time.{Instant, LocalDate}
-import scala.reflect.ClassTag
-import scala.reflect.runtime.universe
-
-object BookkeeperDeltaTable {
- val recordsTable = "bookkeeping"
- val schemasTable = "schemas"
-
- def getFullTableName(databaseOpt: Option[String], tablePrefix: String, tableName: String): String = {
- databaseOpt match {
- case Some(db) => s"$db.$tablePrefix$tableName"
- case None => s"$tablePrefix$tableName"
- }
- }
-}
-
-class BookkeeperDeltaTable(database: Option[String],
- tablePrefix: String,
- batchId: Long)
- (implicit spark: SparkSession) extends BookkeeperDeltaBase(batchId) {
- import BookkeeperDeltaTable._
- import spark.implicits._
-
- private val recordsFullTableName = getFullTableName(database, tablePrefix, recordsTable)
- private val schemasFullTableName = getFullTableName(database, tablePrefix, schemasTable)
-
- init()
-
- override def getBkDf(filter: Column): Dataset[DataChunk] = {
- val df = try {
- spark.table(recordsFullTableName).as[DataChunk]
- } catch {
- case ex: AnalysisException if ex.getMessage().contains("cannot resolve") =>
- // Spark 2 and 3
- migrateModel()
- spark.table(recordsFullTableName).as[DataChunk]
-
- case ex: Throwable if ex.getMessage.contains("UNRESOLVED_COLUMN") =>
- // Spark 3 and 4
- migrateModel()
- spark.table(recordsFullTableName).as[DataChunk]
- }
-
- df.filter(filter)
- .orderBy(col("jobFinished"))
- .as[DataChunk]
- }
-
- override def saveRecordCountDelta(dataChunk: DataChunk): Unit = {
- val df = Seq(dataChunk).toDF()
-
- df.write
- .format("delta")
- .mode(SaveMode.Append)
- .option("mergeSchema", "true")
- .saveAsTable(recordsFullTableName)
- }
-
- override def deleteNonCurrentBatchRecords(table: String, infoDate: LocalDate): Unit = {
- // This is not supported for Delta tables using Spark 2.*
- }
-
- override def getSchemasDeltaDf: Dataset[TableSchemaJson] = {
- spark.table(schemasFullTableName).as[TableSchemaJson]
- }
-
- override def saveSchemaDelta(schema: TableSchema): Unit = {
- val df = Seq(
- TableSchemaJson(schema.tableName, schema.infoDate, schema.schemaJson, Instant.now().toEpochMilli)
- ).toDF()
-
- df.write
- .format("delta")
- .mode(SaveMode.Append)
- .option("mergeSchema", "true")
- .saveAsTable(schemasFullTableName)
- }
-
- override def writeEmptyDataset[T <: Product : universe.TypeTag : ClassTag](pathOrTable: String): Unit = {
- val df = Seq.empty[T].toDS
-
- df.write
- .format("delta")
- .saveAsTable(pathOrTable)
- }
-
- override def deleteTable(tableWithWildcard: String): Seq[String] = ???
-
- def init(): Unit = {
- initRecordsDirectory()
- initSchemasDirectory()
- }
-
- private def initRecordsDirectory(): Unit = {
- if (!spark.catalog.tableExists(recordsFullTableName)) {
- writeEmptyDataset[DataChunk](recordsFullTableName)
- }
- }
-
- private def initSchemasDirectory(): Unit = {
- if (!spark.catalog.tableExists(schemasFullTableName)) {
- writeEmptyDataset[TableSchemaJson](schemasFullTableName)
- }
- }
-
- private def migrateModel(): Unit = {
- migrateModelViaEmptyDataset[DataChunk](recordsFullTableName)
- }
-
- private def migrateModelViaEmptyDataset[T <: Product : universe.TypeTag : ClassTag](pathOrTable: String): Unit = {
- val df = Seq.empty[T].toDS
-
- df.write
- .format("delta")
- .mode(SaveMode.Append)
- .option("mergeSchema", "true")
- .saveAsTable(pathOrTable)
- }
-}
diff --git a/pramen/core/src/main/scala_2.11/za/co/absa/pramen/core/metastore/peristence/MetastorePersistenceIcebergOps.scala b/pramen/core/src/main/scala_2.11/za/co/absa/pramen/core/metastore/peristence/MetastorePersistenceIcebergOps.scala
deleted file mode 100644
index e81c8a64c..000000000
--- a/pramen/core/src/main/scala_2.11/za/co/absa/pramen/core/metastore/peristence/MetastorePersistenceIcebergOps.scala
+++ /dev/null
@@ -1,70 +0,0 @@
-/*
- * Copyright 2022 ABSA Group Limited
- *
- * Licensed under the Apache License, Version 2.0 (the "License");
- * you may not use this file except in compliance with the License.
- * You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package za.co.absa.pramen.core.metastore.peristence
-
-import org.apache.spark.sql.functions._
-import org.apache.spark.sql.types.DateType
-import org.apache.spark.sql.{Column, DataFrame, SaveMode, SparkSession}
-import org.slf4j.LoggerFactory
-import za.co.absa.pramen.api.{CatalogTable, PartitionInfo, PartitionScheme}
-
-import java.sql.Date
-import java.time.LocalDate
-import java.time.format.DateTimeFormatter
-
-object MetastorePersistenceIcebergOps {
- def createIcebergTable(df: DataFrame,
- table: String,
- infoDateColumn: String,
- location: Option[String] = None,
- description: String = "",
- partitionScheme: PartitionScheme,
- tableProperties: Map[String, String],
- writerOptions: Map[String, String]): Unit = {
- throw new UnsupportedOperationException(s"Iceberg format is not supported in Scala 2.11")
- }
-
- def overwriteDailyPartition(infoDate: LocalDate,
- df: DataFrame,
- table: String,
- infoDateColumn: String,
- writerOptions: Map[String, String]): Unit = {
- throw new UnsupportedOperationException(s"Iceberg format is not supported in Scala 2.11")
- }
-
- def overwriteFullTable(df: DataFrame,
- table: String,
- writerOptions: Map[String, String]): Unit = {
- throw new UnsupportedOperationException(s"Iceberg format is not supported in Scala 2.11")
- }
-
- def appendToTable(df: DataFrame,
- table: String,
- writerOptions: Map[String, String]): Unit = {
- throw new UnsupportedOperationException(s"Iceberg format is not supported in Scala 2.11")
- }
-
- def writeRepartitionedDf(df: DataFrame,
- table: String,
- infoDateColumn: String,
- infoDateFrom: LocalDate,
- infoDateTo: LocalDate,
- writerOptions: Map[String, String]): Unit = {
- throw new UnsupportedOperationException(s"Iceberg format is not supported in Scala 2.11")
- }
-
-}
diff --git a/pramen/core/src/test/scala/za/co/absa/pramen/core/metastore/MetastoreSuite.scala b/pramen/core/src/test/scala/za/co/absa/pramen/core/metastore/MetastoreSuite.scala
index 393bd563a..4916a1967 100644
--- a/pramen/core/src/test/scala/za/co/absa/pramen/core/metastore/MetastoreSuite.scala
+++ b/pramen/core/src/test/scala/za/co/absa/pramen/core/metastore/MetastoreSuite.scala
@@ -487,31 +487,19 @@ class MetastoreSuite extends AnyWordSpec with SparkTestBase with TextComparisonF
"withSparkConfig()" should {
"set the config at runtime, and restore the original config afterwards" in {
val sparkConfig = Map(
- "spark.sql.sources.commitProtocolClass" -> "org.apache.spark.internal.io.cloud.PathOutputCommitProtocol",
- "spark.sql.parquet.output.committer.class" -> "org.apache.spark.internal.io.cloud.BindingParquetOutputCommitter",
"spark.pramen.test" -> "test"
)
- var inner1: String = null
- var inner2: String = null
- var inner3: String = null
+ var inner: String = null
MetastoreImpl.withSparkConfig(sparkConfig) {
- inner1 = spark.conf.get("spark.sql.sources.commitProtocolClass")
- inner2 = spark.conf.get("spark.sql.parquet.output.committer.class")
- inner3 = spark.conf.get("spark.pramen.test")
+ inner = spark.conf.get("spark.pramen.test")
}
- val outer1 = spark.conf.get("spark.sql.sources.commitProtocolClass")
- val outer2 = spark.conf.get("spark.sql.parquet.output.committer.class")
- val outer3 = spark.conf.getOption("spark.pramen.test")
+ val outer = spark.conf.getOption("spark.pramen.test")
- assert(inner1 == "org.apache.spark.internal.io.cloud.PathOutputCommitProtocol")
- assert(inner2 == "org.apache.spark.internal.io.cloud.BindingParquetOutputCommitter")
- assert(inner3 == "test")
- assert(outer1 != "org.apache.spark.internal.io.cloud.PathOutputCommitProtocol")
- assert(outer2 != "org.apache.spark.internal.io.cloud.BindingParquetOutputCommitter")
- assert(outer3.isEmpty)
+ assert(inner == "test")
+ assert(outer.isEmpty)
}
}
diff --git a/pramen/core/src/test/scala/za/co/absa/pramen/core/pipeline/SinkJobSuite.scala b/pramen/core/src/test/scala/za/co/absa/pramen/core/pipeline/SinkJobSuite.scala
index a929a9419..a17352e64 100644
--- a/pramen/core/src/test/scala/za/co/absa/pramen/core/pipeline/SinkJobSuite.scala
+++ b/pramen/core/src/test/scala/za/co/absa/pramen/core/pipeline/SinkJobSuite.scala
@@ -42,7 +42,7 @@ class SinkJobSuite extends AnyWordSpec with SparkTestBase with TextComparisonFix
private val conf = ConfigFactory.empty()
private val runReason: TaskRunReason = TaskRunReason.New
- private def exampleDf: DataFrame = List(("A", 1), ("B", 2), ("C", 3)).toDF("a", "b")
+ private def exampleDf: DataFrame = List(("A", 1), ("B", 2), ("C", 3)).toDF("b", "a")
"preRunCheckJob" should {
"return Ready when the input table is available" in {
@@ -83,7 +83,7 @@ class SinkJobSuite extends AnyWordSpec with SparkTestBase with TextComparisonFix
}
"return Skip when the data frame is empty" in {
- val (job, _) = getUseCase(tableDf = exampleDf.filter(col("b") > 10))
+ val (job, _) = getUseCase(tableDf = exampleDf.filter(col("a") > 10))
val result = job.validate(infoDate, runReason, conf)
@@ -138,16 +138,16 @@ class SinkJobSuite extends AnyWordSpec with SparkTestBase with TextComparisonFix
"apply transformations, filters and projections" in {
val expectedData =
"""[ {
- | "a" : "B",
- | "b1" : "2"
+ | "a" : 2,
+ | "b1" : "B"
|}, {
- | "a" : "C",
- | "b1" : "3"
+ | "a" : 3,
+ | "b1" : "C"
|} ]""".stripMargin
val sinkTable = SinkTableFactory.getDummySinkTable(
transformations = Seq(TransformExpression("b1", Some("cast(b as string)"), None)),
- filters = Seq("b > 1"),
+ filters = Seq("a > 1"),
columns = Seq("a", "b1")
)
diff --git a/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/bookkeeper/BookkeeperDeltaTableLongSuite.scala b/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/bookkeeper/BookkeeperDeltaTableLongSuite.scala
index 94e0c12cc..5248d5e3b 100644
--- a/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/bookkeeper/BookkeeperDeltaTableLongSuite.scala
+++ b/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/bookkeeper/BookkeeperDeltaTableLongSuite.scala
@@ -59,11 +59,7 @@ class BookkeeperDeltaTableLongSuite extends BookkeeperCommonSuite with SparkTest
"BookkeeperHadoopDeltaTable" when {
testBookKeeper { batchId =>
- if (spark.version.startsWith("2.")) {
- getBookkeeper(getNewTablePrefix, batchId)
- } else {
- getBookkeeper(bookkeepingTablePrefix, batchId)
- }
+ getBookkeeper(bookkeepingTablePrefix, batchId)
}
"test tables are created properly" in {
diff --git a/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/bookkeeper/BookkeeperTextLongSuite.scala b/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/bookkeeper/BookkeeperTextLongSuite.scala
index 6722c6f51..05e0fa09f 100644
--- a/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/bookkeeper/BookkeeperTextLongSuite.scala
+++ b/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/bookkeeper/BookkeeperTextLongSuite.scala
@@ -64,7 +64,10 @@ class BookkeeperTextLongSuite extends BookkeeperCommonSuite with SparkTestBase w
val actual = bk.getFilter("table1", Some(LocalDate.of(2021, 1, 1)), Some(LocalDate.of(2021, 1, 2)), None).toString()
- assert(actual == "(((tableName = table1) AND (infoDate >= 2021-01-01)) AND (infoDate <= 2021-01-02))")
+ val possible1 = "(((tableName = table1) AND (infoDate >= 2021-01-01)) AND (infoDate <= 2021-01-02))"
+ val possible2 = "and(and(=(tableName, 'table1'), >=(infoDate, '2021-01-01')), <=(infoDate, '2021-01-02'))"
+
+ assert(actual == possible1 || actual == possible2)
}
"get a ranged filter with batch id" in {
@@ -72,7 +75,10 @@ class BookkeeperTextLongSuite extends BookkeeperCommonSuite with SparkTestBase w
val actual = bk.getFilter("table1", Some(LocalDate.of(2021, 1, 1)), Some(LocalDate.of(2021, 1, 2)), Some(123L)).toString()
- assert(actual == "((((tableName = table1) AND (infoDate >= 2021-01-01)) AND (infoDate <= 2021-01-02)) AND (batchId = 123))")
+ val possible1 = "((((tableName = table1) AND (infoDate >= 2021-01-01)) AND (infoDate <= 2021-01-02)) AND (batchId = 123))"
+ val possible2 = "and(and(and(=(tableName, 'table1'), >=(infoDate, '2021-01-01')), <=(infoDate, '2021-01-02')), =(batchId, 123L))"
+
+ assert(actual == possible1 || actual == possible2)
}
"get a from filter" in {
@@ -80,7 +86,10 @@ class BookkeeperTextLongSuite extends BookkeeperCommonSuite with SparkTestBase w
val actual = bk.getFilter("table1", Some(LocalDate.of(2021, 1, 1)), None, None).toString()
- assert(actual == "((tableName = table1) AND (infoDate >= 2021-01-01))")
+ val possible1 = "((tableName = table1) AND (infoDate >= 2021-01-01))"
+ val possible2 = "and(=(tableName, 'table1'), >=(infoDate, '2021-01-01'))"
+
+ assert(actual == possible1 || actual == possible2)
}
"get a to filter" in {
@@ -88,7 +97,10 @@ class BookkeeperTextLongSuite extends BookkeeperCommonSuite with SparkTestBase w
val actual = bk.getFilter("table1", None, Some(LocalDate.of(2021, 1, 2)), None).toString()
- assert(actual == "((tableName = table1) AND (infoDate <= 2021-01-02))")
+ val possible1 = "((tableName = table1) AND (infoDate <= 2021-01-02))"
+ val possible2 = "and(=(tableName, 'table1'), <=(infoDate, '2021-01-02'))"
+
+ assert(actual == possible1 || actual == possible2)
}
"get a batchid filter" in {
@@ -96,7 +108,10 @@ class BookkeeperTextLongSuite extends BookkeeperCommonSuite with SparkTestBase w
val actual = bk.getFilter("table1", None, None, Some(123L)).toString()
- assert(actual == "((tableName = table1) AND (batchId = 123))")
+ val possible1 = "((tableName = table1) AND (batchId = 123))"
+ val possible2 = "and(=(tableName, 'table1'), =(batchId, 123L))"
+
+ assert(actual == possible1 || actual == possible2)
}
"get a table filter" in {
@@ -104,7 +119,10 @@ class BookkeeperTextLongSuite extends BookkeeperCommonSuite with SparkTestBase w
val actual = bk.getFilter("table1", None, None, None).toString()
- assert(actual == "(tableName = table1)")
+ val possible1 = "(tableName = table1)"
+ val possible2 = "=(tableName, 'table1')"
+
+ assert(actual == possible1 || actual == possible2)
}
}
}
diff --git a/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/bookkeeper/OffsetManagerJdbcSuite.scala b/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/bookkeeper/OffsetManagerJdbcSuite.scala
index 4b759cb3f..1724585d3 100644
--- a/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/bookkeeper/OffsetManagerJdbcSuite.scala
+++ b/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/bookkeeper/OffsetManagerJdbcSuite.scala
@@ -80,7 +80,7 @@ class OffsetManagerJdbcSuite extends AnyWordSpec with RelationalDbFixture with B
val offset = actualNonEmpty.head.asInstanceOf[UncommittedOffset]
assert(offset.infoDate == infoDate)
- assert(!offset.createdAt.isBefore(now))
+ assert(offset.createdAt.toEpochMilli >= now.toEpochMilli)
assert(offset.createdAt.isBefore(nextHour))
}
diff --git a/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/utils/SparkCompatUtilsSuite.scala b/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/utils/SparkCompatUtilsSuite.scala
new file mode 100644
index 000000000..409386da8
--- /dev/null
+++ b/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/utils/SparkCompatUtilsSuite.scala
@@ -0,0 +1,128 @@
+/*
+ * Copyright 2022 ABSA Group Limited
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package za.co.absa.pramen.core.tests.utils
+
+import org.scalatest.wordspec.AnyWordSpec
+import za.co.absa.pramen.core.base.SparkTestBase
+import za.co.absa.pramen.core.utils.SparkCompatUtils
+
+class SparkCompatUtilsSuite extends AnyWordSpec with SparkTestBase {
+ "col2expr" should {
+ "convert a column reference to a Catalyst expression" in {
+ val expr = SparkCompatUtils.col2expr(org.apache.spark.sql.functions.col("a"))
+
+ assert(expr != null)
+ assert(expr.isInstanceOf[org.apache.spark.sql.catalyst.expressions.Expression])
+ assert(expr.sql.contains("a"))
+ }
+
+ "convert a literal to a Catalyst expression" in {
+ val expr = SparkCompatUtils.col2expr(org.apache.spark.sql.functions.lit(42))
+
+ assert(expr != null)
+ assert(expr.sql.contains("42"))
+ }
+
+ "convert a complex expression" in {
+ val column = (org.apache.spark.sql.functions.col("a") + 1) * 2
+ val expr = SparkCompatUtils.col2expr(column)
+
+ assert(expr != null)
+ assert(expr.sql.contains("a"))
+ assert(expr.sql.contains("1"))
+ assert(expr.sql.contains("2"))
+ }
+
+ "be usable multiple times since it is a lazy val" in {
+ val expr1 = SparkCompatUtils.col2expr(org.apache.spark.sql.functions.col("a"))
+ val expr2 = SparkCompatUtils.col2expr(org.apache.spark.sql.functions.col("b"))
+
+ assert(expr1.sql.contains("a"))
+ assert(expr2.sql.contains("b"))
+ }
+ }
+
+ "expr2col" should {
+ "convert a Catalyst expression to a column" in {
+ val expr = org.apache.spark.sql.catalyst.expressions.Literal(42)
+ val column = SparkCompatUtils.expr2col(expr)
+
+ assert(column != null)
+ assert(column.isInstanceOf[org.apache.spark.sql.Column])
+ assert(column.toString.contains("42"))
+ }
+
+ "convert an unresolved attribute to a column" in {
+ val expr = org.apache.spark.sql.catalyst.analysis.UnresolvedAttribute(Seq("a"))
+ val column = SparkCompatUtils.expr2col(expr)
+
+ assert(column != null)
+ assert(column.toString.contains("a"))
+ }
+
+ "produce a column that can be used in a query" in {
+ import spark.implicits._
+
+ val df = List((1, "x"), (2, "y")).toDF("id", "name")
+ val column = SparkCompatUtils.expr2col(org.apache.spark.sql.catalyst.expressions.Literal(7))
+
+ val actual = df.select(column.as("value")).collect().map(_.getInt(0)).toSeq
+
+ assert(actual == Seq(7, 7))
+ }
+ }
+
+ "col2expr and expr2col" should {
+ "round trip a simple column" in {
+ val original = org.apache.spark.sql.functions.col("a")
+ val actual = SparkCompatUtils.expr2col(SparkCompatUtils.col2expr(original))
+
+ assert(actual.toString == original.toString)
+ }
+
+ "round trip a literal column" in {
+ val original = org.apache.spark.sql.functions.lit("abc")
+ val actual = SparkCompatUtils.expr2col(SparkCompatUtils.col2expr(original))
+
+ assert(actual.toString.contains("abc"))
+ }
+
+ "round trip a complex column used in a query" in {
+ import spark.implicits._
+
+ val df = List((1, "x"), (2, "y")).toDF("id", "name")
+ val original = (org.apache.spark.sql.functions.col("id") + 1) * 2
+ val roundTripped = SparkCompatUtils.expr2col(SparkCompatUtils.col2expr(original))
+
+ val expected = df.select(original.as("v")).collect().map(_.getInt(0)).toSeq
+ val actual = df.select(roundTripped.as("v")).collect().map(_.getInt(0)).toSeq
+
+ assert(actual == expected)
+ assert(actual == Seq(4, 6))
+ }
+
+ "round trip an expression" in {
+ val original: org.apache.spark.sql.catalyst.expressions.Expression =
+ org.apache.spark.sql.catalyst.expressions.Literal(10)
+ val actual = SparkCompatUtils.col2expr(SparkCompatUtils.expr2col(original))
+
+ assert(actual.sql == original.sql)
+ }
+ }
+
+}
+
diff --git a/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/utils/StringUtilsSuite.scala b/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/utils/StringUtilsSuite.scala
index 15610263e..3580c12cc 100644
--- a/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/utils/StringUtilsSuite.scala
+++ b/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/utils/StringUtilsSuite.scala
@@ -251,9 +251,9 @@ class StringUtilsSuite extends AnyWordSpec {
val actual = renderThreadDumps(ex.threadStackTraces)
assert(actual.startsWith("Stack trace of threads at the moment of the interruption:"))
- assert(actual.contains(" Thread 0"))
+ assert(actual.contains("Thread 0"))
assert(actual.contains("(ScalaTest-dispatcher)"))
- assert(actual.contains(" java.lang.Thread.dumpThreads(Native Method)"))
+ assert(actual.contains("java.lang.Thread.dumpThreads(Native Method)"))
}
}
diff --git a/pramen/examples/combined_example.sh b/pramen/examples/combined_example.sh
index 34226a1a4..145c998d1 100755
--- a/pramen/examples/combined_example.sh
+++ b/pramen/examples/combined_example.sh
@@ -16,7 +16,7 @@
# Prerequisites:
# 1. Download Spark 3.4.1 (Scala 2.12) and install it in /opt/spark/spark-3.4.1 or some other directory
# 2. At repo_root/pramen, run
-# sbt -DSPARK_VERSION="3.4.1" ++2.12.20 assembly
+# sbt -DSPARK_VERSION="3.4.1" ++2.12.21 assembly
# 3. Run
# ./examples/combined_example.sh
diff --git a/pramen/examples/enceladus_single_config/daily_ingestion.sh b/pramen/examples/enceladus_single_config/daily_ingestion.sh
index 6df2339de..a0449c6eb 100644
--- a/pramen/examples/enceladus_single_config/daily_ingestion.sh
+++ b/pramen/examples/enceladus_single_config/daily_ingestion.sh
@@ -16,7 +16,7 @@ ME=`basename "$0" .sh`
cd $(dirname $(readlink -f $0))
-SCALA_VERSION="2.11"
+SCALA_VERSION="2.12"
EXTRAS_JAR="pramen-extras_${SCALA_VERSION}-1.0.0.jar"
RUNNER_JAR="pramen-runner_${SCALA_VERSION}-1.0.0.jar"
diff --git a/pramen/examples/enceladus_single_config/weekly_ingestion.sh b/pramen/examples/enceladus_single_config/weekly_ingestion.sh
index 6df2339de..a0449c6eb 100644
--- a/pramen/examples/enceladus_single_config/weekly_ingestion.sh
+++ b/pramen/examples/enceladus_single_config/weekly_ingestion.sh
@@ -16,7 +16,7 @@ ME=`basename "$0" .sh`
cd $(dirname $(readlink -f $0))
-SCALA_VERSION="2.11"
+SCALA_VERSION="2.12"
EXTRAS_JAR="pramen-extras_${SCALA_VERSION}-1.0.0.jar"
RUNNER_JAR="pramen-runner_${SCALA_VERSION}-1.0.0.jar"
diff --git a/pramen/examples/enceladus_sourcing/daily_ingestion.sh b/pramen/examples/enceladus_sourcing/daily_ingestion.sh
index 69cbce261..5443cd4a6 100644
--- a/pramen/examples/enceladus_sourcing/daily_ingestion.sh
+++ b/pramen/examples/enceladus_sourcing/daily_ingestion.sh
@@ -17,7 +17,7 @@ ME=`basename "$0" .sh`
cd $(dirname $(readlink -f $0))
-SCALA_VERSION="2.11"
+SCALA_VERSION="2.12"
EXTRAS_JAR="pramen-extras_${SCALA_VERSION}-1.0.0.jar"
RUNNER_JAR="pramen-runner_${SCALA_VERSION}-1.0.0.jar"
diff --git a/pramen/examples/enceladus_sourcing/daily_snapshot.sh b/pramen/examples/enceladus_sourcing/daily_snapshot.sh
index bb2213477..e7fb39658 100644
--- a/pramen/examples/enceladus_sourcing/daily_snapshot.sh
+++ b/pramen/examples/enceladus_sourcing/daily_snapshot.sh
@@ -17,7 +17,7 @@ ME=`basename "$0" .sh`
cd $(dirname $(readlink -f $0))
-SCALA_VERSION="2.11"
+SCALA_VERSION="2.12"
EXTRAS_JAR="pramen-extras_${SCALA_VERSION}-1.0.0.jar"
RUNNER_JAR="pramen-runner_${SCALA_VERSION}-1.0.0.jar"
diff --git a/pramen/examples/jdbc_sourcing/daily_ingestion.sh b/pramen/examples/jdbc_sourcing/daily_ingestion.sh
index 6ef6361c4..58bb6df07 100644
--- a/pramen/examples/jdbc_sourcing/daily_ingestion.sh
+++ b/pramen/examples/jdbc_sourcing/daily_ingestion.sh
@@ -17,7 +17,7 @@ ME=`basename "$0" .sh`
cd $(dirname $(readlink -f $0))
-SCALA_VERSION="2.11"
+SCALA_VERSION="2.12"
EXTRAS_JAR="pramen-extras_${SCALA_VERSION}-1.0.0.jar"
RUNNER_JAR="pramen-runner_${SCALA_VERSION}-1.0.0.jar"
diff --git a/pramen/extras/pom.xml b/pramen/extras/pom.xml
index 6c3075e0b..27675e6af 100644
--- a/pramen/extras/pom.xml
+++ b/pramen/extras/pom.xml
@@ -27,7 +27,7 @@
za.co.absa.pramen
pramen
- 1.14.9-SNAPSHOT
+ 1.15.0-SNAPSHOT
diff --git a/pramen/extras/src/main/scala/za/co/absa/pramen/extras/writer/TableWriterKafka.scala b/pramen/extras/src/main/scala/za/co/absa/pramen/extras/writer/TableWriterKafka.scala
index 46f621066..6ebb26deb 100644
--- a/pramen/extras/src/main/scala/za/co/absa/pramen/extras/writer/TableWriterKafka.scala
+++ b/pramen/extras/src/main/scala/za/co/absa/pramen/extras/writer/TableWriterKafka.scala
@@ -23,6 +23,7 @@ import org.slf4j.LoggerFactory
import za.co.absa.abris.avro.functions.to_avro
import za.co.absa.abris.avro.read.confluent.SchemaManagerFactory
import za.co.absa.abris.config.{AbrisConfig, ToAvroConfig}
+import za.co.absa.pramen.core.utils.SparkCompatUtils
import za.co.absa.pramen.extras.avro.AvroUtils.{convertSparkToAvroSchema, fixNullableFields}
import za.co.absa.pramen.extras.source.KafkaAvroSource.KAFKA_TOKENS_TO_REDACT
import za.co.absa.pramen.extras.utils.ConfigUtils
@@ -102,7 +103,7 @@ class TableWriterKafka(topicName: String,
namingStrategy: NamingStrategy,
isKey: Boolean): Int = {
// generate schema
- val expression = columns.expr
+ val expression = SparkCompatUtils.col2expr(columns)
val schema = fixNullableFields(
convertSparkToAvroSchema(expression.dataType)
)
diff --git a/pramen/extras/src/test/scala/za/co/absa/pramen/extras/avro/AvroUtilsSuite.scala b/pramen/extras/src/test/scala/za/co/absa/pramen/extras/avro/AvroUtilsSuite.scala
index 2cfbf37b2..dc10a5b6d 100644
--- a/pramen/extras/src/test/scala/za/co/absa/pramen/extras/avro/AvroUtilsSuite.scala
+++ b/pramen/extras/src/test/scala/za/co/absa/pramen/extras/avro/AvroUtilsSuite.scala
@@ -18,6 +18,7 @@ package za.co.absa.pramen.extras.avro
import org.apache.spark.sql.functions.struct
import org.scalatest.wordspec.AnyWordSpec
+import za.co.absa.pramen.core.utils.SparkCompatUtils
import za.co.absa.pramen.extras.NestedDataFrameFactory
import za.co.absa.pramen.extras.base.SparkTestBase
import za.co.absa.pramen.extras.fixtures.TextComparisonFixture
@@ -29,11 +30,13 @@ class AvroUtilsSuite extends AnyWordSpec with SparkTestBase with TextComparisonF
"convertSparkToAvroSchema" should {
"convert basic schema with nullable values" in {
+ assume(spark.version.split('.').head.toInt < 4, s"Ignored for Spark ${spark.version}")
+
val df = List(("A", 1), ("B", 2), ("C", 3)).toDF("a", "b")
val allColumns = struct(df.columns.map(c => df(c)): _*)
- val avro = AvroUtils.convertSparkToAvroSchema(allColumns.expr.dataType)
+ val avro = AvroUtils.convertSparkToAvroSchema(SparkCompatUtils.col2expr(allColumns).dataType)
val avroWithNullsFixed = AvroUtils.fixNullableFields(avro)
@@ -49,11 +52,13 @@ class AvroUtilsSuite extends AnyWordSpec with SparkTestBase with TextComparisonF
}
"convert nested schema with nullable values" in {
+ assume(spark.version.split('.').head.toInt < 4, s"Ignored for Spark ${spark.version}")
+
val df = NestedDataFrameFactory.getNestedTestCase
val allColumns = struct(df.columns.map(c => df(c)): _*)
- val avro = AvroUtils.convertSparkToAvroSchema(allColumns.expr.dataType)
+ val avro = AvroUtils.convertSparkToAvroSchema(SparkCompatUtils.col2expr(allColumns).dataType)
val avroWithNullsFixed = AvroUtils.fixNullableFields(avro)
@@ -68,11 +73,13 @@ class AvroUtilsSuite extends AnyWordSpec with SparkTestBase with TextComparisonF
}
"convert nested schema with a map" in {
+ assume(spark.version.split('.').head.toInt < 4, s"Ignored for Spark ${spark.version}")
+
val df = NestedDataFrameFactory.getMapTestCase
val allColumns = struct(df.columns.map(c => df(c)): _*)
- val avro = AvroUtils.convertSparkToAvroSchema(allColumns.expr.dataType)
+ val avro = AvroUtils.convertSparkToAvroSchema(SparkCompatUtils.col2expr(allColumns).dataType)
val avroWithNullsFixed = AvroUtils.fixNullableFields(avro)
diff --git a/pramen/pom.xml b/pramen/pom.xml
index bbd0b3a64..8900107db 100644
--- a/pramen/pom.xml
+++ b/pramen/pom.xml
@@ -22,7 +22,7 @@
za.co.absa.pramen
pramen
- 1.14.9-SNAPSHOT
+ 1.15.0-SNAPSHOT
pom
Pramen
@@ -92,14 +92,6 @@
UTF-8
${maven.build.timestamp}
-
-
-
-
-
-
-
-
@@ -109,22 +101,24 @@
- 2.12.20
- 2.12
- 3.5.8
- 3.5
- 3.3.6
- 3.3.2
- 6.4.1
-
-
-
-
-
-
-
-
-
+
+
+
+
+
+
+
+
+
+
+ 2.13.18
+ 2.13
+ 4.1.2
+ 4.1
+ 3.4.3
+ 4.3.1
+ 7.0.0-RC1
+ 1.11.0
yyyy-MM-dd'T'HH:mm:ssX
@@ -337,7 +331,7 @@
org.apache.iceberg
iceberg-spark-runtime-${spark.compat.version}_${scala.compat.version}
- 1.6.1
+ ${iceberg.version}
@@ -353,7 +347,7 @@
1.3.1
-
+
com.github.yruslan
channel_scala_${scala.compat.version}
diff --git a/pramen/project/Versions.scala b/pramen/project/Versions.scala
index 858c1dd50..177a8c7dc 100644
--- a/pramen/project/Versions.scala
+++ b/pramen/project/Versions.scala
@@ -17,9 +17,8 @@
import sbt.*
object Versions {
- val defaultSparkVersionForScala211 = "2.4.8"
- val defaultSparkVersionForScala212 = "3.3.4"
- val defaultSparkVersionForScala213 = "3.4.4"
+ val defaultSparkVersionForScala212 = "3.5.5"
+ val defaultSparkVersionForScala213 = "4.1.2"
val typesafeConfigVersion = "1.4.3"
val postgreSqlDriverVersion = "42.7.9"
@@ -41,7 +40,7 @@ object Versions {
def sparkFallbackVersion(scalaVersion: String): String = {
if (scalaVersion.startsWith("2.11.")) {
- defaultSparkVersionForScala211
+ throw new IllegalArgumentException(s"Scala 2.11 not supported.")
} else if (scalaVersion.startsWith("2.12.")) {
defaultSparkVersionForScala212
} else if (scalaVersion.startsWith("2.13.")) {
@@ -59,17 +58,6 @@ object Versions {
fullVersion.split('.').take(2).mkString(".")
}
- def getSparkVersionRelatedDeps(sparkVersion: String): Seq[ModuleID] = {
- if (sparkVersion.startsWith("2.")) {
- // Seq("com.fasterxml.jackson.core" % "jackson-databind" % "2.6.7.3")
- Nil
- } else if (sparkVersion.startsWith("3.")) {
- Nil
- } else {
- throw new IllegalArgumentException(s"Spark $sparkVersion not supported.")
- }
- }
-
def getDeltaDependency(sparkVersion: String, isCompile: Boolean, isTest: Boolean): ModuleID = {
// According to this: https://docs.delta.io/latest/releases.html
val (deltaArtifact, deltaVersion) = sparkVersion match {
@@ -80,6 +68,8 @@ object Versions {
case version if version.startsWith("3.3.") => ("delta-core", "2.2.0")
case version if version.startsWith("3.4.") => ("delta-core", "2.4.0")
case version if version.startsWith("3.5.") => ("delta-spark", "3.0.0") // 'delta-core' was renamed to 'delta-spark' since 3.0.0.
+ case version if version.startsWith("4.0.") => ("delta-spark_4.0", "4.3.1")
+ case version if version.startsWith("4.1.") => ("delta-spark", "4.3.1")
case _ => throw new IllegalArgumentException(s"Spark $sparkVersion not supported.")
}
if (isTest) {
@@ -104,6 +94,8 @@ object Versions {
case version if version.startsWith("3.3.") => ("1.6.1", "3.3")
case version if version.startsWith("3.4.") => ("1.6.1", "3.4")
case version if version.startsWith("3.5.") => ("1.6.1", "3.5")
+ case version if version.startsWith("4.0.") => ("1.10.1", "4.0")
+ case version if version.startsWith("4.1.") => ("1.11.0", "4.1")
case _ => throw new IllegalArgumentException(s"Spark $sparkVersion not supported.")
}
@@ -130,6 +122,7 @@ object Versions {
case version if version.startsWith("3.1.") => "5.1.1"
case version if version == "3.2.0" => "6.1.1"
case version if version.startsWith("3.") => "6.4.1"
+ case version if version.startsWith("4.") => "7.0.0-RC1"
case _ => throw new IllegalArgumentException(s"Spark $sparkVersion not supported for Abris dependency.")
}