From 2d071b3ff2e8d6f417662879aacf7eb81b1bf858 Mon Sep 17 00:00:00 2001 From: ShengHuang Date: Mon, 9 Feb 2026 20:25:28 +0800 Subject: [PATCH 1/6] feat(model): add built-in auto scheduled segment build with model-level config and scheduler --- .../setting/SettingModel/SettingModel.vue | 79 +++++++++++++++++ .../setting/SettingModel/locales.js | 11 ++- .../kylin/metadata/model/SegmentConfig.java | 3 + .../metadata/project/ProjectInstance.java | 3 +- .../cube/model/NSegmentConfigHelperTest.java | 2 +- .../controller/open/OpenModelController.java | 1 + .../rest/request/ModelConfigRequest.java | 3 + .../rest/response/ModelConfigResponse.java | 3 + .../kylin/rest/service/ModelBuildService.java | 33 +++++-- .../kylin/rest/service/ModelService.java | 67 ++++++++++++++- .../kylin/rest/service/ModelServiceTest.java | 85 +++++++++++++++++++ 11 files changed, 280 insertions(+), 10 deletions(-) diff --git a/kystudio/src/components/setting/SettingModel/SettingModel.vue b/kystudio/src/components/setting/SettingModel/SettingModel.vue index 9e3725d3762..a3dd57c965f 100644 --- a/kystudio/src/components/setting/SettingModel/SettingModel.vue +++ b/kystudio/src/components/setting/SettingModel/SettingModel.vue @@ -46,6 +46,15 @@ +
+ + {{$t('autoSegmentBuild')}}{{formatAutoSegmentBuild(scope.row.auto_segment_build)}} + + + + + +
@@ -341,6 +347,12 @@ export default class SettingStorage extends Vue { // return largestRange || '' return 'DAY' } + get isAutoSegmentBuildEndOfDay () { + return this.modelSettingForm.autoSegmentBuild.data_range_end_time === END_OF_DAY + } + set isAutoSegmentBuildEndOfDay (isEndOfDay) { + this.modelSettingForm.autoSegmentBuild.data_range_end_time = isEndOfDay ? END_OF_DAY : null + } validateSettingItem (rule, value, callback) { const autoMergeRanges = this.activeRow && this.activeRow.auto_merge_time_ranges || [] if (this.step === 'stepOne' && value === 'Retention Threshold' && !autoMergeRanges.length) { diff --git a/kystudio/src/components/setting/SettingModel/locales.js b/kystudio/src/components/setting/SettingModel/locales.js index d8ac1a2f98d..e9328b21e0f 100644 --- a/kystudio/src/components/setting/SettingModel/locales.js +++ b/kystudio/src/components/setting/SettingModel/locales.js @@ -63,6 +63,7 @@ export default { autoSegmentBuildLogicalOffset: 'Logical Date Offset (Day)', autoSegmentBuildRangeStart: 'Range Start Time', autoSegmentBuildRangeEnd: 'Range End Time', + autoSegmentBuildEndOfDay: 'End of logical day (next day 00:00)', autoSegmentBuildTip: 'Build segment automatically every day with logical day offset and configured time range.', autoSegmentBuildSummary: 'Daily {trigger}, D-{offset} {start}~{end}' } diff --git a/src/core-common/src/main/resources/kylin-defaults0.properties b/src/core-common/src/main/resources/kylin-defaults0.properties index 13007960639..8a7af4d8c6c 100644 --- a/src/core-common/src/main/resources/kylin-defaults0.properties +++ b/src/core-common/src/main/resources/kylin-defaults0.properties @@ -449,7 +449,8 @@ kylin.web.session.jdbc-encode-enabled=false kylin.security.user-password-encoder=org.apache.kylin.rest.security.CachedBCryptPasswordEncoder # model -kylin.model.auto-segment-build.dispatcher-cron=*/30 * * * * ? +# System-level internal scan cadence in milliseconds; model trigger times are read from metadata on every scan. +kylin.model.auto-segment-build.dispatcher-interval-ms=30000 kylin.model.recommendation-page-size=500 kylin.model.dimension-measure-name.max-length=300 diff --git a/src/data-loading-service/src/main/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentScheduler.java b/src/data-loading-service/src/main/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentScheduler.java index f7968b31f7b..5ff4ad5f1f1 100644 --- a/src/data-loading-service/src/main/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentScheduler.java +++ b/src/data-loading-service/src/main/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentScheduler.java @@ -20,24 +20,25 @@ import static org.apache.kylin.metadata.model.AutoSegmentBuildConfig.END_OF_DAY; -import java.time.Duration; import java.time.Instant; -import java.time.LocalDate; import java.time.LocalDateTime; import java.time.LocalTime; import java.time.ZoneId; import java.time.ZonedDateTime; import java.time.format.DateTimeFormatter; import java.time.format.DateTimeParseException; +import java.time.format.ResolverStyle; import java.util.List; import java.util.Locale; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import org.apache.commons.lang3.StringUtils; import org.apache.kylin.common.KylinConfig; import org.apache.kylin.common.util.Pair; -import org.apache.kylin.guava30.shaded.common.annotations.VisibleForTesting; +import org.apache.kylin.guava30.shaded.common.cache.Cache; +import org.apache.kylin.guava30.shaded.common.cache.CacheBuilder; import org.apache.kylin.guava30.shaded.common.collect.Lists; import org.apache.kylin.job.execution.AbstractExecutable; import org.apache.kylin.job.execution.ExecutableManager; @@ -54,6 +55,7 @@ import org.apache.kylin.rest.service.params.IncrementBuildSegmentParams; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.beans.factory.annotation.Value; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; @@ -63,17 +65,24 @@ @Slf4j @Component public class AutoBuildSegmentScheduler { - private static final Duration INITIAL_TRIGGER_LOOKBACK = Duration.ofMinutes(1); - private static final DateTimeFormatter TIME_FORMATTER = DateTimeFormatter.ofPattern("HH:mm:ss", Locale.ROOT); + private static final long DEFAULT_DISPATCHER_INTERVAL_MILLIS = 30_000L; + private static final DateTimeFormatter TIME_FORMATTER = DateTimeFormatter.ofPattern("HH:mm:ss", Locale.ROOT) + .withResolverStyle(ResolverStyle.STRICT); @Autowired @Qualifier("modelBuildService") private ModelBuildService modelBuildService; + @Value("${kylin.model.auto-segment-build.dispatcher-interval-ms:30000}") + private long dispatcherIntervalMillis = DEFAULT_DISPATCHER_INTERVAL_MILLIS; + private final AtomicBoolean dispatching = new AtomicBoolean(false); private final AtomicReference lastDispatchTime = new AtomicReference<>(); + // Avoid repeating the same malformed-metadata warning on every dispatcher scan. + private final Cache invalidConfigWarnings = CacheBuilder.newBuilder().maximumSize(10_000) + .expireAfterAccess(1, TimeUnit.DAYS).build(); - @Scheduled(cron = "${kylin.model.auto-segment-build.dispatcher-cron:*/30 * * * * ?}") + @Scheduled(fixedDelayString = "${kylin.model.auto-segment-build.dispatcher-interval-ms:30000}") public void schedulerAutoBuildSegment() { val currentTime = Instant.now(); if (!JobContextUtil.getJobContext(KylinConfig.getInstanceFromEnv()).getJobScheduler().isMaster()) { @@ -97,12 +106,11 @@ public void schedulerAutoBuildSegment() { private Instant getPreviousDispatchTime(Instant currentTime) { val previousTime = lastDispatchTime.get(); if (previousTime == null || previousTime.isAfter(currentTime)) { - return currentTime.minus(INITIAL_TRIGGER_LOOKBACK); + return currentTime.minusMillis(Math.max(1L, dispatcherIntervalMillis)); } return previousTime; } - @VisibleForTesting void dispatch(Instant previousTime, Instant currentTime) { val systemConfig = KylinConfig.readSystemKylinConfig(); val projectManager = NProjectManager.getInstance(systemConfig); @@ -132,15 +140,10 @@ private void dispatchProject(KylinConfig systemConfig, ProjectInstance project, private void dispatchModel(String project, NDataModel model, ZoneId zoneId, Instant previousTime, Instant currentTime) { - val autoSegmentBuild = getEligibleConfig(model); + val autoSegmentBuild = getEligibleConfig(project, model); if (autoSegmentBuild == null) { return; } - if (StringUtils.isBlank(autoSegmentBuild.getTriggerTime())) { - log.warn("Skip auto build segment because trigger_time is blank, project: {}, model: {}", project, - model.getUuid()); - return; - } val scheduledTime = getLatestScheduledTime(autoSegmentBuild.getTriggerTime(), zoneId, currentTime); val scheduledInstant = scheduledTime.toInstant(); @@ -150,16 +153,77 @@ private void dispatchModel(String project, NDataModel model, ZoneId zoneId, Inst submitJob(project, model, autoSegmentBuild, scheduledTime); } - private AutoSegmentBuildConfig getEligibleConfig(NDataModel model) { + private AutoSegmentBuildConfig getEligibleConfig(String project, NDataModel model) { + val configKey = project + "/" + model.getUuid(); if (model.isBroken() || model.isStreaming() || model.isMultiPartitionModel() || PartitionDesc.isEmptyPartitionDesc(model.getPartitionDesc()) || model.getSegmentConfig() == null) { + invalidConfigWarnings.invalidate(configKey); return null; } val autoSegmentBuild = model.getSegmentConfig().getAutoSegmentBuild(); - return autoSegmentBuild != null && autoSegmentBuild.isEnabled() ? autoSegmentBuild : null; + if (autoSegmentBuild == null || !autoSegmentBuild.isEnabled()) { + invalidConfigWarnings.invalidate(configKey); + return null; + } + + val invalidReason = getInvalidConfigReason(autoSegmentBuild); + if (invalidReason == null) { + invalidConfigWarnings.invalidate(configKey); + return autoSegmentBuild; + } + if (!StringUtils.equals(invalidReason, invalidConfigWarnings.getIfPresent(configKey))) { + log.warn("Skip auto build segment because of invalid config: {}, project: {}, model: {}", invalidReason, + project, model.getUuid()); + invalidConfigWarnings.put(configKey, invalidReason); + } + return null; + } + + private String getInvalidConfigReason(AutoSegmentBuildConfig config) { + List missingFields = Lists.newArrayList(); + if (StringUtils.isBlank(config.getTriggerTime())) { + missingFields.add("trigger_time"); + } + if (config.getLogicalDateOffsetDays() == null) { + missingFields.add("logical_date_offset_days"); + } + if (StringUtils.isBlank(config.getDataRangeStartTime())) { + missingFields.add("data_range_start_time"); + } + if (StringUtils.isBlank(config.getDataRangeEndTime())) { + missingFields.add("data_range_end_time"); + } + if (!missingFields.isEmpty()) { + return "missing required field(s): " + StringUtils.join(missingFields, ", "); + } + if (config.getLogicalDateOffsetDays() < 1) { + return "logical_date_offset_days must be >= 1"; + } + + try { + LocalTime.parse(config.getTriggerTime(), TIME_FORMATTER); + } catch (DateTimeParseException e) { + return "invalid trigger_time, expected HH:mm:ss"; + } + + Pair start; + Pair end; + try { + start = parseTime(config.getDataRangeStartTime(), false); + } catch (DateTimeParseException e) { + return "invalid data_range_start_time, expected HH:mm:ss"; + } + try { + end = parseTime(config.getDataRangeEndTime(), true); + } catch (DateTimeParseException e) { + return "invalid data_range_end_time, expected HH:mm:ss or " + END_OF_DAY; + } + if (!end.getSecond() && !start.getFirst().isBefore(end.getFirst())) { + return "data_range_start_time must be earlier than data_range_end_time"; + } + return null; } - @VisibleForTesting ZonedDateTime getLatestScheduledTime(String triggerTime, ZoneId zoneId, Instant currentTime) { val localCurrentTime = currentTime.atZone(zoneId); val parsedTriggerTime = LocalTime.parse(triggerTime, TIME_FORMATTER); @@ -170,7 +234,6 @@ ZonedDateTime getLatestScheduledTime(String triggerTime, ZoneId zoneId, Instant return scheduledTime; } - @VisibleForTesting void submitJob(String project, NDataModel model, AutoSegmentBuildConfig autoSegmentBuild, ZonedDateTime scheduledTime) { val modelId = model.getUuid(); @@ -196,7 +259,7 @@ void submitJob(String project, NDataModel model, AutoSegmentBuildConfig autoSegm val endMillis = String.valueOf(endDateTime.atZone(zoneId).toInstant().toEpochMilli()); val params = new IncrementBuildSegmentParams(project, modelId, startMillis, endMillis, model.getPartitionDesc(), model.getMultiPartitionDesc(), Lists.newArrayList(), true, null); - modelBuildService.incrementBuildSegmentsByScheduler(params, "System"); + modelBuildService.incrementBuildSegmentsByScheduler(params); log.info("Auto build segment submitted, project: {}, model: {}, scheduled time: {}, range: [{}, {})", project, modelId, scheduledTime, startMillis, endMillis); } catch (DateTimeParseException e) { @@ -206,7 +269,6 @@ void submitJob(String project, NDataModel model, AutoSegmentBuildConfig autoSegm } } - @VisibleForTesting boolean hasRunningModelBuildJob(String project, String modelId) { // Segment and model metadata mutations are serialized with any progressing build job on the same model. val executableManager = ExecutableManager.getInstance(KylinConfig.getInstanceFromEnv(), project); diff --git a/src/data-loading-service/src/test/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentSchedulerTest.java b/src/data-loading-service/src/test/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentSchedulerTest.java index fbe61850085..07619d63696 100644 --- a/src/data-loading-service/src/test/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentSchedulerTest.java +++ b/src/data-loading-service/src/test/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentSchedulerTest.java @@ -19,6 +19,7 @@ package org.apache.kylin.rest.scheduler; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; import java.time.Instant; import java.time.LocalDate; @@ -26,6 +27,7 @@ import java.time.LocalTime; import java.time.ZoneId; import java.time.ZonedDateTime; +import java.time.format.DateTimeParseException; import org.apache.kylin.common.KylinConfig; import org.apache.kylin.junit.annotation.MetadataInfo; @@ -59,6 +61,19 @@ void testLatestScheduledTimeUsesProjectTimeZone() { scheduler.getLatestScheduledTime("00:00:00", zoneId, currentTime)); assertEquals(ZonedDateTime.of(2026, 8, 12, 0, 1, 0, 0, zoneId), scheduler.getLatestScheduledTime("00:01:00", zoneId, currentTime)); + assertThrows(DateTimeParseException.class, + () -> scheduler.getLatestScheduledTime("24:00:00", zoneId, currentTime)); + } + + @Test + void testInitialLookbackUsesDispatcherInterval() { + val scheduler = new AutoBuildSegmentScheduler(); + val currentTime = Instant.parse("2026-08-12T16:00:20Z"); + ReflectionTestUtils.setField(scheduler, "dispatcherIntervalMillis", 300_000L); + + Instant previousTime = ReflectionTestUtils.invokeMethod(scheduler, "getPreviousDispatchTime", currentTime); + + assertEquals(currentTime.minusSeconds(300), previousTime); } @Test @@ -76,8 +91,7 @@ void testDispatchModelOnlyOnceForTriggerWindow() throws Exception { scheduler.dispatch(scheduledTime.minusSeconds(1).toInstant(), scheduledTime.plusSeconds(1).toInstant()); val paramsCaptor = ArgumentCaptor.forClass(IncrementBuildSegmentParams.class); - Mockito.verify(modelBuildService).incrementBuildSegmentsByScheduler(paramsCaptor.capture(), - Mockito.eq("System")); + Mockito.verify(modelBuildService).incrementBuildSegmentsByScheduler(paramsCaptor.capture()); val logicalDate = scheduledTime.toLocalDate().minusDays(1); assertEquals(String.valueOf(LocalDateTime.of(logicalDate, LocalTime.MIDNIGHT).atZone(zoneId).toInstant() .toEpochMilli()), paramsCaptor.getValue().getStart()); @@ -121,6 +135,25 @@ void testSkipOfflineModel() throws Exception { Mockito.verifyNoInteractions(modelBuildService); } + @Test + void testSkipIncompleteConfig() throws Exception { + val modelBuildService = Mockito.mock(ModelBuildService.class); + val scheduler = Mockito.spy(new AutoBuildSegmentScheduler()); + ReflectionTestUtils.setField(scheduler, "modelBuildService", modelBuildService); + Mockito.doReturn(false).when(scheduler).hasRunningModelBuildJob(PROJECT, MODEL_ID); + enableAutoSegmentBuild(); + NDataModelManager.getInstance(KylinConfig.getInstanceFromEnv(), PROJECT).updateDataModel(MODEL_ID, + copyForWrite -> copyForWrite.getSegmentConfig().getAutoSegmentBuild() + .setLogicalDateOffsetDays(null)); + + val zoneId = ZoneId.of(NProjectManager.getInstance(KylinConfig.getInstanceFromEnv()).getProject(PROJECT) + .getConfig().getTimeZone()); + val scheduledTime = ZonedDateTime.of(LocalDate.of(2026, 8, 12), LocalTime.of(1, 0), zoneId); + scheduler.dispatch(scheduledTime.minusSeconds(1).toInstant(), scheduledTime.plusSeconds(1).toInstant()); + + Mockito.verifyNoInteractions(modelBuildService); + } + private NDataModel enableAutoSegmentBuild() { val modelManager = NDataModelManager.getInstance(KylinConfig.getInstanceFromEnv(), PROJECT); return modelManager.updateDataModel(MODEL_ID, copyForWrite -> { @@ -129,7 +162,7 @@ private NDataModel enableAutoSegmentBuild() { config.setTriggerTime("01:00:00"); config.setLogicalDateOffsetDays(1); config.setDataRangeStartTime("00:00:00"); - config.setDataRangeEndTime("24:00:00"); + config.setDataRangeEndTime(AutoSegmentBuildConfig.END_OF_DAY); copyForWrite.getSegmentConfig().setAutoSegmentBuild(config); }); } diff --git a/src/modeling-service/src/main/java/org/apache/kylin/rest/service/ModelBuildService.java b/src/modeling-service/src/main/java/org/apache/kylin/rest/service/ModelBuildService.java index 923fe9d8518..7b4451ecb67 100644 --- a/src/modeling-service/src/main/java/org/apache/kylin/rest/service/ModelBuildService.java +++ b/src/modeling-service/src/main/java/org/apache/kylin/rest/service/ModelBuildService.java @@ -103,6 +103,7 @@ @Component("modelBuildService") public class ModelBuildService extends AbstractModelService implements ModelBuildSupporter { + private static final String SCHEDULER_SUBMITTER = "System"; private static final Logger logger = LoggerFactory.getLogger(ModelBuildService.class); @Autowired private ModelService modelService; @@ -284,9 +285,8 @@ public JobInfoResponse incrementBuildSegmentsManually(IncrementBuildSegmentParam return incrementBuildSegmentsInternal(params, getUsername(), true); } - public JobInfoResponse incrementBuildSegmentsByScheduler(IncrementBuildSegmentParams params, String submitter) - throws Exception { - return incrementBuildSegmentsInternal(params, submitter, false); + public JobInfoResponse incrementBuildSegmentsByScheduler(IncrementBuildSegmentParams params) throws Exception { + return incrementBuildSegmentsInternal(params, SCHEDULER_SUBMITTER, false); } private JobInfoResponse incrementBuildSegmentsInternal(IncrementBuildSegmentParams params, String submitter, diff --git a/src/modeling-service/src/main/java/org/apache/kylin/rest/service/ModelService.java b/src/modeling-service/src/main/java/org/apache/kylin/rest/service/ModelService.java index 3388f4fa7a9..6a7a15ffc67 100644 --- a/src/modeling-service/src/main/java/org/apache/kylin/rest/service/ModelService.java +++ b/src/modeling-service/src/main/java/org/apache/kylin/rest/service/ModelService.java @@ -74,6 +74,7 @@ import java.time.LocalTime; import java.time.format.DateTimeFormatter; import java.time.format.DateTimeParseException; +import java.time.format.ResolverStyle; import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; @@ -291,7 +292,7 @@ public class ModelService extends AbstractModelService implements TableModelSupp private static final List MODEL_CONFIG_BLOCK_LIST = Lists.newArrayList("kylin.index.rule-scheduler-data"); private static final Set STRING_TYPE_SET = Sets.newHashSet("STRING", "CHAR", "VARCHAR"); private static final DateTimeFormatter SEGMENT_AUTO_BUILD_TIME_FORMATTER = DateTimeFormatter.ofPattern("HH:mm:ss", - Locale.ROOT); + Locale.ROOT).withResolverStyle(ResolverStyle.STRICT); //The front-end supports only the following formats private static final List SUPPORTED_FORMATS = ImmutableList.of("ZZ", "DD", "D", "Do", "dddd", "ddd", "dd", // "d", "MMM", "MM", "M", "yyyy", "yy", "hh", "hh", "h", "HH", "H", "m", "mm", "ss", "s", "SSS", "SS", "S", // @@ -3394,7 +3395,9 @@ private void checkAutoSegmentBuildConfig(String project, String modelId, AutoSeg || StringUtils.isBlank(autoSegmentBuild.getDataRangeStartTime()) || StringUtils.isBlank(autoSegmentBuild.getDataRangeEndTime()) || autoSegmentBuild.getLogicalDateOffsetDays() == null) { - throw new KylinException(INVALID_PARAMETER, "Invalid auto_segment_build config."); + throw new KylinException(INVALID_PARAMETER, + "auto_segment_build requires trigger_time, logical_date_offset_days, data_range_start_time " + + "and data_range_end_time."); } if (autoSegmentBuild.getLogicalDateOffsetDays() < 1) { throw new KylinException(INVALID_PARAMETER, "logical_date_offset_days must be >= 1."); diff --git a/src/modeling-service/src/test/java/org/apache/kylin/rest/service/ModelServiceTest.java b/src/modeling-service/src/test/java/org/apache/kylin/rest/service/ModelServiceTest.java index 6185d604a8c..17611630989 100644 --- a/src/modeling-service/src/test/java/org/apache/kylin/rest/service/ModelServiceTest.java +++ b/src/modeling-service/src/test/java/org/apache/kylin/rest/service/ModelServiceTest.java @@ -2794,7 +2794,7 @@ public void testUpdateAndGetModelConfig() { autoSegmentBuildConfig.setTriggerTime("01:00:00"); autoSegmentBuildConfig.setLogicalDateOffsetDays(1); autoSegmentBuildConfig.setDataRangeStartTime("00:00:00"); - autoSegmentBuildConfig.setDataRangeEndTime("24:00:00"); + autoSegmentBuildConfig.setDataRangeEndTime(AutoSegmentBuildConfig.END_OF_DAY); modelConfigRequest.setAutoSegmentBuild(autoSegmentBuildConfig); modelService.updateModelConfig(project, model, modelConfigRequest); @@ -4345,11 +4345,22 @@ public void testCheckModelConfigParameters_AutoSegmentBuildInvalidConfig() { ModelConfigRequest request = new ModelConfigRequest(); AutoSegmentBuildConfig config = new AutoSegmentBuildConfig(); config.setEnabled(true); + request.setAutoSegmentBuild(config); + try { + modelService.checkModelConfigParameters(request); + Assert.fail(); + } catch (Exception e) { + Assert.assertTrue(e instanceof KylinException); + Assert.assertTrue(e.getMessage().contains("trigger_time")); + Assert.assertTrue(e.getMessage().contains("logical_date_offset_days")); + Assert.assertTrue(e.getMessage().contains("data_range_start_time")); + Assert.assertTrue(e.getMessage().contains("data_range_end_time")); + } + config.setTriggerTime("01:00:00"); config.setLogicalDateOffsetDays(0); config.setDataRangeStartTime("00:00:00"); - config.setDataRangeEndTime("24:00:00"); - request.setAutoSegmentBuild(config); + config.setDataRangeEndTime(AutoSegmentBuildConfig.END_OF_DAY); try { modelService.checkModelConfigParameters(request); Assert.fail(); @@ -4368,7 +4379,25 @@ public void testCheckModelConfigParameters_AutoSegmentBuildInvalidConfig() { Assert.assertTrue(e.getMessage().contains("trigger_time")); } + config.setTriggerTime("24:00:00"); + try { + modelService.checkModelConfigParameters(request); + Assert.fail(); + } catch (Exception e) { + Assert.assertTrue(e instanceof KylinException); + Assert.assertTrue(e.getMessage().contains("trigger_time")); + } + config.setTriggerTime("01:00:00"); + config.setDataRangeStartTime("24:00:00"); + try { + modelService.checkModelConfigParameters(request); + Assert.fail(); + } catch (Exception e) { + Assert.assertTrue(e instanceof KylinException); + Assert.assertTrue(e.getMessage().contains("data_range_start_time")); + } + config.setDataRangeStartTime("10:00:00"); config.setDataRangeEndTime("09:00:00"); try { @@ -4378,6 +4407,10 @@ public void testCheckModelConfigParameters_AutoSegmentBuildInvalidConfig() { Assert.assertTrue(e instanceof KylinException); Assert.assertTrue(e.getMessage().contains("data_range_start_time")); } + + config.setDataRangeStartTime("00:00:00"); + config.setDataRangeEndTime(AutoSegmentBuildConfig.END_OF_DAY); + modelService.checkModelConfigParameters(request); } @Test @@ -4388,7 +4421,7 @@ public void testCheckModelConfigParameters_AutoSegmentBuildModelConstraints() { config.setTriggerTime("01:00:00"); config.setLogicalDateOffsetDays(1); config.setDataRangeStartTime("00:00:00"); - config.setDataRangeEndTime("24:00:00"); + config.setDataRangeEndTime(AutoSegmentBuildConfig.END_OF_DAY); request.setAutoSegmentBuild(config); NDataModel streamingModel = NDataModelManager.getInstance(getTestConfig(), "streaming_test").listAllModels() From f5e3f16a0ef4f8d78f2a136ce86cee1a3355d0dc Mon Sep 17 00:00:00 2001 From: ShengHuang Date: Tue, 1 Sep 2026 19:46:02 +0800 Subject: [PATCH 6/6] fix(model): harden auto segment build dispatching Add a project-level dispatcher switch and evaluate the latest due schedule cycle on each scan. Handle covered, overlapping, blocked, and retryable segment build states by target range. --- .../setting/SettingBasic/SettingBasic.vue | 10 + .../setting/SettingBasic/handler.js | 2 + .../setting/SettingBasic/locales.js | 2 + .../rest/request/SegmentConfigRequest.java | 2 + .../rest/response/ProjectConfigResponse.java | 2 + .../kylin/rest/service/ProjectService.java | 8 + .../kylin/metadata/model/SegmentConfig.java | 5 + .../metadata/project/ProjectInstance.java | 2 +- .../cube/model/NSegmentConfigHelperTest.java | 2 +- .../scheduler/AutoBuildSegmentScheduler.java | 287 ++++++++++++++---- .../AutoBuildSegmentSchedulerTest.java | 204 +++++++++++-- .../controller/NProjectControllerTest.java | 1 + .../rest/service/ProjectServiceTest.java | 3 + 13 files changed, 436 insertions(+), 94 deletions(-) diff --git a/kystudio/src/components/setting/SettingBasic/SettingBasic.vue b/kystudio/src/components/setting/SettingBasic/SettingBasic.vue index a0268f50670..7d1dd631c9b 100644 --- a/kystudio/src/components/setting/SettingBasic/SettingBasic.vue +++ b/kystudio/src/components/setting/SettingBasic/SettingBasic.vue @@ -93,6 +93,16 @@ @submit="(scb, ecb) => handleSubmit('segment-settings', scb, ecb)" @cancel="(scb, ecb) => handleResetForm('segment-settings', scb, ecb)"> +
+ {{$t('autoSegmentBuild')}} + + + +
{{$t('autoSegmentBuildDesc')}}
+
{{$t('segmentMerge')}} overrideProps) { diff --git a/src/core-metadata/src/test/java/org/apache/kylin/metadata/cube/model/NSegmentConfigHelperTest.java b/src/core-metadata/src/test/java/org/apache/kylin/metadata/cube/model/NSegmentConfigHelperTest.java index 26a5ed736ed..5b29bb0355b 100644 --- a/src/core-metadata/src/test/java/org/apache/kylin/metadata/cube/model/NSegmentConfigHelperTest.java +++ b/src/core-metadata/src/test/java/org/apache/kylin/metadata/cube/model/NSegmentConfigHelperTest.java @@ -63,7 +63,7 @@ public void testGetSegmentConfig() { // 2. MODEL_BASED && model segmentConfig is not empty, get mergedSegmentConfig of project segmentConfig and model SegmentConfig dataModelManager.updateDataModel(model, copyForWrite -> { copyForWrite.setSegmentConfig( - new SegmentConfig(false, Lists.newArrayList(AutoMergeTimeEnum.WEEK), null, null, false, null)); + new SegmentConfig(false, Lists.newArrayList(AutoMergeTimeEnum.WEEK), null, null, false, null, null)); }); segmentConfig = NSegmentConfigHelper.getModelSegmentConfig(DEFAULT_PROJECT, model); Assert.assertEquals(false, segmentConfig.getAutoMergeEnabled()); diff --git a/src/data-loading-service/src/main/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentScheduler.java b/src/data-loading-service/src/main/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentScheduler.java index 5ff4ad5f1f1..ccb872b6841 100644 --- a/src/data-loading-service/src/main/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentScheduler.java +++ b/src/data-loading-service/src/main/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentScheduler.java @@ -28,14 +28,17 @@ import java.time.format.DateTimeFormatter; import java.time.format.DateTimeParseException; import java.time.format.ResolverStyle; +import java.util.ArrayList; +import java.util.Comparator; import java.util.List; import java.util.Locale; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; -import java.util.concurrent.atomic.AtomicReference; +import java.util.stream.Collectors; import org.apache.commons.lang3.StringUtils; import org.apache.kylin.common.KylinConfig; +import org.apache.kylin.common.util.DateFormat; import org.apache.kylin.common.util.Pair; import org.apache.kylin.guava30.shaded.common.cache.Cache; import org.apache.kylin.guava30.shaded.common.cache.CacheBuilder; @@ -45,17 +48,19 @@ import org.apache.kylin.job.execution.ExecutableState; import org.apache.kylin.job.execution.JobTypeEnum; import org.apache.kylin.job.util.JobContextUtil; +import org.apache.kylin.metadata.cube.model.NDataSegment; import org.apache.kylin.metadata.cube.model.NDataflowManager; import org.apache.kylin.metadata.model.AutoSegmentBuildConfig; import org.apache.kylin.metadata.model.NDataModel; import org.apache.kylin.metadata.model.PartitionDesc; +import org.apache.kylin.metadata.model.SegmentRange; +import org.apache.kylin.metadata.model.SegmentStatusEnum; import org.apache.kylin.metadata.project.NProjectManager; import org.apache.kylin.metadata.project.ProjectInstance; import org.apache.kylin.rest.service.ModelBuildService; import org.apache.kylin.rest.service.params.IncrementBuildSegmentParams; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; -import org.springframework.beans.factory.annotation.Value; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; @@ -65,7 +70,7 @@ @Slf4j @Component public class AutoBuildSegmentScheduler { - private static final long DEFAULT_DISPATCHER_INTERVAL_MILLIS = 30_000L; + private static final int MAX_TRACKED_CYCLES = 10_000; private static final DateTimeFormatter TIME_FORMATTER = DateTimeFormatter.ofPattern("HH:mm:ss", Locale.ROOT) .withResolverStyle(ResolverStyle.STRICT); @@ -73,20 +78,18 @@ public class AutoBuildSegmentScheduler { @Qualifier("modelBuildService") private ModelBuildService modelBuildService; - @Value("${kylin.model.auto-segment-build.dispatcher-interval-ms:30000}") - private long dispatcherIntervalMillis = DEFAULT_DISPATCHER_INTERVAL_MILLIS; - private final AtomicBoolean dispatching = new AtomicBoolean(false); - private final AtomicReference lastDispatchTime = new AtomicReference<>(); + // A completed cycle is re-evaluated after failover/restart, with segment metadata providing idempotency. + private final Cache cycleStates = CacheBuilder.newBuilder().maximumSize(MAX_TRACKED_CYCLES) + .expireAfterAccess(2, TimeUnit.DAYS).build(); // Avoid repeating the same malformed-metadata warning on every dispatcher scan. - private final Cache invalidConfigWarnings = CacheBuilder.newBuilder().maximumSize(10_000) - .expireAfterAccess(1, TimeUnit.DAYS).build(); + private final Cache invalidConfigWarnings = CacheBuilder.newBuilder() + .maximumSize(MAX_TRACKED_CYCLES).expireAfterAccess(1, TimeUnit.DAYS).build(); @Scheduled(fixedDelayString = "${kylin.model.auto-segment-build.dispatcher-interval-ms:30000}") public void schedulerAutoBuildSegment() { val currentTime = Instant.now(); if (!JobContextUtil.getJobContext(KylinConfig.getInstanceFromEnv()).getJobScheduler().isMaster()) { - lastDispatchTime.set(currentTime); return; } if (!dispatching.compareAndSet(false, true)) { @@ -94,43 +97,36 @@ public void schedulerAutoBuildSegment() { return; } - val previousTime = getPreviousDispatchTime(currentTime); try { - dispatch(previousTime, currentTime); + dispatch(currentTime); } finally { - lastDispatchTime.set(currentTime); dispatching.set(false); } } - private Instant getPreviousDispatchTime(Instant currentTime) { - val previousTime = lastDispatchTime.get(); - if (previousTime == null || previousTime.isAfter(currentTime)) { - return currentTime.minusMillis(Math.max(1L, dispatcherIntervalMillis)); - } - return previousTime; - } - - void dispatch(Instant previousTime, Instant currentTime) { + void dispatch(Instant currentTime) { val systemConfig = KylinConfig.readSystemKylinConfig(); val projectManager = NProjectManager.getInstance(systemConfig); for (ProjectInstance project : projectManager.listAllProjects()) { try { - dispatchProject(systemConfig, project, previousTime, currentTime); + dispatchProject(systemConfig, project, currentTime); } catch (Exception e) { log.error("Auto build segment dispatch failed for project: {}", project.getName(), e); } } } - private void dispatchProject(KylinConfig systemConfig, ProjectInstance project, Instant previousTime, - Instant currentTime) { + private void dispatchProject(KylinConfig systemConfig, ProjectInstance project, Instant currentTime) { + if (project.getSegmentConfig() == null + || !Boolean.TRUE.equals(project.getSegmentConfig().getAutoSegmentBuildEnabled())) { + return; + } val projectName = project.getName(); val zoneId = ZoneId.of(project.getConfig().getTimeZone()); val dataflowManager = NDataflowManager.getInstance(systemConfig, projectName); - for (NDataModel model : dataflowManager.listOnlineDataModels()) { + for (NDataModel model : dataflowManager.listUnderliningDataModels()) { try { - dispatchModel(projectName, model, zoneId, previousTime, currentTime); + dispatchModel(projectName, model, zoneId, currentTime); } catch (Exception e) { log.error("Auto build segment dispatch failed, project: {}, model: {}", projectName, model.getUuid(), e); @@ -138,19 +134,17 @@ private void dispatchProject(KylinConfig systemConfig, ProjectInstance project, } } - private void dispatchModel(String project, NDataModel model, ZoneId zoneId, Instant previousTime, - Instant currentTime) { + private void dispatchModel(String project, NDataModel model, ZoneId zoneId, Instant currentTime) { val autoSegmentBuild = getEligibleConfig(project, model); if (autoSegmentBuild == null) { return; } val scheduledTime = getLatestScheduledTime(autoSegmentBuild.getTriggerTime(), zoneId, currentTime); - val scheduledInstant = scheduledTime.toInstant(); - if (!scheduledInstant.isAfter(previousTime) || scheduledInstant.isAfter(currentTime)) { + if (scheduledTime.toInstant().isAfter(currentTime)) { return; } - submitJob(project, model, autoSegmentBuild, scheduledTime); + dispatchTarget(project, model, autoSegmentBuild, scheduledTime); } private AutoSegmentBuildConfig getEligibleConfig(String project, NDataModel model) { @@ -234,49 +228,214 @@ ZonedDateTime getLatestScheduledTime(String triggerTime, ZoneId zoneId, Instant return scheduledTime; } - void submitJob(String project, NDataModel model, AutoSegmentBuildConfig autoSegmentBuild, + private void dispatchTarget(String project, NDataModel model, AutoSegmentBuildConfig autoSegmentBuild, ZonedDateTime scheduledTime) { val modelId = model.getUuid(); - if (hasRunningModelBuildJob(project, modelId)) { - log.info("Skip auto build segment because the model has a progressing build job, project: {}, model: {}", - project, modelId); + val target = createBuildTarget(model, autoSegmentBuild, scheduledTime); + val cycleKey = createCycleKey(project, modelId, scheduledTime, target.range); + val previousState = cycleStates.getIfPresent(cycleKey); + if (previousState != null && previousState.isComplete()) { return; } - try { - val zoneId = scheduledTime.getZone(); - val logicalDate = scheduledTime.toLocalDate().minusDays(autoSegmentBuild.getLogicalDateOffsetDays()); - val start = parseTime(autoSegmentBuild.getDataRangeStartTime(), false); - val end = parseTime(autoSegmentBuild.getDataRangeEndTime(), true); - val startDateTime = LocalDateTime.of(logicalDate, start.getFirst()); - val endDate = end.getSecond() ? logicalDate.plusDays(1) : logicalDate; - val endDateTime = LocalDateTime.of(endDate, end.getFirst()); - if (!startDateTime.isBefore(endDateTime)) { - log.warn("Skip auto build segment because of an invalid range, project: {}, model: {}", project, - modelId); - return; + + val dataflow = NDataflowManager.getInstance(KylinConfig.getInstanceFromEnv(), project).getDataflow(modelId); + List segments = dataflow == null ? Lists.newArrayList() : dataflow.getSegments(); + val jobs = getModelBuildJobs(project, modelId); + TargetState state = evaluateTargetState(segments, jobs, target.range); + if (state == TargetState.NO_OVERLAP) { + try { + submitJob(project, model, target); + state = TargetState.SUBMITTED; + } catch (Exception e) { + state = TargetState.RETRYABLE_FAILURE; + if (previousState != TargetState.RETRYABLE_FAILURE) { + log.error("Auto build segment submit failed, project: {}, model: {}, range: {}", project, modelId, + target.range, e); + } else { + log.debug("Auto build segment submit still failing, project: {}, model: {}, range: {}", project, + modelId, target.range, e); + } } - val startMillis = String.valueOf(startDateTime.atZone(zoneId).toInstant().toEpochMilli()); - val endMillis = String.valueOf(endDateTime.atZone(zoneId).toInstant().toEpochMilli()); - val params = new IncrementBuildSegmentParams(project, modelId, startMillis, endMillis, - model.getPartitionDesc(), model.getMultiPartitionDesc(), Lists.newArrayList(), true, null); - modelBuildService.incrementBuildSegmentsByScheduler(params); - log.info("Auto build segment submitted, project: {}, model: {}, scheduled time: {}, range: [{}, {})", - project, modelId, scheduledTime, startMillis, endMillis); - } catch (DateTimeParseException e) { - log.error("Invalid time in auto build segment config, project: {}, model: {}", project, modelId, e); - } catch (Exception e) { - log.error("Auto build segment submit failed, project: {}, model: {}", project, modelId, e); } + recordCycleState(cycleKey, previousState, state, project, modelId, scheduledTime, target.range, segments); } - boolean hasRunningModelBuildJob(String project, String modelId) { - // Segment and model metadata mutations are serialized with any progressing build job on the same model. + private BuildTarget createBuildTarget(NDataModel model, AutoSegmentBuildConfig autoSegmentBuild, + ZonedDateTime scheduledTime) { + val zoneId = scheduledTime.getZone(); + val logicalDate = scheduledTime.toLocalDate().minusDays(autoSegmentBuild.getLogicalDateOffsetDays()); + val start = parseTime(autoSegmentBuild.getDataRangeStartTime(), false); + val end = parseTime(autoSegmentBuild.getDataRangeEndTime(), true); + val startDateTime = LocalDateTime.of(logicalDate, start.getFirst()); + val endDate = end.getSecond() ? logicalDate.plusDays(1) : logicalDate; + val endDateTime = LocalDateTime.of(endDate, end.getFirst()); + if (!startDateTime.isBefore(endDateTime)) { + throw new IllegalArgumentException("Auto build segment range start must be earlier than end"); + } + + val startMillis = String.valueOf(startDateTime.atZone(zoneId).toInstant().toEpochMilli()); + val endMillis = String.valueOf(endDateTime.atZone(zoneId).toInstant().toEpochMilli()); + val partitionDateFormat = model.getPartitionDesc().getPartitionDateFormat(); + val normalizedStart = DateFormat.getFormatTimeStamp(startMillis, partitionDateFormat); + val normalizedEnd = DateFormat.getFormatTimeStamp(endMillis, partitionDateFormat); + if (normalizedStart >= normalizedEnd) { + throw new IllegalArgumentException( + "Auto build segment range is empty after partition format normalization"); + } + return new BuildTarget(startMillis, endMillis, + new SegmentRange.TimePartitionedSegmentRange(normalizedStart, normalizedEnd)); + } + + private String createCycleKey(String project, String modelId, ZonedDateTime scheduledTime, + SegmentRange targetRange) { + return project + "/" + modelId + "/" + scheduledTime.toInstant().toEpochMilli() + "/" + + targetRange.getStart() + "-" + targetRange.getEnd(); + } + + private void submitJob(String project, NDataModel model, BuildTarget target) throws Exception { + val params = new IncrementBuildSegmentParams(project, model.getUuid(), target.startMillis, target.endMillis, + model.getPartitionDesc(), model.getMultiPartitionDesc(), Lists.newArrayList(), true, null); + modelBuildService.incrementBuildSegmentsByScheduler(params); + } + + List getModelBuildJobs(String project, String modelId) { val executableManager = ExecutableManager.getInstance(KylinConfig.getInstanceFromEnv(), project); JobTypeEnum[] buildJobTypes = JobTypeEnum.getJobTypeByCategory(JobTypeEnum.Category.BUILD) .toArray(new JobTypeEnum[0]); - List jobs = executableManager.listExecByModelAndStatus(modelId, - ExecutableState::isProgressing, buildJobTypes); - return !jobs.isEmpty(); + return executableManager.listExecByModelAndStatus(modelId, ExecutableState::isRunning, buildJobTypes); + } + + TargetState evaluateTargetState(List segments, List jobs, + SegmentRange targetRange) { + val coverageState = getCoverageState(segments, targetRange); + if (coverageState != null) { + return coverageState; + } + + List buildingSegments = segments.stream() + .filter(segment -> segment.getStatus() == SegmentStatusEnum.NEW) + .filter(segment -> segment.getSegRange().overlaps(targetRange)).collect(Collectors.toList()); + if (!buildingSegments.isEmpty()) { + boolean hasProgressingJob = false; + for (NDataSegment segment : buildingSegments) { + boolean hasRelatedJob = false; + for (AbstractExecutable job : jobs) { + if (job.getTargetSegments() == null || !job.getTargetSegments().contains(segment.getId())) { + continue; + } + val jobState = job.getStatusInMem(); + if (!jobState.isRunning()) { + continue; + } + hasRelatedJob = true; + if (!jobState.isProgressing()) { + return TargetState.BLOCKED; + } + hasProgressingJob = true; + } + if (!hasRelatedJob) { + return TargetState.BLOCKED; + } + } + return hasProgressingJob ? TargetState.IN_PROGRESS : TargetState.BLOCKED; + } + + boolean hasAvailableOverlap = segments.stream() + .filter(segment -> segment.getStatus() == SegmentStatusEnum.READY + || segment.getStatus() == SegmentStatusEnum.WARNING) + .anyMatch(segment -> segment.getSegRange().overlaps(targetRange)); + return hasAvailableOverlap ? TargetState.PARTIAL_OVERLAP : TargetState.NO_OVERLAP; + } + + private TargetState getCoverageState(List segments, SegmentRange targetRange) { + List availableSegments = new ArrayList<>(); + for (NDataSegment segment : segments) { + if ((segment.getStatus() == SegmentStatusEnum.READY || segment.getStatus() == SegmentStatusEnum.WARNING) + && segment.getSegRange().overlaps(targetRange)) { + availableSegments.add(segment); + } + } + availableSegments.sort(Comparator.comparingLong(segment -> (Long) segment.getSegRange().getStart())); + + long coveredUntil = targetRange.getStart(); + boolean includesWarning = false; + for (NDataSegment segment : availableSegments) { + long segmentStart = (Long) segment.getSegRange().getStart(); + long segmentEnd = (Long) segment.getSegRange().getEnd(); + if (segmentEnd <= coveredUntil) { + continue; + } + if (segmentStart > coveredUntil) { + break; + } + coveredUntil = Math.max(coveredUntil, segmentEnd); + includesWarning |= segment.getStatus() == SegmentStatusEnum.WARNING; + if (coveredUntil >= targetRange.getEnd()) { + return includesWarning ? TargetState.COVERED_WITH_WARNING : TargetState.COVERED; + } + } + return null; + } + + private void recordCycleState(String cycleKey, TargetState previousState, TargetState state, String project, + String modelId, ZonedDateTime scheduledTime, SegmentRange targetRange, + List segments) { + cycleStates.put(cycleKey, state); + if (state == previousState) { + return; + } + + val relatedSegments = segments.stream().filter(segment -> segment.getSegRange().overlaps(targetRange)) + .map(segment -> segment.getId() + ":" + segment.getStatus() + segment.getSegRange()) + .collect(Collectors.joining(",", "[", "]")); + switch (state) { + case COVERED: + log.info("Auto build segment already covered, project: {}, model: {}, range: {}", project, modelId, + targetRange); + break; + case COVERED_WITH_WARNING: + log.warn("Auto build segment covered by warning segment, project: {}, model: {}, range: {}, segments: {}", + project, modelId, targetRange, relatedSegments); + break; + case IN_PROGRESS: + log.info("Auto build segment waiting for overlapping build job, project: {}, model: {}, range: {}, " + + "segments: {}", project, modelId, targetRange, relatedSegments); + break; + case BLOCKED: + log.warn("Auto build segment blocked by non-progressing job or orphan segment, project: {}, model: {}, " + + "range: {}, segments: {}", project, modelId, targetRange, relatedSegments); + break; + case PARTIAL_OVERLAP: + log.warn("Auto build segment has a partial overlap, project: {}, model: {}, range: {}, segments: {}", + project, modelId, targetRange, relatedSegments); + break; + case SUBMITTED: + log.info("Auto build segment submitted, project: {}, model: {}, scheduled time: {}, range: {}", project, + modelId, scheduledTime, targetRange); + break; + default: + break; + } + } + + private static class BuildTarget { + private final String startMillis; + private final String endMillis; + private final SegmentRange range; + + BuildTarget(String startMillis, String endMillis, SegmentRange range) { + this.startMillis = startMillis; + this.endMillis = endMillis; + this.range = range; + } + } + + enum TargetState { + NO_OVERLAP, COVERED, COVERED_WITH_WARNING, IN_PROGRESS, BLOCKED, PARTIAL_OVERLAP, SUBMITTED, RETRYABLE_FAILURE; + + boolean isComplete() { + return this == COVERED || this == COVERED_WITH_WARNING; + } } private Pair parseTime(String time, boolean allow24Hour) { diff --git a/src/data-loading-service/src/test/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentSchedulerTest.java b/src/data-loading-service/src/test/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentSchedulerTest.java index 07619d63696..c8deff1348a 100644 --- a/src/data-loading-service/src/test/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentSchedulerTest.java +++ b/src/data-loading-service/src/test/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentSchedulerTest.java @@ -28,13 +28,21 @@ import java.time.ZoneId; import java.time.ZonedDateTime; import java.time.format.DateTimeParseException; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; import org.apache.kylin.common.KylinConfig; +import org.apache.kylin.job.execution.AbstractExecutable; +import org.apache.kylin.job.execution.ExecutableState; import org.apache.kylin.junit.annotation.MetadataInfo; +import org.apache.kylin.metadata.cube.model.NDataSegment; import org.apache.kylin.metadata.cube.model.NDataflowManager; import org.apache.kylin.metadata.model.AutoSegmentBuildConfig; import org.apache.kylin.metadata.model.NDataModel; import org.apache.kylin.metadata.model.NDataModelManager; +import org.apache.kylin.metadata.model.SegmentRange; +import org.apache.kylin.metadata.model.SegmentStatusEnum; import org.apache.kylin.metadata.project.NProjectManager; import org.apache.kylin.metadata.realization.RealizationStatusEnum; import org.apache.kylin.rest.service.ModelBuildService; @@ -66,29 +74,20 @@ void testLatestScheduledTimeUsesProjectTimeZone() { } @Test - void testInitialLookbackUsesDispatcherInterval() { - val scheduler = new AutoBuildSegmentScheduler(); - val currentTime = Instant.parse("2026-08-12T16:00:20Z"); - ReflectionTestUtils.setField(scheduler, "dispatcherIntervalMillis", 300_000L); - - Instant previousTime = ReflectionTestUtils.invokeMethod(scheduler, "getPreviousDispatchTime", currentTime); - - assertEquals(currentTime.minusSeconds(300), previousTime); - } - - @Test - void testDispatchModelOnlyOnceForTriggerWindow() throws Exception { + void testDispatchCatchesUpLatestCycleAndBuildsConfiguredRange() throws Exception { val modelBuildService = Mockito.mock(ModelBuildService.class); val scheduler = Mockito.spy(new AutoBuildSegmentScheduler()); ReflectionTestUtils.setField(scheduler, "modelBuildService", modelBuildService); - Mockito.doReturn(false).when(scheduler).hasRunningModelBuildJob(PROJECT, MODEL_ID); + Mockito.doReturn(Collections.emptyList()).when(scheduler).getModelBuildJobs(PROJECT, MODEL_ID); + Mockito.doReturn(AutoBuildSegmentScheduler.TargetState.NO_OVERLAP).when(scheduler) + .evaluateTargetState(Mockito.anyList(), Mockito.anyList(), Mockito.any()); enableAutoSegmentBuild(); val projectConfig = NProjectManager.getInstance(KylinConfig.getInstanceFromEnv()).getProject(PROJECT) .getConfig(); val zoneId = ZoneId.of(projectConfig.getTimeZone()); val scheduledTime = ZonedDateTime.of(LocalDate.of(2026, 8, 12), LocalTime.of(1, 0), zoneId); - scheduler.dispatch(scheduledTime.minusSeconds(1).toInstant(), scheduledTime.plusSeconds(1).toInstant()); + scheduler.dispatch(scheduledTime.plusHours(6).toInstant()); val paramsCaptor = ArgumentCaptor.forClass(IncrementBuildSegmentParams.class); Mockito.verify(modelBuildService).incrementBuildSegmentsByScheduler(paramsCaptor.capture()); @@ -97,32 +96,36 @@ void testDispatchModelOnlyOnceForTriggerWindow() throws Exception { .toEpochMilli()), paramsCaptor.getValue().getStart()); assertEquals(String.valueOf(LocalDateTime.of(logicalDate.plusDays(1), LocalTime.MIDNIGHT).atZone(zoneId) .toInstant().toEpochMilli()), paramsCaptor.getValue().getEnd()); - - Mockito.clearInvocations(modelBuildService); - scheduler.dispatch(scheduledTime.toInstant(), scheduledTime.plusSeconds(30).toInstant()); - Mockito.verifyNoInteractions(modelBuildService); } @Test - void testSkipWhenModelHasProgressingBuildJob() throws Exception { + void testCompletedCycleIsNotEvaluatedAgain() throws Exception { val modelBuildService = Mockito.mock(ModelBuildService.class); val scheduler = Mockito.spy(new AutoBuildSegmentScheduler()); ReflectionTestUtils.setField(scheduler, "modelBuildService", modelBuildService); - Mockito.doReturn(true).when(scheduler).hasRunningModelBuildJob(PROJECT, MODEL_ID); - val model = enableAutoSegmentBuild(); - val config = model.getSegmentConfig().getAutoSegmentBuild(); + Mockito.doReturn(Collections.emptyList()).when(scheduler).getModelBuildJobs(PROJECT, MODEL_ID); + Mockito.doReturn(AutoBuildSegmentScheduler.TargetState.COVERED).when(scheduler) + .evaluateTargetState(Mockito.anyList(), Mockito.anyList(), Mockito.any()); + enableAutoSegmentBuild(); - scheduler.submitJob(PROJECT, model, config, ZonedDateTime.of(2026, 8, 12, 1, 0, 0, 0, ZoneId.of("UTC"))); + val zoneId = ZoneId.of(NProjectManager.getInstance(KylinConfig.getInstanceFromEnv()).getProject(PROJECT) + .getConfig().getTimeZone()); + val scheduledTime = ZonedDateTime.of(2026, 8, 12, 1, 0, 0, 0, zoneId); + scheduler.dispatch(scheduledTime.plusMinutes(1).toInstant()); + scheduler.dispatch(scheduledTime.plusMinutes(2).toInstant()); Mockito.verifyNoInteractions(modelBuildService); + Mockito.verify(scheduler).evaluateTargetState(Mockito.anyList(), Mockito.anyList(), Mockito.any()); } @Test - void testSkipOfflineModel() throws Exception { + void testBuildOfflineModel() throws Exception { val modelBuildService = Mockito.mock(ModelBuildService.class); val scheduler = Mockito.spy(new AutoBuildSegmentScheduler()); ReflectionTestUtils.setField(scheduler, "modelBuildService", modelBuildService); - Mockito.doReturn(false).when(scheduler).hasRunningModelBuildJob(PROJECT, MODEL_ID); + Mockito.doReturn(Collections.emptyList()).when(scheduler).getModelBuildJobs(PROJECT, MODEL_ID); + Mockito.doReturn(AutoBuildSegmentScheduler.TargetState.NO_OVERLAP).when(scheduler) + .evaluateTargetState(Mockito.anyList(), Mockito.anyList(), Mockito.any()); enableAutoSegmentBuild(); NDataflowManager.getInstance(KylinConfig.getInstanceFromEnv(), PROJECT).updateDataflowStatus(MODEL_ID, RealizationStatusEnum.OFFLINE); @@ -130,7 +133,23 @@ void testSkipOfflineModel() throws Exception { val zoneId = ZoneId.of(NProjectManager.getInstance(KylinConfig.getInstanceFromEnv()).getProject(PROJECT) .getConfig().getTimeZone()); val scheduledTime = ZonedDateTime.of(LocalDate.of(2026, 8, 12), LocalTime.of(1, 0), zoneId); - scheduler.dispatch(scheduledTime.minusSeconds(1).toInstant(), scheduledTime.plusSeconds(1).toInstant()); + scheduler.dispatch(scheduledTime.plusSeconds(1).toInstant()); + + Mockito.verify(modelBuildService).incrementBuildSegmentsByScheduler(Mockito.any()); + } + + @Test + void testSkipProjectWhenAutoSegmentBuildDisabled() throws Exception { + val modelBuildService = Mockito.mock(ModelBuildService.class); + val scheduler = Mockito.spy(new AutoBuildSegmentScheduler()); + ReflectionTestUtils.setField(scheduler, "modelBuildService", modelBuildService); + enableAutoSegmentBuild(); + setProjectAutoSegmentBuildEnabled(false); + + val zoneId = ZoneId.of(NProjectManager.getInstance(KylinConfig.getInstanceFromEnv()).getProject(PROJECT) + .getConfig().getTimeZone()); + val scheduledTime = ZonedDateTime.of(LocalDate.of(2026, 8, 12), LocalTime.of(1, 0), zoneId); + scheduler.dispatch(scheduledTime.plusSeconds(1).toInstant()); Mockito.verifyNoInteractions(modelBuildService); } @@ -140,7 +159,6 @@ void testSkipIncompleteConfig() throws Exception { val modelBuildService = Mockito.mock(ModelBuildService.class); val scheduler = Mockito.spy(new AutoBuildSegmentScheduler()); ReflectionTestUtils.setField(scheduler, "modelBuildService", modelBuildService); - Mockito.doReturn(false).when(scheduler).hasRunningModelBuildJob(PROJECT, MODEL_ID); enableAutoSegmentBuild(); NDataModelManager.getInstance(KylinConfig.getInstanceFromEnv(), PROJECT).updateDataModel(MODEL_ID, copyForWrite -> copyForWrite.getSegmentConfig().getAutoSegmentBuild() @@ -149,12 +167,137 @@ void testSkipIncompleteConfig() throws Exception { val zoneId = ZoneId.of(NProjectManager.getInstance(KylinConfig.getInstanceFromEnv()).getProject(PROJECT) .getConfig().getTimeZone()); val scheduledTime = ZonedDateTime.of(LocalDate.of(2026, 8, 12), LocalTime.of(1, 0), zoneId); - scheduler.dispatch(scheduledTime.minusSeconds(1).toInstant(), scheduledTime.plusSeconds(1).toInstant()); + scheduler.dispatch(scheduledTime.plusSeconds(1).toInstant()); Mockito.verifyNoInteractions(modelBuildService); } + @Test + void testTargetCoveredByExistingSegments() { + val scheduler = new AutoBuildSegmentScheduler(); + val target = range(100L, 200L); + + assertEquals(AutoBuildSegmentScheduler.TargetState.COVERED, + scheduler.evaluateTargetState(Collections.singletonList(segment("superset", 50L, 250L, + SegmentStatusEnum.READY)), Collections.emptyList(), target)); + assertEquals(AutoBuildSegmentScheduler.TargetState.COVERED, + scheduler.evaluateTargetState(Arrays.asList(segment("first", 100L, 150L, SegmentStatusEnum.READY), + segment("second", 150L, 200L, SegmentStatusEnum.READY)), Collections.emptyList(), target)); + assertEquals(AutoBuildSegmentScheduler.TargetState.COVERED_WITH_WARNING, + scheduler.evaluateTargetState(Collections.singletonList(segment("warning", 100L, 200L, + SegmentStatusEnum.WARNING)), Collections.emptyList(), target)); + } + + @Test + void testProgressingJobIsCheckedByRange() { + val scheduler = new AutoBuildSegmentScheduler(); + val target = range(100L, 200L); + + NDataSegment larger = segment("larger", 50L, 250L, SegmentStatusEnum.NEW); + assertEquals(AutoBuildSegmentScheduler.TargetState.IN_PROGRESS, scheduler.evaluateTargetState( + Collections.singletonList(larger), Collections.singletonList(job(ExecutableState.RUNNING, "larger")), + target)); + + NDataSegment smaller = segment("smaller", 100L, 150L, SegmentStatusEnum.NEW); + assertEquals(AutoBuildSegmentScheduler.TargetState.IN_PROGRESS, scheduler.evaluateTargetState( + Collections.singletonList(smaller), Collections.singletonList(job(ExecutableState.PENDING, "smaller")), + target)); + + NDataSegment disjoint = segment("disjoint", 200L, 300L, SegmentStatusEnum.NEW); + assertEquals(AutoBuildSegmentScheduler.TargetState.NO_OVERLAP, scheduler.evaluateTargetState( + Collections.singletonList(disjoint), + Collections.singletonList(job(ExecutableState.RUNNING, "disjoint")), target)); + } + + @Test + void testErrorPausedAndOrphanSegmentsBlockTarget() { + val scheduler = new AutoBuildSegmentScheduler(); + val target = range(100L, 200L); + NDataSegment building = segment("building", 100L, 200L, SegmentStatusEnum.NEW); + + assertEquals(AutoBuildSegmentScheduler.TargetState.BLOCKED, scheduler.evaluateTargetState( + Collections.singletonList(building), Collections.singletonList(job(ExecutableState.ERROR, "building")), + target)); + assertEquals(AutoBuildSegmentScheduler.TargetState.BLOCKED, scheduler.evaluateTargetState( + Collections.singletonList(building), Collections.singletonList(job(ExecutableState.PAUSED, "building")), + target)); + assertEquals(AutoBuildSegmentScheduler.TargetState.BLOCKED, scheduler.evaluateTargetState( + Collections.singletonList(building), Collections.emptyList(), target)); + } + + @Test + void testPartialAvailableSegmentIsConflict() { + val scheduler = new AutoBuildSegmentScheduler(); + val target = range(100L, 200L); + List segments = Collections + .singletonList(segment("partial", 100L, 150L, SegmentStatusEnum.READY)); + + assertEquals(AutoBuildSegmentScheduler.TargetState.PARTIAL_OVERLAP, + scheduler.evaluateTargetState(segments, Collections.emptyList(), target)); + } + + @Test + void testInProgressCycleIsReevaluated() throws Exception { + val modelBuildService = Mockito.mock(ModelBuildService.class); + val scheduler = Mockito.spy(new AutoBuildSegmentScheduler()); + ReflectionTestUtils.setField(scheduler, "modelBuildService", modelBuildService); + Mockito.doReturn(Collections.emptyList()).when(scheduler).getModelBuildJobs(PROJECT, MODEL_ID); + Mockito.doReturn(AutoBuildSegmentScheduler.TargetState.IN_PROGRESS, + AutoBuildSegmentScheduler.TargetState.NO_OVERLAP).when(scheduler) + .evaluateTargetState(Mockito.anyList(), Mockito.anyList(), Mockito.any()); + enableAutoSegmentBuild(); + + val zoneId = ZoneId.of(NProjectManager.getInstance(KylinConfig.getInstanceFromEnv()).getProject(PROJECT) + .getConfig().getTimeZone()); + val scheduledTime = ZonedDateTime.of(LocalDate.of(2026, 8, 12), LocalTime.of(1, 0), zoneId); + scheduler.dispatch(scheduledTime.plusMinutes(1).toInstant()); + scheduler.dispatch(scheduledTime.plusMinutes(2).toInstant()); + + Mockito.verify(modelBuildService).incrementBuildSegmentsByScheduler(Mockito.any()); + } + + @Test + void testRetryableSubmissionFailureIsReevaluated() throws Exception { + val modelBuildService = Mockito.mock(ModelBuildService.class); + val scheduler = Mockito.spy(new AutoBuildSegmentScheduler()); + ReflectionTestUtils.setField(scheduler, "modelBuildService", modelBuildService); + Mockito.doReturn(Collections.emptyList()).when(scheduler).getModelBuildJobs(PROJECT, MODEL_ID); + Mockito.doReturn(AutoBuildSegmentScheduler.TargetState.NO_OVERLAP).when(scheduler) + .evaluateTargetState(Mockito.anyList(), Mockito.anyList(), Mockito.any()); + Mockito.doThrow(new IllegalStateException("transient failure")).doReturn(null).when(modelBuildService) + .incrementBuildSegmentsByScheduler(Mockito.any()); + enableAutoSegmentBuild(); + + val zoneId = ZoneId.of(NProjectManager.getInstance(KylinConfig.getInstanceFromEnv()).getProject(PROJECT) + .getConfig().getTimeZone()); + val scheduledTime = ZonedDateTime.of(LocalDate.of(2026, 8, 12), LocalTime.of(1, 0), zoneId); + scheduler.dispatch(scheduledTime.plusMinutes(1).toInstant()); + scheduler.dispatch(scheduledTime.plusMinutes(2).toInstant()); + + Mockito.verify(modelBuildService, Mockito.times(2)).incrementBuildSegmentsByScheduler(Mockito.any()); + } + + private SegmentRange range(long start, long end) { + return new SegmentRange.TimePartitionedSegmentRange(start, end); + } + + private NDataSegment segment(String id, long start, long end, SegmentStatusEnum status) { + NDataSegment segment = Mockito.mock(NDataSegment.class); + Mockito.when(segment.getId()).thenReturn(id); + Mockito.when(segment.getStatus()).thenReturn(status); + Mockito.when(segment.getSegRange()).thenReturn(range(start, end)); + return segment; + } + + private AbstractExecutable job(ExecutableState state, String... segmentIds) { + AbstractExecutable job = Mockito.mock(AbstractExecutable.class); + Mockito.when(job.getStatusInMem()).thenReturn(state); + Mockito.when(job.getTargetSegments()).thenReturn(Arrays.asList(segmentIds)); + return job; + } + private NDataModel enableAutoSegmentBuild() { + setProjectAutoSegmentBuildEnabled(true); val modelManager = NDataModelManager.getInstance(KylinConfig.getInstanceFromEnv(), PROJECT); return modelManager.updateDataModel(MODEL_ID, copyForWrite -> { AutoSegmentBuildConfig config = new AutoSegmentBuildConfig(); @@ -166,4 +309,9 @@ private NDataModel enableAutoSegmentBuild() { copyForWrite.getSegmentConfig().setAutoSegmentBuild(config); }); } + + private void setProjectAutoSegmentBuildEnabled(boolean enabled) { + NProjectManager.getInstance(KylinConfig.getInstanceFromEnv()).updateProject(PROJECT, + copyForWrite -> copyForWrite.getSegmentConfig().setAutoSegmentBuildEnabled(enabled)); + } } diff --git a/src/metadata-server/src/test/java/org/apache/kylin/rest/controller/NProjectControllerTest.java b/src/metadata-server/src/test/java/org/apache/kylin/rest/controller/NProjectControllerTest.java index 996fd65ffe5..cd04cf4f3fa 100644 --- a/src/metadata-server/src/test/java/org/apache/kylin/rest/controller/NProjectControllerTest.java +++ b/src/metadata-server/src/test/java/org/apache/kylin/rest/controller/NProjectControllerTest.java @@ -328,6 +328,7 @@ public void testUpdateSegmentConfig() throws Exception { request.setAutoMergeEnabled(true); request.setAutoMergeTimeRanges(Arrays.asList(AutoMergeTimeEnum.DAY)); request.setCreateEmptySegmentEnabled(true); + request.setAutoSegmentBuildEnabled(true); Mockito.doNothing().when(projectService).updateSegmentConfig("default", request); mockMvc.perform(MockMvcRequestBuilders.put("/api/projects/{project}/segment_config", "default") diff --git a/src/modeling-service/src/test/java/org/apache/kylin/rest/service/ProjectServiceTest.java b/src/modeling-service/src/test/java/org/apache/kylin/rest/service/ProjectServiceTest.java index 4745d6c63ed..6cc4e4a2c55 100644 --- a/src/modeling-service/src/test/java/org/apache/kylin/rest/service/ProjectServiceTest.java +++ b/src/modeling-service/src/test/java/org/apache/kylin/rest/service/ProjectServiceTest.java @@ -474,9 +474,11 @@ public void testUpdateProjectConfig() throws IOException { val segmentConfigRequest = new SegmentConfigRequest(); segmentConfigRequest.setAutoMergeEnabled(false); segmentConfigRequest.setAutoMergeTimeRanges(Collections.singletonList(AutoMergeTimeEnum.DAY)); + segmentConfigRequest.setAutoSegmentBuildEnabled(true); projectService.updateSegmentConfig(project, segmentConfigRequest); response = projectService.getProjectConfig(project); Assert.assertFalse(response.isAutoMergeEnabled()); + Assert.assertTrue(response.isAutoSegmentBuildEnabled()); val pushDownConfigRequest = new PushDownConfigRequest(); pushDownConfigRequest.setPushDownEnabled(false); @@ -835,6 +837,7 @@ public void testResetProjectConfig() { response = projectService.resetProjectConfig(PROJECT, "segment_config"); Assert.assertFalse(response.isAutoMergeEnabled()); + Assert.assertFalse(response.isAutoSegmentBuildEnabled()); Assert.assertEquals(4, response.getAutoMergeTimeRanges().size()); response = projectService.resetProjectConfig(PROJECT, "storage_quota_config");