From 29fb4a4f3e305ea42d8441d787d12bf01a4201fb Mon Sep 17 00:00:00 2001 From: Sergey Nuyanzin Date: Mon, 3 Aug 2026 12:39:58 +0200 Subject: [PATCH 1/5] [FLINK-40307][table] `RestoreTestCompleteness` was never executed in CI --- .../batch/SortAggregateBatchRestoreTest.java | 2 +- .../stream/AsyncCorrelateRestoreTest.java | 2 +- .../ProcessTableFunctionRestoreTests.java | 2 +- .../stream/WatermarkAssignerRestoreTest.java | 2 +- ....java => RestoreTestCompletenessTest.java} | 78 ++++++++++++------ .../plan/async-correlate-catalog-func.json | 0 .../savepoint/_metadata | Bin .../plan/async-correlate-exception.json | 0 .../savepoint/_metadata | Bin .../plan/async-correlate-join-filter.json | 0 .../savepoint/_metadata | Bin .../plan/async-correlate-left-join.json | 0 .../savepoint/_metadata | Bin .../plan/async-correlate-system-func.json | 0 .../savepoint/_metadata | Bin 15 files changed, 57 insertions(+), 29 deletions(-) rename flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/testutils/{RestoreTestCompleteness.java => RestoreTestCompletenessTest.java} (65%) rename flink-table/flink-table-planner/src/test/resources/restore-tests/{stream-exec-correlate_1 => stream-exec-async-correlate_1}/async-correlate-catalog-func/plan/async-correlate-catalog-func.json (100%) rename flink-table/flink-table-planner/src/test/resources/restore-tests/{stream-exec-correlate_1 => stream-exec-async-correlate_1}/async-correlate-catalog-func/savepoint/_metadata (100%) rename flink-table/flink-table-planner/src/test/resources/restore-tests/{stream-exec-correlate_1 => stream-exec-async-correlate_1}/async-correlate-exception/plan/async-correlate-exception.json (100%) rename flink-table/flink-table-planner/src/test/resources/restore-tests/{stream-exec-correlate_1 => stream-exec-async-correlate_1}/async-correlate-exception/savepoint/_metadata (100%) rename flink-table/flink-table-planner/src/test/resources/restore-tests/{stream-exec-correlate_1 => stream-exec-async-correlate_1}/async-correlate-join-filter/plan/async-correlate-join-filter.json (100%) rename flink-table/flink-table-planner/src/test/resources/restore-tests/{stream-exec-correlate_1 => stream-exec-async-correlate_1}/async-correlate-join-filter/savepoint/_metadata (100%) rename flink-table/flink-table-planner/src/test/resources/restore-tests/{stream-exec-correlate_1 => stream-exec-async-correlate_1}/async-correlate-left-join/plan/async-correlate-left-join.json (100%) rename flink-table/flink-table-planner/src/test/resources/restore-tests/{stream-exec-correlate_1 => stream-exec-async-correlate_1}/async-correlate-left-join/savepoint/_metadata (100%) rename flink-table/flink-table-planner/src/test/resources/restore-tests/{stream-exec-correlate_1 => stream-exec-async-correlate_1}/async-correlate-system-func/plan/async-correlate-system-func.json (100%) rename flink-table/flink-table-planner/src/test/resources/restore-tests/{stream-exec-correlate_1 => stream-exec-async-correlate_1}/async-correlate-system-func/savepoint/_metadata (100%) 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 d329a2d04ecb13..a4581a32cea8ed 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 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 2d52118b225262..6d8214e1caa8ca 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 @@ 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 03ef42de0e142a..caccbfc02190db 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 @@ /** 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 0a684e4c13360f..3d63a1523e14eb 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 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 ea183bf3dbad39..ec8cff2bddf0fa 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.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.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>> SKIP_EXEC_NODES = - new HashSet>>() { - { - /** 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> getExecNode(Class restoreTest) throws NoSuchMethodException, @@ -87,7 +98,7 @@ private List>> getChildExecNodes(Class restoreTes } @Test - public void testMissingRestoreTest() + void testMissingRestoreTest() throws IOException, NoSuchMethodException, InstantiationException, @@ -97,12 +108,14 @@ public void testMissingRestoreTest() ExecNodeMetadataUtil.getVersionedExecNodes(); Set 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>> execNodesWithRestoreTests = new HashSet<>(); @@ -118,18 +131,33 @@ public void testMissingRestoreTest() } } + Set>> productionExecNodes = ExecNodeMetadataUtil.execNodes(); for (Map.Entry>> entry : versionedExecNodes.entrySet()) { ExecNodeNameVersion execNodeNameVersion = entry.getKey(); Class> 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 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 From 28715de77db636a673ce616f51026e64d17c9cb8 Mon Sep 17 00:00:00 2001 From: Sergey Nuyanzin Date: Mon, 3 Aug 2026 13:03:12 +0200 Subject: [PATCH 2/5] [FLINK-40285][table] `MLPredictSemanticTests` fails because of `ON CONFLICT` --- .../exec/stream/MLPredictTestPrograms.java | 19 +++++++++++++------ 1 file changed, 13 insertions(+), 6 deletions(-) diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/MLPredictTestPrograms.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/MLPredictTestPrograms.java index 26d4903a3e41de..205aa672e25d6d 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/MLPredictTestPrograms.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/MLPredictTestPrograms.java @@ -20,6 +20,7 @@ import org.apache.flink.table.api.DataTypes; import org.apache.flink.table.api.Expressions; +import org.apache.flink.table.api.InsertConflictStrategy; import org.apache.flink.table.api.ModelDescriptor; import org.apache.flink.table.api.Schema; import org.apache.flink.table.api.config.ExecutionConfigOptions; @@ -182,7 +183,8 @@ public class MLPredictTestPrograms { env.from("features").asArgument("INPUT"), env.fromModel("chatgpt").asArgument("MODEL"), descriptor("feature").asArgument("ARGS")), - "sink") + "sink", + InsertConflictStrategy.deduplicate()) .build(); public static final TableTestProgram ASYNC_ML_PREDICT_TABLE_API = @@ -209,7 +211,8 @@ public class MLPredictTestPrograms { DataTypes.STRING()) .notNull()) .asArgument("CONFIG")), - "sink") + "sink", + InsertConflictStrategy.deduplicate()) .build(); public static final TableTestProgram ASYNC_ML_PREDICT_TABLE_API_MAP_EXPRESSION_CONFIG = @@ -235,7 +238,8 @@ public class MLPredictTestPrograms { "max-concurrent-operations", "10") .asArgument("CONFIG")), - "sink") + "sink", + InsertConflictStrategy.deduplicate()) .build(); public static final TableTestProgram ML_PREDICT_MODEL_API = @@ -248,7 +252,8 @@ public class MLPredictTestPrograms { env.fromModel("chatgpt") .predict( env.from("features"), ColumnList.of("feature")), - "sink") + "sink", + InsertConflictStrategy.deduplicate()) .build(); public static final TableTestProgram ASYNC_ML_PREDICT_MODEL_API = @@ -270,7 +275,8 @@ public class MLPredictTestPrograms { "true", "max-concurrent-operations", "10")), - "sink") + "sink", + InsertConflictStrategy.deduplicate()) .build(); public static final TableTestProgram ML_PREDICT_ANON_MODEL_API = @@ -304,6 +310,7 @@ public class MLPredictTestPrograms { .build()) .predict( env.from("features"), ColumnList.of("feature")), - "sink") + "sink", + InsertConflictStrategy.deduplicate()) .build(); } From 331f68d6a81ce2f04456430f4c1b26ad96ab9ffe Mon Sep 17 00:00:00 2001 From: Sergey Nuyanzin Date: Mon, 3 Aug 2026 20:13:56 +0200 Subject: [PATCH 3/5] [FLINK-40317][tests] Make DeletesByKeySemanticTests more stable --- .../exec/stream/DeletesByKeyPrograms.java | 18 ++++-------------- 1 file changed, 4 insertions(+), 14 deletions(-) diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java index 6599f649950551..a66837d1c45421 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java @@ -237,13 +237,8 @@ public final class DeletesByKeyPrograms { "`value` INT") .addOption("changelog-mode", "I,UA,D") .addOption("sink.supports-delete-by-key", "false") - .consumedValues( - "+I[1, Alice, 10]", - "+I[2, Bob, 20]", - "+I[3, Emily, 30]", - "-D[1, Alice, 10]", - "+U[3, Emily, 40]", - "+U[2, BOB, 20]") + .testMaterializedData() + .consumedValues("+I[3, Emily, 40]", "+I[2, BOB, 20]") .build()) .runSql( "INSERT INTO sink_t SELECT l.id, r.name, l.`value` FROM left_t l JOIN right_t r ON l.id = r.id") @@ -291,13 +286,8 @@ public final class DeletesByKeyPrograms { "`value` INT") .addOption("changelog-mode", "I,UA,D") .addOption("sink.supports-delete-by-key", "true") - .consumedValues( - "+I[1, Alice, 10]", - "+I[2, Bob, 20]", - "+I[3, Emily, 30]", - "-D[1, Alice, null]", - "+U[3, Emily, 40]", - "+U[2, BOB, 20]") + .testMaterializedData() + .consumedValues("+I[2, BOB, 20]", "+I[3, Emily, 40]") .build()) .runSql( "INSERT INTO sink_t SELECT l.id, r.name, l.`value` FROM left_t l JOIN right_t r ON l.id = r.id") From 37b31b8b4d1cf64a6febbf69c0cccbeb786cec53 Mon Sep 17 00:00:00 2001 From: Sergey Nuyanzin Date: Mon, 3 Aug 2026 12:46:46 +0200 Subject: [PATCH 4/5] [FLINK-40284][tests] Make tests ending with `Tests` executing in CI --- pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index 49a425f9455b04..b7eb7fc8bbd3e1 100644 --- a/pom.xml +++ b/pom.xml @@ -221,7 +221,7 @@ under the License. 256 1.0 - **/*Test.* + **/*Test.*,**/*Tests.* 1.1.10.7 3.18.0 From 1630fd4a078bc3eb64275c515a30eafb28ba4018 Mon Sep 17 00:00:00 2001 From: Sergey Nuyanzin Date: Mon, 3 Aug 2026 12:35:20 +0200 Subject: [PATCH 5/5] [FLINK-40284][tests] Archunit should fail in case of tests not matching name requirements This closes #28887. --- .../TestCodeArchitectureTestBase.java | 3 + .../architecture/rules/TestNamingRules.java | 86 +++++++++++++++++++ 2 files changed, 89 insertions(+) create mode 100644 flink-architecture-tests/flink-architecture-tests-test/src/main/java/org/apache/flink/architecture/rules/TestNamingRules.java diff --git a/flink-architecture-tests/flink-architecture-tests-test/src/main/java/org/apache/flink/architecture/TestCodeArchitectureTestBase.java b/flink-architecture-tests/flink-architecture-tests-test/src/main/java/org/apache/flink/architecture/TestCodeArchitectureTestBase.java index a33ad9b6d6653f..9b7984d002b75e 100644 --- a/flink-architecture-tests/flink-architecture-tests-test/src/main/java/org/apache/flink/architecture/TestCodeArchitectureTestBase.java +++ b/flink-architecture-tests/flink-architecture-tests-test/src/main/java/org/apache/flink/architecture/TestCodeArchitectureTestBase.java @@ -19,6 +19,7 @@ package org.apache.flink.architecture; import org.apache.flink.architecture.rules.ITCaseRules; +import org.apache.flink.architecture.rules.TestNamingRules; import com.tngtech.archunit.junit.ArchTest; import com.tngtech.archunit.junit.ArchTests; @@ -33,4 +34,6 @@ public class TestCodeArchitectureTestBase { @ArchTest public static final ArchTests ITCASE = ArchTests.in(ITCaseRules.class); + + @ArchTest public static final ArchTests TEST_NAMING = ArchTests.in(TestNamingRules.class); } diff --git a/flink-architecture-tests/flink-architecture-tests-test/src/main/java/org/apache/flink/architecture/rules/TestNamingRules.java b/flink-architecture-tests/flink-architecture-tests-test/src/main/java/org/apache/flink/architecture/rules/TestNamingRules.java new file mode 100644 index 00000000000000..79810fa043f729 --- /dev/null +++ b/flink-architecture-tests/flink-architecture-tests-test/src/main/java/org/apache/flink/architecture/rules/TestNamingRules.java @@ -0,0 +1,86 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.architecture.rules; + +import com.tngtech.archunit.base.DescribedPredicate; +import com.tngtech.archunit.core.domain.JavaClass; +import com.tngtech.archunit.junit.ArchTest; +import com.tngtech.archunit.lang.ArchRule; + +import java.util.Arrays; +import java.util.List; + +import static com.tngtech.archunit.core.domain.JavaModifier.ABSTRACT; +import static org.apache.flink.architecture.common.GivenJavaClasses.javaClassesThat; + +/** + * Rules ensuring executable test classes are named so the build actually runs them. + * + *

