From fa9cc38a9b64230f6884d6b9c486837ab57db3d4 Mon Sep 17 00:00:00 2001 From: hutiefang Date: Mon, 22 Jun 2026 04:01:01 +0800 Subject: [PATCH 1/2] [ISSUE-4339][Bug] Fix latest savepoint lookup --- .../impl/FlinkSavepointServiceImpl.java | 3 ++ .../service/FlinkSavepointServiceTest.java | 29 +++++++++++++++++++ 2 files changed, 32 insertions(+) diff --git a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/service/impl/FlinkSavepointServiceImpl.java b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/service/impl/FlinkSavepointServiceImpl.java index 743cbe6156..c580956fd1 100644 --- a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/service/impl/FlinkSavepointServiceImpl.java +++ b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/service/impl/FlinkSavepointServiceImpl.java @@ -131,6 +131,9 @@ public FlinkSavepoint getLatest(Long id) { return this.lambdaQuery() .eq(FlinkSavepoint::getAppId, id) .eq(FlinkSavepoint::getLatest, true) + .orderByDesc(FlinkSavepoint::getTriggerTime) + .orderByDesc(FlinkSavepoint::getId) + .last("limit 1") .one(); } diff --git a/streampark-console/streampark-console-service/src/test/java/org/apache/streampark/console/core/service/FlinkSavepointServiceTest.java b/streampark-console/streampark-console-service/src/test/java/org/apache/streampark/console/core/service/FlinkSavepointServiceTest.java index 54eece310c..485e099529 100644 --- a/streampark-console/streampark-console-service/src/test/java/org/apache/streampark/console/core/service/FlinkSavepointServiceTest.java +++ b/streampark-console/streampark-console-service/src/test/java/org/apache/streampark/console/core/service/FlinkSavepointServiceTest.java @@ -26,8 +26,10 @@ import org.apache.streampark.console.core.entity.FlinkApplicationConfig; import org.apache.streampark.console.core.entity.FlinkEffective; import org.apache.streampark.console.core.entity.FlinkEnv; +import org.apache.streampark.console.core.entity.FlinkSavepoint; import org.apache.streampark.console.core.enums.ConfigFileTypeEnum; import org.apache.streampark.console.core.enums.EffectiveTypeEnum; +import org.apache.streampark.console.core.mapper.FlinkSavepointMapper; import org.apache.streampark.console.core.service.application.FlinkApplicationConfigService; import org.apache.streampark.console.core.service.application.FlinkApplicationManageService; import org.apache.streampark.console.core.service.impl.FlinkSavepointServiceImpl; @@ -38,6 +40,8 @@ import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; +import java.util.Date; + import static org.apache.flink.configuration.CheckpointingOptions.SAVEPOINT_DIRECTORY; import static org.apache.flink.streaming.api.environment.ExecutionCheckpointingOptions.CHECKPOINTING_INTERVAL; import static org.assertj.core.api.Assertions.assertThat; @@ -66,6 +70,9 @@ class FlinkSavepointServiceTest extends SpringUnitTestBase { @Autowired FlinkApplicationManageService applicationManageService; + @Autowired + private FlinkSavepointMapper savepointMapper; + @AfterEach void cleanTestRecordsInDatabase() { savepointService.remove(new QueryWrapper<>()); @@ -92,6 +99,19 @@ void testGetSavepointFromDynamicProps() { .isEmpty(); } + @Test + void testGetLatestReturnsNewestSavepointWhenDuplicateLatestRecordsExist() { + Long appId = 1L; + FlinkSavepoint older = latestSavepoint(appId, "hdfs:///older", new Date(1000)); + FlinkSavepoint newer = latestSavepoint(appId, "hdfs:///newer", new Date(2000)); + savepointMapper.insert(older); + savepointMapper.insert(newer); + + FlinkSavepoint latest = savepointService.getLatest(appId); + + assertThat(latest.getPath()).isEqualTo("hdfs:///newer"); + } + @Test void testGetSavepointFromAppCfgIfStreamParkOrSQLJob() { FlinkSavepointServiceImpl savepointServiceImpl = (FlinkSavepointServiceImpl) savepointService; @@ -190,4 +210,13 @@ void testGetSavepointFromDeployLayer() throws JsonProcessingException { // Test for it with the configured non-empty target value } + + private FlinkSavepoint latestSavepoint(Long appId, String path, Date triggerTime) { + FlinkSavepoint savepoint = new FlinkSavepoint(); + savepoint.setAppId(appId); + savepoint.setLatest(true); + savepoint.setPath(path); + savepoint.setTriggerTime(triggerTime); + return savepoint; + } } From 5c8709ef0f122dc247975595901aa6d1b600a6df Mon Sep 17 00:00:00 2001 From: benjobs Date: Thu, 25 Jun 2026 11:10:04 +0800 Subject: [PATCH 2/2] Update FlinkSavepointServiceTest.java trigger ci --- .../console/core/service/FlinkSavepointServiceTest.java | 5 ----- 1 file changed, 5 deletions(-) diff --git a/streampark-console/streampark-console-service/src/test/java/org/apache/streampark/console/core/service/FlinkSavepointServiceTest.java b/streampark-console/streampark-console-service/src/test/java/org/apache/streampark/console/core/service/FlinkSavepointServiceTest.java index 485e099529..925adc1280 100644 --- a/streampark-console/streampark-console-service/src/test/java/org/apache/streampark/console/core/service/FlinkSavepointServiceTest.java +++ b/streampark-console/streampark-console-service/src/test/java/org/apache/streampark/console/core/service/FlinkSavepointServiceTest.java @@ -204,11 +204,6 @@ void testGetSavepointFromDeployLayer() throws JsonProcessingException { assertThatThrownBy(() -> savepointServiceImpl.getSavepointFromDeployLayer(application)) .isInstanceOf(NullPointerException.class); - // Ignored. - // Test for it with empty config - // Test for it with the configured empty target value - // Test for it with the configured non-empty target value - } private FlinkSavepoint latestSavepoint(Long appId, String path, Date triggerTime) {