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..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 @@ -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; @@ -184,10 +204,14 @@ 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) { + FlinkSavepoint savepoint = new FlinkSavepoint(); + savepoint.setAppId(appId); + savepoint.setLatest(true); + savepoint.setPath(path); + savepoint.setTriggerTime(triggerTime); + return savepoint; } }