Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions .github/workflows/jacoco.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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: |
Expand Down Expand Up @@ -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
Expand Down
16 changes: 6 additions & 10 deletions .github/workflows/scala.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -43,7 +39,7 @@ jobs:
uses: actions/[email protected]
with:
distribution: temurin
java-version: 8
java-version: 17
cache: sbt
- name: Build and run unit tests
working-directory: ./pramen
Expand Down
17 changes: 8 additions & 9 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down Expand Up @@ -203,15 +203,15 @@ Pramen for Python transformers is available in PyPi: [![PyPI](https://badge.fury
Pramen is released as a set of thin JAR libraries. When running on a specific environment you might want to include all
dependencies in an uber jar that you can build for your Scala version. You can do that by either
- Downloading pre-compiled version of Pramen runners at the [Releases](https://github.com/AbsaOSS/pramen/releases) 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
Expand All @@ -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
Expand Down
2 changes: 1 addition & 1 deletion pramen/api/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@
<parent>
<groupId>za.co.absa.pramen</groupId>
<artifactId>pramen</artifactId>
<version>1.14.9-SNAPSHOT</version>
<version>1.15.0-SNAPSHOT</version>
</parent>

<properties>
Expand Down
52 changes: 38 additions & 14 deletions pramen/build.sbt
Original file line number Diff line number Diff line change
Expand Up @@ -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")

Expand Down Expand Up @@ -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
Expand All @@ -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,
Expand All @@ -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}")
Expand All @@ -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}")
Expand All @@ -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"),
Expand All @@ -146,22 +170,22 @@ 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}")
sparkVersion(scalaVersion.value)
},
(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
Expand Down
4 changes: 2 additions & 2 deletions pramen/core/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@
<parent>
<groupId>za.co.absa.pramen</groupId>
<artifactId>pramen</artifactId>
<version>1.14.9-SNAPSHOT</version>
<version>1.15.0-SNAPSHOT</version>
</parent>

<properties>
Expand Down Expand Up @@ -128,7 +128,7 @@
<artifactId>javax.mail</artifactId>
</dependency>

<!-- https://search.maven.org/artifact/com.github.yruslan/channel_scala_2.11 -->
<!-- https://search.maven.org/artifact/com.github.yruslan/channel_scala_2.12 -->
<dependency>
<groupId>com.lihaoyi</groupId>
<artifactId>requests_${scala.compat.version}</artifactId>
Expand Down
Original file line number Diff line number Diff line change
@@ -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]
}
}
}
Loading
Loading