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
Original file line number Diff line number Diff line change
Expand Up @@ -649,6 +649,7 @@ class BeamModulePlugin implements Plugin<Project> {
def solace_version = "10.21.0"
def spark2_version = "2.4.8"
def spark3_version = "3.5.0"
def spark4_version = "4.0.2"
def spotbugs_version = "4.8.3"
def testcontainers_version = "1.21.4"
// [bomupgrader] determined by: org.apache.arrow:arrow-memory-core, consistent with: google_cloud_platform_libraries_bom
Expand All @@ -658,6 +659,7 @@ class BeamModulePlugin implements Plugin<Project> {

// Export Spark versions, so they are defined in a single place only
project.ext.spark3_version = spark3_version
project.ext.spark4_version = spark4_version
// version for BigQueryMetastore catalog (used by sdks:java:io:iceberg:bqms)
// TODO: remove this and download the jar normally when the catalog gets
// open-sourced (https://github.com/apache/iceberg/pull/11039)
Expand Down Expand Up @@ -820,6 +822,7 @@ class BeamModulePlugin implements Plugin<Project> {
jackson_datatype_jsr310 : "com.fasterxml.jackson.datatype:jackson-datatype-jsr310:$jackson_version",
jackson_module_scala_2_11 : "com.fasterxml.jackson.module:jackson-module-scala_2.11:$jackson_version",
jackson_module_scala_2_12 : "com.fasterxml.jackson.module:jackson-module-scala_2.12:$jackson_version",
jackson_module_scala_2_13 : "com.fasterxml.jackson.module:jackson-module-scala_2.13:$jackson_version",
jamm : 'com.github.jbellis:jamm:0.4.0',
jaxb_api : "jakarta.xml.bind:jakarta.xml.bind-api:$jaxb_api_version",
jaxb_impl : "com.sun.xml.bind:jaxb-impl:$jaxb_api_version",
Expand Down
3 changes: 3 additions & 0 deletions runners/spark/job-server/spark_job_server.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,10 @@ apply plugin: 'application'
// we need to set mainClassName before applying shadow plugin
mainClassName = "org.apache.beam.runners.spark.SparkJobServerDriver"

def sparkVersion = project.findProperty('spark_version') ?: ''

applyJavaNature(
requireJavaVersion: (sparkVersion.startsWith("4") ? org.gradle.api.JavaVersion.VERSION_17 : null),
automaticModuleName: 'org.apache.beam.runners.spark.jobserver',
archivesBaseName: project.hasProperty('archives_base_name') ? archives_base_name : archivesBaseName,
validateShadowJar: false,
Expand Down
22 changes: 17 additions & 5 deletions runners/spark/spark_runner.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -19,8 +19,20 @@
import groovy.json.JsonOutput

apply plugin: 'org.apache.beam.module'

// Numeric version comparison (lexicographic string compare was fragile — e.g. "3.10.0" < "3.5.0").
def isSparkAtLeast = { String minVersion ->
def parts = spark_version.tokenize('.-').findAll { it.isInteger() }*.toInteger()
def minParts = minVersion.tokenize('.-').findAll { it.isInteger() }*.toInteger()
for (int i = 0; i < Math.min(parts.size(), minParts.size()); i++) {
if (parts[i] != minParts[i]) return parts[i] > minParts[i]
}
return parts.size() >= minParts.size()
}

applyJavaNature(
enableStrictDependencies: true,
requireJavaVersion: (isSparkAtLeast("4.0.0") ? org.gradle.api.JavaVersion.VERSION_17 : null),
automaticModuleName: 'org.apache.beam.runners.spark',
archivesBaseName: (project.hasProperty('archives_base_name') ? archives_base_name : archivesBaseName),
exportJavadoc: (project.hasProperty('exportJavadoc') ? exportJavadoc : true),
Expand Down Expand Up @@ -240,7 +252,7 @@ dependencies {
spark.components.each { component ->
provided "$component:$spark_version"
}
if ("$spark_version" >= "3.5.0") {
if (isSparkAtLeast("3.5.0")) {
implementation "org.apache.spark:spark-common-utils_$spark_scala_version:$spark_version"
implementation "org.apache.spark:spark-sql-api_$spark_scala_version:$spark_version"
}
Expand Down Expand Up @@ -270,7 +282,7 @@ dependencies {
testImplementation library.java.mockito_core
testImplementation "org.assertj:assertj-core:3.11.1"
testImplementation "org.apache.zookeeper:zookeeper:3.4.11"
if ("$spark_version" >= "3.5.0") {
if (isSparkAtLeast("3.5.0")) {
testImplementation "org.apache.spark:spark-common-utils_$spark_scala_version:$spark_version"
testImplementation "org.apache.spark:spark-sql-api_$spark_scala_version:$spark_version"
}
Expand All @@ -284,7 +296,7 @@ dependencies {
"hadoopVersion$kv.key" "org.apache.hadoop:hadoop-common:$kv.value"
// Force paranamer 2.8 to avoid issues when using Scala 2.12
"hadoopVersion$kv.key" "com.thoughtworks.paranamer:paranamer:2.8"
if ("$spark_version" >= "3.5.0") {
if (isSparkAtLeast("3.5.0")) {
// Add log4j 2.x dependencies as Spark 3.5+ uses slf4j with log4j 2.x backend
"hadoopVersion$kv.key" library.java.log4j2_api
"hadoopVersion$kv.key" library.java.log4j2_core
Expand All @@ -310,7 +322,7 @@ configurations.validatesRunner {
// Exclude to make sure log4j binding is used
exclude group: "org.slf4j", module: "slf4j-simple"

if ("$spark_version" >= "3.5.0") {
if (isSparkAtLeast("3.5.0")) {
// Exclude log4j 1.x dependencies to prevent conflict with log4j 2.x used by spark 3.5+
exclude group: "log4j", module: "log4j"
}
Expand All @@ -321,7 +333,7 @@ hadoopVersions.each { kv ->
resolutionStrategy {
force "org.apache.hadoop:hadoop-common:$kv.value"
}
if ("$spark_version" >= "3.5.0") {
if (isSparkAtLeast("3.5.0")) {
// Exclude log4j 1.x dependencies to prevent conflict with log4j 2.x used by spark 3.5+
exclude group: "log4j", module: "log4j"
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import scala.Option;
import scala.collection.JavaConversions;
import scala.collection.JavaConverters;

/** Classes implementing Beam {@link Source} {@link RDD}s. */
@SuppressWarnings({
Expand All @@ -75,7 +75,7 @@ public static class Bounded<T> extends RDD<WindowedValue<T>> {

// to satisfy Scala API.
private static final scala.collection.immutable.Seq<Dependency<?>> NIL =
JavaConversions.asScalaBuffer(Collections.<Dependency<?>>emptyList()).toList();
JavaConverters.asScalaBuffer(Collections.<Dependency<?>>emptyList()).toList();

public Bounded(
SparkContext sc,
Expand Down Expand Up @@ -148,7 +148,7 @@ public scala.collection.Iterator<WindowedValue<T>> compute(
final Iterator<WindowedValue<T>> readerIterator =
new ReaderToIteratorAdapter<>(metricsContainer, reader);

return new InterruptibleIterator<>(context, JavaConversions.asScalaIterator(readerIterator));
return new InterruptibleIterator<>(context, JavaConverters.asScalaIterator(readerIterator));
}

/**
Expand Down Expand Up @@ -299,7 +299,7 @@ public static class Unbounded<T, CheckpointMarkT extends UnboundedSource.Checkpo

// to satisfy Scala API.
private static final scala.collection.immutable.List<Dependency<?>> NIL =
JavaConversions.asScalaBuffer(Collections.<Dependency<?>>emptyList()).toList();
JavaConverters.asScalaBuffer(Collections.<Dependency<?>>emptyList()).toList();

public Unbounded(
SparkContext sc,
Expand Down Expand Up @@ -344,7 +344,7 @@ public scala.collection.Iterator<scala.Tuple2<Source<T>, CheckpointMarkT>> compu
(CheckpointableSourcePartition<T, CheckpointMarkT>) split;
scala.Tuple2<Source<T>, CheckpointMarkT> tuple2 =
new scala.Tuple2<>(partition.getSource(), partition.checkpointMark);
return JavaConversions.asScalaIterator(Collections.singleton(tuple2).iterator());
return JavaConverters.asScalaIterator(Collections.singleton(tuple2).iterator());
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -186,7 +186,7 @@ public Duration slideDuration() {

@Override
public scala.collection.immutable.List<DStream<?>> dependencies() {
return scala.collection.JavaConversions.asScalaBuffer(
return scala.collection.JavaConverters.asScalaBuffer(
Collections.<DStream<?>>singletonList(parent))
.toList();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@
import scala.Tuple2;
import scala.Tuple3;
import scala.collection.Iterator;
import scala.collection.JavaConversions;
import scala.collection.JavaConverters;
import scala.collection.Seq;
import scala.runtime.AbstractFunction1;

Expand Down Expand Up @@ -238,7 +238,7 @@ private Collection<TimerInternals.TimerData> filterTimersEligibleForProcessing(
// new input for key.
try {
final Iterable<WindowedValue<InputT>> elements =
FluentIterable.from(JavaConversions.asJavaIterable(encodedElements))
FluentIterable.from(JavaConverters.asJavaIterable(encodedElements))
.transform(bytes -> CoderHelpers.fromByteArray(bytes, wvCoder));

LOG.trace("{}: input elements: {}", logPrefix, elements);
Expand Down Expand Up @@ -410,7 +410,7 @@ private Collection<TimerInternals.TimerData> filterTimersEligibleForProcessing(
droppedDueToClosedWindow.inc(-droppedDueToClosedWindow.getCumulative());
}

return scala.collection.JavaConversions.asScalaIterator(
return JavaConverters.asScalaIterator(
new UpdateStateByKeyOutputIterator(input, reduceFn, droppedDueToLateness));
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@
package org.apache.beam.runners.spark.structuredstreaming;

import static org.apache.beam.runners.core.metrics.MetricsContainerStepMap.asAttemptedOnlyMetricResults;
import static org.sparkproject.guava.base.Objects.firstNonNull;
import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects.firstNonNull;

import java.io.IOException;
import java.util.concurrent.ExecutionException;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
*/
package org.apache.beam.runners.spark.structuredstreaming.translation.batch;

import java.io.Serializable;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
Expand Down Expand Up @@ -49,7 +50,6 @@
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Maps;
import org.joda.time.Instant;
import scala.Serializable;

/**
* Factory to create a {@link DoFnRunner}. The factory supports fusing multiple {@link DoFnRunner
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -330,7 +330,8 @@ private static <T> void translateFlatten(
}
}
// Unify streams into a single stream.
unifiedStreams = context.getStreamingContext().union(JavaConverters.asScalaBuffer(dStreams));
unifiedStreams =
context.getStreamingContext().union(JavaConverters.asScalaBuffer(dStreams).toList());
}

context.pushDataset(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@

import java.io.Serializable;
import java.util.Collection;
import java.util.Collections;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
Expand Down Expand Up @@ -62,7 +63,6 @@
import org.checkerframework.checker.nullness.qual.Nullable;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.sparkproject.guava.collect.Iterators;
import scala.Option;
import scala.Tuple2;
import scala.runtime.AbstractFunction3;
Expand Down Expand Up @@ -236,7 +236,7 @@ public TimerInternals timerInternals() {
final byte[] byteValue = serializedValue.get();
@Nullable WindowedValue<ValueT> windowedValue;
@Nullable WindowedValue<KV<KeyT, ValueT>> keyedWindowedValue;
Iterator<WindowedValue<KV<KeyT, ValueT>>> iterator = Iterators.emptyIterator();
Iterator<WindowedValue<KV<KeyT, ValueT>>> iterator = Collections.emptyIterator();
if (byteValue.length > 0) {
windowedValue = CoderHelpers.fromByteArray(byteValue, this.wvCoder);
keyedWindowedValue = windowedValue.withValue(KV.of(key, windowedValue.getValue()));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -306,7 +306,7 @@ public void evaluate(Flatten.PCollections<T> transform, EvaluationContext contex
}
// start by unifying streams into a single stream.
JavaDStream<WindowedValue<T>> unifiedStreams =
context.getStreamingContext().union(JavaConverters.asScalaBuffer(dStreams));
context.getStreamingContext().union(JavaConverters.asScalaBuffer(dStreams).toList());
context.putDataset(transform, new UnboundedDataset<>(unifiedStreams, streamingSources));
}

Expand Down
Loading