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 envMap = (Map) theEnvironmentField.get(null); + envMap.clear(); - 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); + Class variableClass = null; + Class valueClass = null; + try { + variableClass = Class.forName("java.lang.ProcessEnvironment$Variable"); + valueClass = Class.forName("java.lang.ProcessEnvironment$Value"); + } catch (ClassNotFoundException e) { + // Windows: theEnvironment uses String keys. + } + if (variableClass != null && valueClass != null) { + Method valueOfVariable = variableClass.getDeclaredMethod("valueOf", String.class); + Method valueOfValue = valueClass.getDeclaredMethod("valueOf", String.class); + valueOfVariable.setAccessible(true); + valueOfValue.setAccessible(true); + for (Map.Entry entry : env.entrySet()) { + envMap.put( + valueOfVariable.invoke(null, entry.getKey()), + valueOfValue.invoke(null, entry.getValue())); + } + } else { + envMap.putAll(env); + } - modifiersField.setInt(envField, envField.getModifiers() & Modifier.FINAL); - modifiersField.setAccessible(false); + if (System.getProperty("os.name", "").toLowerCase(Locale.ROOT).startsWith("windows")) { + Field ciEnvField = processEnv.getDeclaredField("theCaseInsensitiveEnvironment"); + ciEnvField.setAccessible(true); + Map ciEnvMap = (Map) ciEnvField.get(null); + if (ciEnvMap != null) { + ciEnvMap.clear(); + ciEnvMap.putAll(env); + } + } } /** Prevents the CliFrontend from calling System.exit. */ diff --git a/runners/flink/src/test/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/UnboundedSourceWrapperTest.java b/runners/flink/src/test/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/UnboundedSourceWrapperTest.java index f57198e08e3e..0033243d9255 100644 --- a/runners/flink/src/test/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/UnboundedSourceWrapperTest.java +++ b/runners/flink/src/test/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/UnboundedSourceWrapperTest.java @@ -666,10 +666,7 @@ private static void testSourceDoesNotShutdown(boolean shouldHaveReaders) throws if (!shouldHaveReaders) { // The expected state is for finalizeSource to sleep instead of exiting while (true) { - StackTraceElement[] callStack = thread.getStackTrace(); - if (callStack.length >= 2 - && "sleep".equals(callStack[0].getMethodName()) - && "finalizeSource".equals(callStack[1].getMethodName())) { + if (isInFinalizeSource(thread.getStackTrace())) { break; } Thread.sleep(10); @@ -694,6 +691,15 @@ private static void testSourceDoesNotShutdown(boolean shouldHaveReaders) throws assertThat(thread.isAlive(), is(false)); } + private static boolean isInFinalizeSource(StackTraceElement[] callStack) { + for (StackTraceElement frame : callStack) { + if ("finalizeSource".equals(frame.getMethodName())) { + return true; + } + } + return false; + } + @Test public void testSequentialReadingFromBoundedSource() throws Exception { UnboundedReadFromBoundedSource.BoundedToUnboundedSourceAdapter source =