This is an automated email from the ASF dual-hosted git repository. snuyanzin pushed a commit to branch release-1.20 in repository https://gitbox.apache.org/repos/asf/flink.git
commit 4df0326b784c7fb609fd6219dc6d6e9b9700ed41 Author: Sergey Nuyanzin <[email protected]> AuthorDate: Tue Aug 4 09:06:33 2026 +0200 [FLINK-40307][table] `RestoreTestCompletenessTest` was never executed in CI --- ...eness.java => RestoreTestCompletenessTest.java} | 86 ++++++++++++++-------- 1 file changed, 54 insertions(+), 32 deletions(-) diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/testutils/RestoreTestCompleteness.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/testutils/RestoreTestCompletenessTest.java similarity index 64% rename from flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/testutils/RestoreTestCompleteness.java rename to flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/testutils/RestoreTestCompletenessTest.java index ca2d28b7166..01fba17c8c6 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/testutils/RestoreTestCompleteness.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/testutils/RestoreTestCompletenessTest.java @@ -19,6 +19,9 @@ package org.apache.flink.table.planner.plan.nodes.exec.testutils; import org.apache.flink.table.planner.plan.nodes.exec.ExecNode; +import org.apache.flink.table.planner.plan.nodes.exec.batch.BatchExecHashAggregate; +import org.apache.flink.table.planner.plan.nodes.exec.batch.BatchExecNestedLoopJoin; +import org.apache.flink.table.planner.plan.nodes.exec.batch.BatchExecSortAggregate; import org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecGlobalWindowAggregate; import org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecLocalWindowAggregate; import org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecOverAggregate; @@ -34,41 +37,48 @@ import org.apache.flink.table.planner.plan.utils.ExecNodeMetadataUtil.ExecNodeNa import org.apache.flink.shaded.guava31.com.google.common.reflect.ClassPath; -import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; import java.io.IOException; import java.lang.reflect.InvocationTargetException; import java.lang.reflect.Method; +import java.util.Arrays; import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Set; import java.util.stream.Collectors; +import static org.assertj.core.api.Assertions.fail; + /** Validate restore tests exists for Exec Nodes. */ -public class RestoreTestCompleteness { +class RestoreTestCompletenessTest { private static final Set<Class<? extends ExecNode<?>>> SKIP_EXEC_NODES = - new HashSet<Class<? extends ExecNode<?>>>() { - { - /** TODO: Remove after FLINK-33676 is merged. */ - add(StreamExecWindowAggregate.class); - add(StreamExecLocalWindowAggregate.class); - add(StreamExecGlobalWindowAggregate.class); - - /** TODO: Remove after FLINK-33805 is merged. */ - add(StreamExecOverAggregate.class); - - /** Ignoring python based exec nodes temporarily. */ - add(StreamExecPythonCalc.class); - add(StreamExecPythonCorrelate.class); - add(StreamExecPythonOverAggregate.class); - add(StreamExecPythonGroupAggregate.class); - add(StreamExecPythonGroupTableAggregate.class); - add(StreamExecPythonGroupWindowAggregate.class); - } - }; + new HashSet( + Arrays.asList( + /* Ignoring python based exec nodes temporarily. */ + StreamExecPythonCalc.class, + StreamExecPythonCorrelate.class, + StreamExecPythonOverAggregate.class, + StreamExecPythonGroupAggregate.class, + StreamExecPythonGroupTableAggregate.class, + StreamExecPythonGroupWindowAggregate.class, + + /** TODO: Remove after FLINK-33676 is merged. */ + StreamExecWindowAggregate.class, + StreamExecLocalWindowAggregate.class, + StreamExecGlobalWindowAggregate.class, + + /** TODO: Remove after FLINK-33805 is merged. */ + StreamExecOverAggregate.class, + + // There is jira for these 2 batch tests + // https://issues.apache.org/jira/browse/FLINK-40306 + BatchExecHashAggregate.class, + BatchExecNestedLoopJoin.class, + // restore tests for batch sort aggregate was added in later releases + BatchExecSortAggregate.class)); private Class<? extends ExecNode<?>> getExecNode(Class<?> restoreTest) throws NoSuchMethodException, InvocationTargetException, InstantiationException, @@ -93,19 +103,17 @@ public class RestoreTestCompleteness { } @Test - public void testMissingRestoreTest() + void testMissingRestoreTest() throws IOException, NoSuchMethodException, InstantiationException, IllegalAccessException, InvocationTargetException { Map<ExecNodeNameVersion, Class<? extends ExecNode<?>>> versionedExecNodes = ExecNodeMetadataUtil.getVersionedExecNodes(); Set<ClassPath.ClassInfo> classesInPackage = - ClassPath.from(this.getClass().getClassLoader()) - .getTopLevelClassesRecursive( - "org.apache.flink.table.planner.plan.nodes.exec.stream") - .stream() - .filter(x -> RestoreTestBase.class.isAssignableFrom(x.load())) - .collect(Collectors.toSet()); + new HashSet<>( + gatherClasses( + RestoreTestBase.class, + "org.apache.flink.table.planner.plan.nodes.exec.stream")); Set<Class<? extends ExecNode<?>>> execNodesWithRestoreTests = new HashSet<>(); @@ -121,18 +129,32 @@ public class RestoreTestCompleteness { } } + Set<Class<? extends ExecNode<?>>> productionExecNodes = ExecNodeMetadataUtil.execNodes(); for (Map.Entry<ExecNodeNameVersion, Class<? extends ExecNode<?>>> entry : versionedExecNodes.entrySet()) { ExecNodeNameVersion execNodeNameVersion = entry.getKey(); Class<? extends ExecNode<?>> execNode = entry.getValue(); - if (!SKIP_EXEC_NODES.contains(execNode)) { - final String msg = + // Ignore test-only nodes that other tests leak into the shared LOOKUP_MAP via + // addTestNode(). + if (!productionExecNodes.contains(execNode)) { + continue; + } + if (!SKIP_EXEC_NODES.contains(execNode) + && !execNodesWithRestoreTests.contains(execNode)) { + fail( "Missing restore test for " + execNodeNameVersion + "\nPlease add a restore test for " - + execNode.toString(); - Assertions.assertTrue(execNodesWithRestoreTests.contains(execNode), msg); + + execNode.toString()); } } } + + private Set<ClassPath.ClassInfo> gatherClasses(Class<?> clazz, String packageName) + throws IOException { + return ClassPath.from(this.getClass().getClassLoader()) + .getTopLevelClassesRecursive(packageName).stream() + .filter(x -> clazz.isAssignableFrom(x.load())) + .collect(Collectors.toSet()); + } }
