Skip to content

Commit f20fca2

Browse files
authored
Add JVM --add-opens flags for Flink PVR on Java 21 (#39303)
1 parent ad6762a commit f20fca2

2 files changed

Lines changed: 25 additions & 2 deletions

File tree

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
{
22
"comment": "Modify this file in a trivial way to cause this test suite to run",
3-
"modification": 1,
3+
"modification": 2,
44
"https://github.com/apache/beam/pull/32440": "test new datastream runner for batch"
55
}

runners/flink/job-server/flink_job_server.gradle

Lines changed: 24 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -239,6 +239,27 @@ tasks.register("validatesPortableRunner") {
239239
dependsOn validatesPortableRunnerStreamingCheckpoint
240240
}
241241

242+
def flinkJobServerJvmArgs() {
243+
def testJavaVer = project.findProperty('testJavaVersion') ? (project.property('testJavaVersion') as int) : JavaVersion.current().majorVersion.toInteger()
244+
if (testJavaVer >= 17) {
245+
return [
246+
"--add-opens=java.base/sun.nio.ch=ALL-UNNAMED",
247+
"--add-opens=java.base/java.nio=ALL-UNNAMED",
248+
"--add-opens=java.base/java.util=ALL-UNNAMED",
249+
"--add-opens=java.base/java.lang.invoke=ALL-UNNAMED",
250+
"--add-opens=java.base/java.lang=ALL-UNNAMED",
251+
"-Djava.security.manager=allow",
252+
]
253+
}
254+
return []
255+
}
256+
257+
[validatesPortableRunnerDocker, validatesPortableRunnerBatchDataSet, validatesPortableRunnerBatch, validatesPortableRunnerStreaming, validatesPortableRunnerStreamingCheckpoint].each { task ->
258+
task.configure {
259+
jvmArgs += flinkJobServerJvmArgs()
260+
}
261+
}
262+
242263
def jobPort = BeamModulePlugin.getRandomPort()
243264
def artifactPort = BeamModulePlugin.getRandomPort()
244265

@@ -270,7 +291,9 @@ def setupTask = project.tasks.register("flinkJobServerSetup", Exec) {
270291
}
271292

272293
executable 'sh'
273-
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}\""
294+
def jvmArgs = flinkJobServerJvmArgs().join(' ')
295+
def jvmArgsOpt = jvmArgs ? "--jvm_args \"${jvmArgs}\"" : ""
296+
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}"
274297
}
275298

276299
def cleanupTask = project.tasks.register("flinkJobServerCleanup", Exec) {

0 commit comments

Comments
 (0)