diff --git a/.github/trigger_files/beam_PostCommit_XVR_Spark3.json b/.github/trigger_files/beam_PostCommit_XVR_Spark3.json index 03bb52a7ef46..1be4456862dc 100644 --- a/.github/trigger_files/beam_PostCommit_XVR_Spark3.json +++ b/.github/trigger_files/beam_PostCommit_XVR_Spark3.json @@ -1,3 +1,3 @@ { - "trigger-2026-04-04": "portable_runner expand_sdf opt-in" + "trigger-2026-07-08": "portable_runner expand_sdf opt-in 2" } diff --git a/runners/spark/job-server/spark_job_server.gradle b/runners/spark/job-server/spark_job_server.gradle index b2180ebf70af..e72b59444de8 100644 --- a/runners/spark/job-server/spark_job_server.gradle +++ b/runners/spark/job-server/spark_job_server.gradle @@ -267,13 +267,28 @@ tasks.register("validatesRunnerSickbay", Test) { def jobPort = BeamModulePlugin.getRandomPort() def artifactPort = BeamModulePlugin.getRandomPort() +def sparkJobServerJvmArgs() { + def testJavaVer = project.findProperty('testJavaVersion') ? (project.property('testJavaVersion') as int) : JavaVersion.current().majorVersion.toInteger() + if (testJavaVer >= 17) { + return [ + "--add-opens=java.base/sun.nio.ch=ALL-UNNAMED", + "--add-opens=java.base/java.nio=ALL-UNNAMED", + "--add-opens=java.base/java.util=ALL-UNNAMED", + "--add-opens=java.base/java.lang.invoke=ALL-UNNAMED" + ] + } + return [] +} + def setupTask = project.tasks.register("sparkJobServerSetup", Exec) { dependsOn shadowJar def pythonDir = project.project(":sdks:python").projectDir def sparkJobServerJar = shadowJar.archivePath + def jvmArgs = sparkJobServerJvmArgs().join(' ') + def jvmArgsOpt = jvmArgs ? "--jvm_args \"${jvmArgs}\"" : "" executable 'sh' - args '-c', "$pythonDir/scripts/run_job_server.sh stop --group_id ${project.name} && $pythonDir/scripts/run_job_server.sh start --group_id ${project.name} --job_port ${jobPort} --artifact_port ${artifactPort} --job_server_jar ${sparkJobServerJar}" + args '-c', "$pythonDir/scripts/run_job_server.sh stop --group_id ${project.name} && $pythonDir/scripts/run_job_server.sh start --group_id ${project.name} --job_port ${jobPort} --artifact_port ${artifactPort} --job_server_jar ${sparkJobServerJar} ${jvmArgsOpt}" } def cleanupTask = project.tasks.register("sparkJobServerCleanup", Exec) { diff --git a/sdks/python/scripts/run_job_server.sh b/sdks/python/scripts/run_job_server.sh index ce2ed9d08bc6..bd8ecd8c81a1 100755 --- a/sdks/python/scripts/run_job_server.sh +++ b/sdks/python/scripts/run_job_server.sh @@ -23,6 +23,7 @@ Options: --job_port [port for job endpoint, default 8099] --artifact_port [port for artifact service, default 8098] --job_server_jar [path to job server jar] + --jvm_args [additional JVM arguments, e.g. --add-opens flags] END JOB_PORT=8099 @@ -61,6 +62,11 @@ while [[ $# -gt 0 ]]; do shift shift ;; + --jvm_args) + JVM_ARGS="$2" + shift + shift + ;; start) STARTSTOP="$1" shift @@ -107,7 +113,7 @@ case $STARTSTOP in fi echo "Launching job server @ $JOB_PORT ..." - "$JAVA_CMD" -jar $JOB_SERVER_JAR --job-port=$JOB_PORT --artifact-port=$ARTIFACT_PORT --expansion-port=0 $ADDITIONAL_ARGS >$TEMP_DIR/$FILE_BASE.log 2>&1 $TEMP_DIR/$FILE_BASE.log 2>&1 /dev/null 2>&1; then echo $mypid >> $pid