Skip to content

Commit 96018bb

Browse files
committed
Fix Flink runner tests on Java 21
1 parent 477c747 commit 96018bb

6 files changed

Lines changed: 77 additions & 31 deletions

File tree

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
11
{
22
"comment": "Modify this file in a trivial way to cause this test suite to run.",
3-
"modification": 1
3+
"modification": 2
44
}

runners/flink/2.0/src/test/java/org/apache/beam/runners/flink/FlinkSubmissionTest.java

Lines changed: 17 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,6 @@
1919

2020
import java.io.File;
2121
import java.lang.reflect.Field;
22-
import java.lang.reflect.Modifier;
2322
import java.nio.charset.StandardCharsets;
2423
import java.nio.file.Files;
2524
import java.security.Permission;
@@ -211,20 +210,24 @@ private static void restoreEnvironment() throws Exception {
211210
* because Flink's CliFrontend requires a Flink configuration file for which the location can only
212211
* be set using the {@code ConfigConstants.ENV_FLINK_CONF_DIR} environment variable.
213212
*/
213+
@SuppressWarnings("unchecked")
214214
private static void modifyEnv(Map<String, String> env) throws Exception {
215-
Class processEnv = Class.forName("java.lang.ProcessEnvironment");
216-
Field envField = processEnv.getDeclaredField("theUnmodifiableEnvironment");
217-
218-
Field modifiersField = Field.class.getDeclaredField("modifiers");
219-
modifiersField.setAccessible(true);
220-
modifiersField.setInt(envField, envField.getModifiers() & ~Modifier.FINAL);
221-
222-
envField.setAccessible(true);
223-
envField.set(null, env);
224-
envField.setAccessible(false);
225-
226-
modifiersField.setInt(envField, envField.getModifiers() & Modifier.FINAL);
227-
modifiersField.setAccessible(false);
215+
Class<?> processEnvironment = Class.forName("java.lang.ProcessEnvironment");
216+
Field theEnvironmentField = processEnvironment.getDeclaredField("theEnvironment");
217+
theEnvironmentField.setAccessible(true);
218+
Map<String, String> envMap = (Map<String, String>) theEnvironmentField.get(null);
219+
envMap.clear();
220+
envMap.putAll(env);
221+
222+
Field theCaseInsensitiveEnvironmentField =
223+
processEnvironment.getDeclaredField("theCaseInsensitiveEnvironment");
224+
theCaseInsensitiveEnvironmentField.setAccessible(true);
225+
Map<String, String> ciEnvMap =
226+
(Map<String, String>) theCaseInsensitiveEnvironmentField.get(null);
227+
if (ciEnvMap != null) {
228+
ciEnvMap.clear();
229+
ciEnvMap.putAll(env);
230+
}
228231
}
229232

230233
/** Prevents the CliFrontend from calling System.exit. */

runners/flink/flink_runner.gradle

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -180,13 +180,29 @@ if (use_override) {
180180
}
181181
}
182182

