diff --git a/paimon-e2e-tests/src/test/java/org/apache/paimon/tests/SparkE2eTest.java b/paimon-e2e-tests/src/test/java/org/apache/paimon/tests/SparkE2eTest.java index 2f87456d7db3..e78b82e6071b 100644 --- a/paimon-e2e-tests/src/test/java/org/apache/paimon/tests/SparkE2eTest.java +++ b/paimon-e2e-tests/src/test/java/org/apache/paimon/tests/SparkE2eTest.java @@ -36,6 +36,8 @@ public class SparkE2eTest extends E2eReaderTestBase { private static final Logger LOG = LoggerFactory.getLogger(SparkE2eTest.class); + private static final String COARSE_GRAINED_SCHEDULER_SHUTDOWN_ERROR = + "ERROR Utils: Uncaught exception in thread dispatcher-CoarseGrainedScheduler"; public SparkE2eTest() { super(false, false, true); @@ -75,13 +77,25 @@ public void testFlinkWriteAndSparkRead() throws Exception { LOG.info(execResult.getStderr()); throw new AssertionError("Failed when running spark sql."); } - return Arrays.stream(execResult.getStdout().split("\n")) + String stdout = + stripCoarseGrainedSchedulerShutdownError(execResult.getStdout()); + return Arrays.stream(stdout.split("\n")) .filter(s -> !s.contains("WARN")) .collect(Collectors.joining("\n")) + "\n"; }); } + private static String stripCoarseGrainedSchedulerShutdownError(String stdout) { + int errorIndex = stdout.indexOf(COARSE_GRAINED_SCHEDULER_SHUTDOWN_ERROR); + if (errorIndex < 0) { + return stdout; + } + + int errorLineStart = stdout.lastIndexOf('\n', errorIndex); + return errorLineStart < 0 ? "" : stdout.substring(0, errorLineStart); + } + private ContainerState getSpark() { return environment.getContainerByServiceName("spark-master-1").get(); }