diff --git a/.github/trigger_files/beam_PostCommit_Java_PVR_Flink_Streaming.json b/.github/trigger_files/beam_PostCommit_Java_PVR_Flink_Streaming.json index 2dd3a2471d89..d7f4c03aaf5b 100644 --- a/.github/trigger_files/beam_PostCommit_Java_PVR_Flink_Streaming.json +++ b/.github/trigger_files/beam_PostCommit_Java_PVR_Flink_Streaming.json @@ -1,5 +1,5 @@ { "comment": "Modify this file in a trivial way to cause this test suite to run", - "modification": 1, + "modification": 2, "https://github.com/apache/beam/pull/32440": "test new datastream runner for batch" } diff --git a/runners/flink/job-server/flink_job_server.gradle b/runners/flink/job-server/flink_job_server.gradle index a6abca5e8586..46b1b5a607fe 100644 --- a/runners/flink/job-server/flink_job_server.gradle +++ b/runners/flink/job-server/flink_job_server.gradle @@ -239,6 +239,27 @@ tasks.register("validatesPortableRunner") { dependsOn validatesPortableRunnerStreamingCheckpoint } +def flinkJobServerJvmArgs() { + 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", + "--add-opens=java.base/java.lang=ALL-UNNAMED", + "-Djava.security.manager=allow", + ] + } + return [] +} + +[validatesPortableRunnerDocker, validatesPortableRunnerBatchDataSet, validatesPortableRunnerBatch, validatesPortableRunnerStreaming, validatesPortableRunnerStreamingCheckpoint].each { task -> + task.configure { + jvmArgs += flinkJobServerJvmArgs() + } +} + def jobPort = BeamModulePlugin.getRandomPort() def artifactPort = BeamModulePlugin.getRandomPort() @@ -270,7 +291,9 @@ def setupTask = project.tasks.register("flinkJobServerSetup", Exec) { } 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 ${flinkJobServerJar} --additional_args \"${additionalArgs}\"" + def jvmArgs = flinkJobServerJvmArgs().join(' ') + def jvmArgsOpt = jvmArgs ? "--jvm_args \"${jvmArgs}\"" : "" + 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 ${flinkJobServerJar} --additional_args \"${additionalArgs}\" ${jvmArgsOpt}" } def cleanupTask = project.tasks.register("flinkJobServerCleanup", Exec) {