This is an automated email from the ASF dual-hosted git repository. snuyanzin pushed a commit to branch release-2.1 in repository https://gitbox.apache.org/repos/asf/flink.git
commit 91a604795c70ce24b7a95949555d2412289550d1 Author: Sergey Nuyanzin <[email protected]> AuthorDate: Tue Aug 4 09:06:33 2026 +0200 [FLINK-40307][table] was never executed in CI --- ...eness.java => RestoreTestCompletenessTest.java} | 84 +++++++++++++++------- 1 file changed, 60 insertions(+), 24 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 60% 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 9e6e124449a..99ac5b09182 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,13 @@ 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.StreamExecDeltaJoin; +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.StreamExecProcessTableFunction; import org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecPythonCalc; import org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecPythonCorrelate; import org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecPythonGroupAggregate; @@ -30,7 +37,6 @@ import org.apache.flink.table.planner.plan.utils.ExecNodeMetadataUtil.ExecNodeNa import org.apache.flink.shaded.guava33.com.google.common.reflect.ClassPath; -import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; import java.io.IOException; @@ -42,21 +48,34 @@ 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<?>>>() { - { - /** Ignoring python based exec nodes temporarily. */ - add(StreamExecPythonCalc.class); - add(StreamExecPythonCorrelate.class); - add(StreamExecPythonOverAggregate.class); - add(StreamExecPythonGroupAggregate.class); - add(StreamExecPythonGroupTableAggregate.class); - add(StreamExecPythonGroupWindowAggregate.class); - } - }; + Set.of( + /* Ignoring python based exec nodes temporarily. */ + StreamExecPythonCalc.class, + StreamExecPythonCorrelate.class, + StreamExecPythonOverAggregate.class, + StreamExecPythonGroupAggregate.class, + StreamExecPythonGroupTableAggregate.class, + StreamExecPythonGroupWindowAggregate.class, + + // Covered by tests in WindowAggregateEventTimeRestoreTest + StreamExecLocalWindowAggregate.class, + StreamExecGlobalWindowAggregate.class, + // restore tests for delta join and PTF were added in later releases + StreamExecDeltaJoin.class, + StreamExecProcessTableFunction.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, @@ -85,7 +104,7 @@ public class RestoreTestCompleteness { } @Test - public void testMissingRestoreTest() + void testMissingRestoreTest() throws IOException, NoSuchMethodException, InstantiationException, @@ -95,12 +114,14 @@ public class RestoreTestCompleteness { 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")); + classesInPackage.addAll( + gatherClasses( + BatchRestoreTestBase.class, + "org.apache.flink.table.planner.plan.nodes.exec.batch")); Set<Class<? extends ExecNode<?>>> execNodesWithRestoreTests = new HashSet<>(); @@ -116,18 +137,33 @@ 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()); + } }
