This is an automated email from the ASF dual-hosted git repository. snuyanzin pushed a commit to branch release-2.3 in repository https://gitbox.apache.org/repos/asf/flink.git
commit 09a8bb03e87037fadf098504b48f4bd755f7bf09 Author: Sergey Nuyanzin <[email protected]> AuthorDate: Mon Aug 3 12:39:58 2026 +0200 [FLINK-40307][table] `RestoreTestCompleteness` was never executed in CI --- .../exec/batch/SortAggregateBatchRestoreTest.java | 2 +- .../exec/stream/AsyncCorrelateRestoreTest.java | 2 +- .../stream/ProcessTableFunctionRestoreTests.java | 2 +- .../exec/stream/WatermarkAssignerRestoreTest.java | 2 +- ...eness.java => RestoreTestCompletenessTest.java} | 78 ++++++++++++++------- .../plan/async-correlate-catalog-func.json | 0 .../savepoint/_metadata | Bin .../plan/async-correlate-exception.json | 0 .../async-correlate-exception/savepoint/_metadata | Bin .../plan/async-correlate-join-filter.json | 0 .../savepoint/_metadata | Bin .../plan/async-correlate-left-join.json | 0 .../async-correlate-left-join/savepoint/_metadata | Bin .../plan/async-correlate-system-func.json | 0 .../savepoint/_metadata | Bin 15 files changed, 57 insertions(+), 29 deletions(-) diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/batch/SortAggregateBatchRestoreTest.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/batch/SortAggregateBatchRestoreTest.java index d329a2d04ec..a4581a32cea 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/batch/SortAggregateBatchRestoreTest.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/batch/SortAggregateBatchRestoreTest.java @@ -24,7 +24,7 @@ import org.apache.flink.table.test.program.TableTestProgram; import java.util.List; /** Batch Compiled Plan tests for {@link BatchExecSortAggregate}. */ -class SortAggregateBatchRestoreTest extends BatchRestoreTestBase { +public class SortAggregateBatchRestoreTest extends BatchRestoreTestBase { public SortAggregateBatchRestoreTest() { super(BatchExecSortAggregate.class); diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/AsyncCorrelateRestoreTest.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/AsyncCorrelateRestoreTest.java index 2d52118b225..6d8214e1caa 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/AsyncCorrelateRestoreTest.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/AsyncCorrelateRestoreTest.java @@ -28,7 +28,7 @@ import java.util.List; public class AsyncCorrelateRestoreTest extends RestoreTestBase { public AsyncCorrelateRestoreTest() { - super(StreamExecCorrelate.class); + super(StreamExecAsyncCorrelate.class); } @Override diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionRestoreTests.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionRestoreTests.java index 03ef42de0e1..caccbfc0219 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionRestoreTests.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionRestoreTests.java @@ -26,7 +26,7 @@ import java.util.List; /** Restore tests for {@link StreamExecProcessTableFunction}. */ public class ProcessTableFunctionRestoreTests extends RestoreTestBase { - protected ProcessTableFunctionRestoreTests() { + public ProcessTableFunctionRestoreTests() { super(StreamExecProcessTableFunction.class); } diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/WatermarkAssignerRestoreTest.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/WatermarkAssignerRestoreTest.java index 0a684e4c133..3d63a1523e1 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/WatermarkAssignerRestoreTest.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/WatermarkAssignerRestoreTest.java @@ -24,7 +24,7 @@ import org.apache.flink.table.test.program.TableTestProgram; import java.util.List; /** Restore tests for {@link StreamExecWatermarkAssigner}. */ -class WatermarkAssignerRestoreTest extends RestoreTestBase { +public class WatermarkAssignerRestoreTest extends RestoreTestBase { public WatermarkAssignerRestoreTest() { super(StreamExecWatermarkAssigner.class); 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 65% 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 ea183bf3dba..ec8cff2bddf 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,10 @@ 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.stream.StreamExecGlobalWindowAggregate; +import org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecLocalWindowAggregate; import org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecPythonAsyncCalc; import org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecPythonCalc; import org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecPythonCorrelate; @@ -31,7 +35,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; @@ -43,22 +46,30 @@ 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); - add(StreamExecPythonAsyncCalc.class); - } - }; + Set.of( + /* Ignoring python based exec nodes temporarily. */ + StreamExecPythonCalc.class, + StreamExecPythonCorrelate.class, + StreamExecPythonOverAggregate.class, + StreamExecPythonGroupAggregate.class, + StreamExecPythonGroupTableAggregate.class, + StreamExecPythonGroupWindowAggregate.class, + StreamExecPythonAsyncCalc.class, + + // Covered by tests in WindowAggregateEventTimeRestoreTest + StreamExecLocalWindowAggregate.class, + StreamExecGlobalWindowAggregate.class, + + // There is jira for these 2 batch tests + // https://issues.apache.org/jira/browse/FLINK-40306 + BatchExecHashAggregate.class, + BatchExecNestedLoopJoin.class); private Class<? extends ExecNode<?>> getExecNode(Class<?> restoreTest) throws NoSuchMethodException, @@ -87,7 +98,7 @@ public class RestoreTestCompleteness { } @Test - public void testMissingRestoreTest() + void testMissingRestoreTest() throws IOException, NoSuchMethodException, InstantiationException, @@ -97,12 +108,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<>(); @@ -118,18 +131,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()); + } } diff --git a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-catalog-func/plan/async-correlate-catalog-func.json b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-catalog-func/plan/async-correlate-catalog-func.json similarity index 100% rename from flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-catalog-func/plan/async-correlate-catalog-func.json rename to flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-catalog-func/plan/async-correlate-catalog-func.json diff --git a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-catalog-func/savepoint/_metadata b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-catalog-func/savepoint/_metadata similarity index 100% rename from flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-catalog-func/savepoint/_metadata rename to flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-catalog-func/savepoint/_metadata diff --git a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-exception/plan/async-correlate-exception.json b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-exception/plan/async-correlate-exception.json similarity index 100% rename from flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-exception/plan/async-correlate-exception.json rename to flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-exception/plan/async-correlate-exception.json diff --git a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-exception/savepoint/_metadata b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-exception/savepoint/_metadata similarity index 100% rename from flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-exception/savepoint/_metadata rename to flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-exception/savepoint/_metadata diff --git a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-join-filter/plan/async-correlate-join-filter.json b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-join-filter/plan/async-correlate-join-filter.json similarity index 100% rename from flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-join-filter/plan/async-correlate-join-filter.json rename to flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-join-filter/plan/async-correlate-join-filter.json diff --git a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-join-filter/savepoint/_metadata b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-join-filter/savepoint/_metadata similarity index 100% rename from flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-join-filter/savepoint/_metadata rename to flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-join-filter/savepoint/_metadata diff --git a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-left-join/plan/async-correlate-left-join.json b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-left-join/plan/async-correlate-left-join.json similarity index 100% rename from flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-left-join/plan/async-correlate-left-join.json rename to flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-left-join/plan/async-correlate-left-join.json diff --git a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-left-join/savepoint/_metadata b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-left-join/savepoint/_metadata similarity index 100% rename from flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-left-join/savepoint/_metadata rename to flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-left-join/savepoint/_metadata diff --git a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-system-func/plan/async-correlate-system-func.json b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-system-func/plan/async-correlate-system-func.json similarity index 100% rename from flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-system-func/plan/async-correlate-system-func.json rename to flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-system-func/plan/async-correlate-system-func.json diff --git a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-system-func/savepoint/_metadata b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-system-func/savepoint/_metadata similarity index 100% rename from flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-system-func/savepoint/_metadata rename to flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-system-func/savepoint/_metadata
