diff --git a/.github/trigger_files/beam_PreCommit_Java.json b/.github/trigger_files/beam_PreCommit_Java.json
index 5abe02fc09c7..0e9f1cacf9bc 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": 9
}
diff --git a/runners/flink/flink_runner.gradle b/runners/flink/flink_runner.gradle
index 837561ec71b7..165d5c2d93d7 100644
--- a/runners/flink/flink_runner.gradle
+++ b/runners/flink/flink_runner.gradle
@@ -180,6 +180,22 @@ 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",
+ "-Djava.security.manager=allow",
+ ]
+ } else {
+ return []
+ }
+}
+
test {
systemProperty "log4j.configuration", "log4j-test.properties"
// Change log level to debug:
@@ -187,6 +203,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..8e5ef8a3445a 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,11 +19,12 @@
import java.io.File;
import java.lang.reflect.Field;
-import java.lang.reflect.Modifier;
+import java.lang.reflect.Method;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.security.Permission;
import java.util.Collection;
+import java.util.Locale;
import java.util.Map;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
@@ -230,21 +231,49 @@ private static void restoreEnvironment() throws Exception {
* We modify the JVM's environment variables here. This is necessary for the end-to-end test
* 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.
+ *
+ *
On Unix, {@code theEnvironment} uses {@code Variable}/{@code Value} keys; {@code putAll}
+ * with {@code String} keys does not work for {@code System.getenv(String)} lookups.
*/
+ @SuppressWarnings("unchecked")
private static void modifyEnv(Map env) throws Exception {
Class processEnv = Class.forName("java.lang.ProcessEnvironment");
- Field envField = processEnv.getDeclaredField("theUnmodifiableEnvironment");
+ Field theEnvironmentField = processEnv.getDeclaredField("theEnvironment");
+ theEnvironmentField.setAccessible(true);
+ Map