183+
def flinkTestJvmArgs() {
184+
def testJavaVer = project.findProperty('testJavaVersion') ? (project.property('testJavaVersion') as int) : JavaVersion.current().majorVersion.toInteger()
185+
if (testJavaVer >= 17) {
186+
return [
187+
"--add-opens=java.base/sun.nio.ch=ALL-UNNAMED",
188+
"--add-opens=java.base/java.nio=ALL-UNNAMED",
189+
"--add-opens=java.base/java.util=ALL-UNNAMED",
190+
"--add-opens=java.base/java.lang.invoke=ALL-UNNAMED",
191+
"--add-opens=java.base/java.lang=ALL-UNNAMED",
192+
]
193+
} else {
194+
return []
195+
}
196+
}
197+
183198
test {
184199
systemProperty "log4j.configuration", "log4j-test.properties"
185200
// Change log level to debug:
186201
// systemProperty "org.slf4j.simpleLogger.defaultLogLevel", "debug"
187202
// Change log level to debug only for the package and nested packages:
188203
// systemProperty "org.slf4j.simpleLogger.log.org.apache.beam.runners.flink.translation.wrappers.streaming", "debug"
189204
jvmArgs "-XX:-UseGCOverheadLimit"
205+
jvmArgs += flinkTestJvmArgs()
190206
if (System.getProperty("beamSurefireArgline")) {
191207
jvmArgs System.getProperty("beamSurefireArgline")
192208
}
@@ -315,6 +331,7 @@ def createValidatesRunnerTask(Map m) {
315331
group = "Verification"
316332
// Disable gradle cache
317333
outputs.upToDateWhen { false }
334+
jvmArgs += flinkTestJvmArgs()
318335
def runnerType = config.streaming ? "streaming" : "batch"
319336
description = "Validates the ${runnerType} runner"
320337
def pipelineOptionsArray = ["--runner=TestFlinkRunner",
@@ -421,6 +438,7 @@ tasks.register("validatesRunnerSickbay", Test) {
421438
systemProperty "beamTestPipelineOptions", JsonOutput.toJson([
422439
"--runner=TestFlinkRunner",
423440
])
441+
jvmArgs += flinkTestJvmArgs()
424442

425443
classpath = configurations.validatesRunner
426444
testClassesDirs = files(project(":sdks:java:core").sourceSets.test.output.classesDirs)
@@ -443,6 +461,7 @@ tasks.register("examplesIntegrationTest", Test) {
443461
group = "Verification"
444462
// Disable gradle cache
445463
outputs.upToDateWhen { false }
464+
jvmArgs += flinkTestJvmArgs()
446465
def pipelineOptionsArray = ["--runner=TestFlinkRunner",
447466
"--parallelism=2",
448467
"--tempLocation=${tempLocation}",

runners/flink/job-server/flink_job_server.gradle

Lines changed: 16 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -242,12 +242,27 @@ tasks.register("validatesPortableRunner") {
242242
def jobPort = BeamModulePlugin.getRandomPort()
243243
def artifactPort = BeamModulePlugin.getRandomPort()
244244

245+
def flinkJobServerJvmArgs() {
246+
def testJavaVer = project.findProperty('testJavaVersion') ? (project.property('testJavaVersion') as int) : JavaVersion.current().majorVersion.toInteger()
247+
if (testJavaVer >= 17) {
248+
return [
249+
"--add-opens=java.base/sun.nio.ch=ALL-UNNAMED",
250+
"--add-opens=java.base/java.nio=ALL-UNNAMED",
251+
"--add-opens=java.base/java.util=ALL-UNNAMED",
252+
"--add-opens=java.base/java.lang.invoke=ALL-UNNAMED"
253+
]
254+
}
255+
return []
256+
}
257+
245258
def setupTask = project.tasks.register("flinkJobServerSetup", Exec) {
246259
dependsOn shadowJar
247260
def pythonDir = project.project(":sdks:python").projectDir
248261
def flinkJobServerJar = shadowJar.archivePath
249262
def flinkDir = project.project(":runners:flink").projectDir
250263
def additionalArgs = ""
264+
def jvmArgs = flinkJobServerJvmArgs().join(' ')
265+
def jvmArgsOpt = jvmArgs ? "--jvm_args \"${jvmArgs}\"" : ""
251266

252267
if (project.hasProperty('flinkConfDir')) {
253268
additionalArgs += " --flink-conf-dir=${project.property('flinkConfDir')}"
@@ -270,7 +285,7 @@ def setupTask = project.tasks.register("flinkJobServerSetup", Exec) {
270285
}
271286

272287
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}\""
288+
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} ${jvmArgsOpt} --additional_args \"${additionalArgs}\""
274289
}
275290

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

runners/flink/src/test/java/org/apache/beam/runners/flink/FlinkSubmissionTest.java

Lines changed: 17 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,6 @@
1919

2020
import java.io.File;
2121
import java.lang.reflect.Field;
22-
import java.lang.reflect.Modifier;
2322
import java.nio.charset.StandardCharsets;
2423
import java.nio.file.Files;
2524
import java.security.Permission;
@@ -231,20 +230,24 @@ private static void restoreEnvironment() throws Exception {
231230
* because Flink's CliFrontend requires a Flink configuration file for which the location can only
232231
* be set using the {@code ConfigConstants.ENV_FLINK_CONF_DIR} environment variable.
233232
*/
233+
@SuppressWarnings("unchecked")
234234
private static void modifyEnv(Map<String, String> env) throws Exception {
235-
Class processEnv = Class.forName("java.lang.ProcessEnvironment");
236-
Field envField = processEnv.getDeclaredField("theUnmodifiableEnvironment");
237-
238-
Field modifiersField = Field.class.getDeclaredField("modifiers");
239-
modifiersField.setAccessible(true);
240-
modifiersField.setInt(envField, envField.getModifiers() & ~Modifier.FINAL);
241-
242-
envField.setAccessible(true);
243-
envField.set(null, env);
244-
envField.setAccessible(false);
245-
246-
modifiersField.setInt(envField, envField.getModifiers() & Modifier.FINAL);
247-
modifiersField.setAccessible(false);
235+
Class<?> processEnvironment = Class.forName("java.lang.ProcessEnvironment");
236+
Field theEnvironmentField = processEnvironment.getDeclaredField("theEnvironment");
237+
theEnvironmentField.setAccessible(true);
238+
Map<String, String> envMap = (Map<String, String>) theEnvironmentField.get(null);
239+
envMap.clear();
240+
envMap.putAll(env);
241+
242+
Field theCaseInsensitiveEnvironmentField =
243+
processEnvironment.getDeclaredField("theCaseInsensitiveEnvironment");
244+
theCaseInsensitiveEnvironmentField.setAccessible(true);
245+
Map<String, String> ciEnvMap =
246+
(Map<String, String>) theCaseInsensitiveEnvironmentField.get(null);
247+
if (ciEnvMap != null) {
248+
ciEnvMap.clear();
249+
ciEnvMap.putAll(env);
250+
}
248251
}
249252

250253
/** Prevents the CliFrontend from calling System.exit. */

sdks/python/scripts/run_job_server.sh

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ Options:
2323
--job_port [port for job endpoint, default 8099]
2424
--artifact_port [port for artifact service, default 8098]
2525
--job_server_jar [path to job server jar]
26+
--jvm_args [additional JVM arguments, e.g. --add-opens flags]
2627
END
2728

2829
JOB_PORT=8099
@@ -56,6 +57,11 @@ while [[ $# -gt 0 ]]; do
5657
shift
5758
shift
5859
;;
60+
--jvm_args)
61+
JVM_ARGS="$2"
62+
shift
63+
shift
64+
;;
5965
--additional_args)
6066
ADDITIONAL_ARGS="$2"
6167
shift
@@ -107,7 +113,7 @@ case $STARTSTOP in
107113
fi
108114

109115
echo "Launching job server @ $JOB_PORT ..."
110-
"$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 </dev/null &
116+
"$JAVA_CMD" $JVM_ARGS -jar $JOB_SERVER_JAR --job-port=$JOB_PORT --artifact-port=$ARTIFACT_PORT --expansion-port=0 $ADDITIONAL_ARGS >$TEMP_DIR/$FILE_BASE.log 2>&1 </dev/null &
111117
mypid=$!
112118
if kill -0 $mypid >/dev/null 2>&1; then
113119
echo $mypid >> $pid

0 commit comments

Comments
 (0)