From 50d4fe1fc1455f321b43d0780a0b3120437e3e1c Mon Sep 17 00:00:00 2001 From: Martijn Visser <2989614+MartijnVisser@users.noreply.github.com> Date: Tue, 16 Jun 2026 11:30:28 -0700 Subject: [PATCH] [FLINK-40070][tests] Fix flaky DynamicParameterITCase reading rolled JobManager logs The distribution log4j configuration rolls the log file on startup, so the JobManager startup banner frequently lands in a rolled .log.N file that FlinkDistribution.searchAllLogs skips; the test then either spins unboundedly waiting for the banner (multi-hour e2e_4 hang) or parses a half-written arguments block ("Missing required option: c"). Search rolled logs for the startup banner and bound the wait so a missing banner fails fast. Generated-by: Claude Opus 4.8 (1M context) (cherry picked from commit f9e0756fbfd15c633ae416c2e76f634da46101e8) --- .../tests/util/flink/FlinkDistribution.java | 17 +++++- .../flink/dist/DynamicParameterITCase.java | 54 +++++++++++++++---- 2 files changed, 60 insertions(+), 11 deletions(-) diff --git a/flink-end-to-end-tests/flink-end-to-end-tests-common/src/main/java/org/apache/flink/tests/util/flink/FlinkDistribution.java b/flink-end-to-end-tests/flink-end-to-end-tests-common/src/main/java/org/apache/flink/tests/util/flink/FlinkDistribution.java index 404b1ed455593a..f4192a53d01036 100644 --- a/flink-end-to-end-tests/flink-end-to-end-tests-common/src/main/java/org/apache/flink/tests/util/flink/FlinkDistribution.java +++ b/flink-end-to-end-tests/flink-end-to-end-tests-common/src/main/java/org/apache/flink/tests/util/flink/FlinkDistribution.java @@ -435,13 +435,28 @@ public void setTaskExecutorHosts(Collection taskExecutorHosts) throws IO public Stream searchAllLogs(Pattern pattern, Function matchProcessor) throws IOException { + return searchAllLogs(pattern, matchProcessor, false); + } + + /** + * Searches the distribution's log files for lines matching the given pattern. + * + * @param includeRolledLogs whether to also search rolled log files ({@code .log.N}); the + * distribution rolls logs on startup, which can move the startup banner out of the live + * {@code .log} file + */ + public Stream searchAllLogs( + Pattern pattern, Function matchProcessor, boolean includeRolledLogs) + throws IOException { final List matches = new ArrayList<>(2); try (Stream logFilesStream = Files.list(log)) { final Iterator logFiles = logFilesStream.iterator(); while (logFiles.hasNext()) { final Path logFile = logFiles.next(); - if (!logFile.getFileName().toString().endsWith(".log")) { + final String fileName = logFile.getFileName().toString(); + final boolean isRolledLog = fileName.matches(".*\\.log\\.\\d+"); + if (!fileName.endsWith(".log") && !(includeRolledLogs && isRolledLog)) { // ignore logs for previous runs that have a number suffix continue; } diff --git a/flink-end-to-end-tests/flink-end-to-end-tests-common/src/test/java/org/apache/flink/dist/DynamicParameterITCase.java b/flink-end-to-end-tests/flink-end-to-end-tests-common/src/test/java/org/apache/flink/dist/DynamicParameterITCase.java index 88ffb1973fe57e..49a4f081248331 100644 --- a/flink-end-to-end-tests/flink-end-to-end-tests-common/src/test/java/org/apache/flink/dist/DynamicParameterITCase.java +++ b/flink-end-to-end-tests/flink-end-to-end-tests-common/src/test/java/org/apache/flink/dist/DynamicParameterITCase.java @@ -22,6 +22,7 @@ import org.apache.flink.configuration.GlobalConfiguration; import org.apache.flink.configuration.JobManagerOptions; import org.apache.flink.configuration.RestOptions; +import org.apache.flink.core.testutils.CommonTestUtils; import org.apache.flink.runtime.entrypoint.ClusterEntrypoint; import org.apache.flink.runtime.entrypoint.EntrypointClusterConfiguration; import org.apache.flink.runtime.entrypoint.FlinkParseException; @@ -36,7 +37,9 @@ import org.junit.jupiter.api.io.TempDir; import java.io.IOException; +import java.io.UncheckedIOException; import java.nio.file.Path; +import java.time.Duration; import java.util.ArrayList; import java.util.List; import java.util.function.Predicate; @@ -52,6 +55,8 @@ class DynamicParameterITCase { private static final Pattern ENTRYPOINT_CLASSPATH_LOG_PATTERN = Pattern.compile(".*ClusterEntrypoint +\\[] - +Classpath:.*"); + private static final Duration LOG_TIMEOUT = Duration.ofMinutes(1); + private static final String HOST = "test_host"; private static final int PORT = 8082; @@ -129,16 +134,14 @@ private static void assertParameterPassing( dist.callJobManagerScript(args.toArray(new String[0])); - while (!allProgramArgumentsLogged(dist)) { - Thread.sleep(500); - } - - try (Stream lines = - dist.searchAllLogs(ENTRYPOINT_LOG_PATTERN, matcher -> matcher.group(1))) { + // The backgrounded JobManager writes its log asynchronously; wait until the complete + // program arguments block has been flushed, otherwise the parser may observe a + // partially-written block or never observe it at all. + final String[] programArguments = waitForProgramArguments(dist); + try { final EntrypointClusterConfiguration entrypointConfig = - ClusterEntrypoint.parseArguments( - lines.filter(new ProgramArgumentsFilter()).toArray(String[]::new)); + ClusterEntrypoint.parseArguments(programArguments); final Configuration configuration = loadConfiguration(entrypointConfig); @@ -158,6 +161,34 @@ private static void assertParameterPassing( } } + private static String[] waitForProgramArguments(FlinkDistribution dist) throws Exception { + final List captured = new ArrayList<>(1); + CommonTestUtils.waitUtil( + () -> { + // The "Classpath:" line is logged after the program arguments, so its presence + // means the complete arguments block has been flushed. + if (!allProgramArgumentsLogged(dist)) { + return false; + } + captured.clear(); + captured.add(readProgramArguments(dist)); + return true; + }, + LOG_TIMEOUT, + Duration.ofMillis(500), + "The JobManager did not log its complete program arguments in time."); + return captured.get(0); + } + + private static String[] readProgramArguments(FlinkDistribution dist) { + try (Stream lines = + dist.searchAllLogs(ENTRYPOINT_LOG_PATTERN, matcher -> matcher.group(1), true)) { + return lines.filter(new ProgramArgumentsFilter()).toArray(String[]::new); + } catch (IOException e) { + throw new UncheckedIOException(e); + } + } + private static Configuration loadConfiguration( EntrypointClusterConfiguration entrypointClusterConfiguration) { final Configuration dynamicProperties = @@ -170,11 +201,14 @@ private static Configuration loadConfiguration( return configuration; } - private static boolean allProgramArgumentsLogged(FlinkDistribution dist) throws IOException { + private static boolean allProgramArgumentsLogged(FlinkDistribution dist) { // the classpath is logged after the program arguments try (Stream lines = - dist.searchAllLogs(ENTRYPOINT_CLASSPATH_LOG_PATTERN, matcher -> matcher.group(0))) { + dist.searchAllLogs( + ENTRYPOINT_CLASSPATH_LOG_PATTERN, matcher -> matcher.group(0), true)) { return lines.iterator().hasNext(); + } catch (IOException e) { + throw new UncheckedIOException(e); } }