From 993e0d0d232b3764b5a33ecf28c3f3a8ac546f52 Mon Sep 17 00:00:00 2001 From: aIbrahiim Date: Thu, 9 Jul 2026 15:34:05 +0300 Subject: [PATCH] Fix Flink runner tests on Java 21 --- .../trigger_files/beam_PreCommit_Java.json | 2 +- .../runners/flink/FlinkSubmissionTest.java | 33 +++++++++++-------- runners/flink/flink_runner.gradle | 16 +++++++++ .../runners/flink/FlinkSubmissionTest.java | 33 +++++++++++-------- 4 files changed, 57 insertions(+), 27 deletions(-) diff --git a/.github/trigger_files/beam_PreCommit_Java.json b/.github/trigger_files/beam_PreCommit_Java.json index 5abe02fc09c7..3a009261f4f9 100644 --- a/.github/trigger_files/beam_PreCommit_Java.json +++ b/.github/trigger_files/beam_PreCommit_Java.json @@ -1,4 +1,4 @@ { "comment": "Modify this file in a trivial way to cause this test suite to run.", - "modification": 1 + "modification": 2 } diff --git a/runners/flink/2.0/src/test/java/org/apache/beam/runners/flink/FlinkSubmissionTest.java b/runners/flink/2.0/src/test/java/org/apache/beam/runners/flink/FlinkSubmissionTest.java index 66079f855a77..b4e1b3f4649b 100644 --- a/runners/flink/2.0/src/test/java/org/apache/beam/runners/flink/FlinkSubmissionTest.java +++ b/runners/flink/2.0/src/test/java/org/apache/beam/runners/flink/FlinkSubmissionTest.java @@ -19,7 +19,6 @@ import java.io.File; import java.lang.reflect.Field; -import java.lang.reflect.Modifier; import java.nio.charset.StandardCharsets; import java.nio.file.Files; import java.security.Permission; @@ -211,20 +210,28 @@ private static void restoreEnvironment() throws Exception { * because Flink's CliFrontend requires a Flink configuration file for which the location can only * be set using the {@code ConfigConstants.ENV_FLINK_CONF_DIR} environment variable. */ + @SuppressWarnings("unchecked") private static void modifyEnv(Map env) throws Exception { - Class processEnv = Class.forName("java.lang.ProcessEnvironment"); - Field envField = processEnv.getDeclaredField("theUnmodifiableEnvironment"); + Class processEnvironment = Class.forName("java.lang.ProcessEnvironment"); + Field theEnvironmentField = processEnvironment.getDeclaredField("theEnvironment"); + theEnvironmentField.setAccessible(true); + Map envMap = (Map) theEnvironmentField.get(null); + envMap.clear(); + envMap.putAll(env); - Field modifiersField = Field.class.getDeclaredField("modifiers"); - modifiersField.setAccessible(true); - modifiersField.setInt(envField, envField.getModifiers() & ~Modifier.FINAL); - - envField.setAccessible(true); - envField.set(null, env); - envField.setAccessible(false); - - modifiersField.setInt(envField, envField.getModifiers() & Modifier.FINAL); - modifiersField.setAccessible(false); + try { + Field theCaseInsensitiveEnvironmentField = + processEnvironment.getDeclaredField("theCaseInsensitiveEnvironment"); + theCaseInsensitiveEnvironmentField.setAccessible(true); + Map ciEnvMap = + (Map) theCaseInsensitiveEnvironmentField.get(null); + if (ciEnvMap != null) { + ciEnvMap.clear(); + ciEnvMap.putAll(env); + } + } catch (NoSuchFieldException e) { + // theCaseInsensitiveEnvironment is not present on all JDK platforms (e.g. Linux). + } } /** Prevents the CliFrontend from calling System.exit. */ diff --git a/runners/flink/flink_runner.gradle b/runners/flink/flink_runner.gradle index 837561ec71b7..ed9ce8b33d1b 100644 --- a/runners/flink/flink_runner.gradle +++ b/runners/flink/flink_runner.gradle @@ -180,6 +180,21 @@ if (use_override) { } } +def flinkTestJvmArgs() { + 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", + ] + } else { + return [] + } +} + test { systemProperty "log4j.configuration", "log4j-test.properties" // Change log level to debug: @@ -187,6 +202,7 @@ test { // Change log level to debug only for the package and nested packages: // systemProperty "org.slf4j.simpleLogger.log.org.apache.beam.runners.flink.translation.wrappers.streaming", "debug" jvmArgs "-XX:-UseGCOverheadLimit" + jvmArgs += flinkTestJvmArgs() if (System.getProperty("beamSurefireArgline")) { jvmArgs System.getProperty("beamSurefireArgline") } diff --git a/runners/flink/src/test/java/org/apache/beam/runners/flink/FlinkSubmissionTest.java b/runners/flink/src/test/java/org/apache/beam/runners/flink/FlinkSubmissionTest.java index 8e4c3255fac5..dc305a6984ef 100644 --- a/runners/flink/src/test/java/org/apache/beam/runners/flink/FlinkSubmissionTest.java +++ b/runners/flink/src/test/java/org/apache/beam/runners/flink/FlinkSubmissionTest.java @@ -19,7 +19,6 @@ import java.io.File; import java.lang.reflect.Field; -import java.lang.reflect.Modifier; import java.nio.charset.StandardCharsets; import java.nio.file.Files; import java.security.Permission; @@ -231,20 +230,28 @@ private static void restoreEnvironment() throws Exception { * because Flink's CliFrontend requires a Flink configuration file for which the location can only * be set using the {@code ConfigConstants.ENV_FLINK_CONF_DIR} environment variable. */ + @SuppressWarnings("unchecked") private static void modifyEnv(Map env) throws Exception { - Class processEnv = Class.forName("java.lang.ProcessEnvironment"); - Field envField = processEnv.getDeclaredField("theUnmodifiableEnvironment"); + Class processEnvironment = Class.forName("java.lang.ProcessEnvironment"); + Field theEnvironmentField = processEnvironment.getDeclaredField("theEnvironment"); + theEnvironmentField.setAccessible(true); + Map envMap = (Map) theEnvironmentField.get(null); + envMap.clear(); + envMap.putAll(env); - Field modifiersField = Field.class.getDeclaredField("modifiers"); - modifiersField.setAccessible(true); - modifiersField.setInt(envField, envField.getModifiers() & ~Modifier.FINAL); - - envField.setAccessible(true); - envField.set(null, env); - envField.setAccessible(false); - - modifiersField.setInt(envField, envField.getModifiers() & Modifier.FINAL); - modifiersField.setAccessible(false); + try { + Field theCaseInsensitiveEnvironmentField = + processEnvironment.getDeclaredField("theCaseInsensitiveEnvironment"); + theCaseInsensitiveEnvironmentField.setAccessible(true); + Map ciEnvMap = + (Map) theCaseInsensitiveEnvironmentField.get(null); + if (ciEnvMap != null) { + ciEnvMap.clear(); + ciEnvMap.putAll(env); + } + } catch (NoSuchFieldException e) { + // theCaseInsensitiveEnvironment is not present on all JDK platforms (e.g. Linux). + } } /** Prevents the CliFrontend from calling System.exit. */