Surefire only runs the unit include pattern {@code **}{@code /*Test.*} in the {@code test} + * phase; integration tests follow the {@code *ITCase} convention. A concrete class that carries (or + * inherits) JUnit test methods but is named otherwise (e.g. {@code *Tests}) is silently skipped by + * the unit run. This rule flags such classes so they are renamed to {@code *Test} or {@code + * *ITCase}. + */ +public class TestNamingRules { + + /** JUnit 5 and (for modules still mid-migration) JUnit 4 test method annotations. */ + private static final List TEST_METHOD_ANNOTATIONS = + Arrays.asList( + "org.junit.jupiter.api.Test", + "org.junit.jupiter.api.TestTemplate", + "org.junit.jupiter.api.RepeatedTest", + "org.junit.jupiter.api.TestFactory", + "org.junit.jupiter.params.ParameterizedTest", + "org.junit.Test"); + + /** + * A class JUnit would execute: it declares or inherits a test method. {@code getAllMethods()} + * covers inherited {@code @TestTemplate} methods, e.g. semantic-test suites that only extend a + * base and add no annotation themselves. + */ + private static final DescribedPredicate ARE_EXECUTABLE_TEST_CLASSES = + DescribedPredicate.describe( + "are executable JUnit test classes", + clazz -> + clazz.getAllMethods().stream() + .anyMatch( + method -> + TEST_METHOD_ANNOTATIONS.stream() + .anyMatch(method::isAnnotatedWith))); + + @ArchTest + public static final ArchRule TEST_CLASSES_SHOULD_BE_NAMED_TEST_OR_ITCASE = + javaClassesThat() + .areTopLevelClasses() + .and() + .doNotHaveModifier(ABSTRACT) + .and(ARE_EXECUTABLE_TEST_CLASSES) + .should() + .haveSimpleNameEndingWith("Test") + .orShould() + .haveSimpleNameEndingWith("Tests") + .orShould() + .haveSimpleNameEndingWith("ITCase") + // not every module has such classes + .allowEmptyShould(true) + .as( + "Executable test classes must be named *Test[s] or *ITCase so the surefire " + + "include pattern runs them"); +}