diff --git a/kystudio/src/components/setting/SettingModel/SettingModel.vue b/kystudio/src/components/setting/SettingModel/SettingModel.vue
index 9e3725d3762..29c02e63023 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)}}
+
+
+
+
+
+
@@ -141,6 +150,38 @@
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
G
@@ -196,12 +237,15 @@ import { handleSuccess, transToGmtTime, kylinConfirm } from '../../../util/busin
import { handleSuccessAsync, handleError, objectClone, ArrayFlat } from '../../../util/index'
import { retentionTypes } from '../handler'
+const END_OF_DAY = '24:00:00'
+
const initialSettingForm = JSON.stringify({
name: '',
settingItem: '',
autoMerge: [],
volatileRange: {volatile_range_number: 0, volatile_range_type: '', volatile_range_enabled: true},
retentionThreshold: {retention_range_number: 0, retention_range_type: '', retention_range_enabled: true},
+ autoSegmentBuild: {enabled: true, trigger_time: '01:00:00', logical_date_offset_days: 1, data_range_start_time: '00:00:00', data_range_end_time: END_OF_DAY},
'spark.executor.cores': null,
'spark.executor.instances': null,
'spark.executor.memory': null,
@@ -255,6 +299,7 @@ export default class SettingStorage extends Vue {
'Auto-merge': 'Auto-merge',
'Volatile Range': 'Volatile Range',
'Retention Threshold': 'Retention Threshold',
+ 'Auto Segment Build': 'Auto Segment Build',
'spark.executor.cores': 'kylin.engine.spark-conf.spark.executor.cores',
'spark.executor.instances': 'kylin.engine.spark-conf.spark.executor.instances',
'spark.executor.memory': 'kylin.engine.spark-conf.spark.executor.memory',
@@ -275,6 +320,7 @@ export default class SettingStorage extends Vue {
'Auto-merge',
'Volatile Range',
'Retention Threshold',
+ 'Auto Segment Build',
'spark.executor.cores',
'spark.executor.instances',
'spark.executor.memory',
@@ -308,6 +354,7 @@ export default class SettingStorage extends Vue {
'Auto-merge': this.$t('autoMergeTip'),
'Volatile Range': this.$t('volatileTip'),
'Retention Threshold': this.$t('retentionThresholdDesc'),
+ 'Auto Segment Build': this.$t('autoSegmentBuildTip'),
'kylin.engine.spark-conf.spark.executor.cores': this.$t('sparkCores'),
'kylin.engine.spark-conf.spark.executor.instances': this.$t('sparkInstances'),
'kylin.engine.spark-conf.spark.executor.memory': this.$t('sparkMemory'),
@@ -346,6 +393,8 @@ export default class SettingStorage extends Vue {
return true
} else if (this.modelSettingForm.settingItem === 'Retention Threshold' && !(this.modelSettingForm.retentionThreshold.retention_range_number >= 0 && this.modelSettingForm.retentionThreshold.retention_range_number !== '' && this.modelSettingForm.retentionThreshold.retention_range_type)) {
return true
+ } else if (this.modelSettingForm.settingItem === 'Auto Segment Build' && !this.isValidAutoSegmentBuild()) {
+ return true
} else if (this.modelSettingForm.settingItem.indexOf('spark.') !== -1 && !this.modelSettingForm[this.modelSettingForm.settingItem]) {
return true
} else if (this.modelSettingForm.settingItem === 'is-base-cuboid-always-valid' && this.modelSettingForm[this.modelSettingForm.settingItem] === '') {
@@ -439,6 +488,15 @@ export default class SettingStorage extends Vue {
this.isEdit = true
this.editModelSetting = true
}
+ editAutoSegmentBuildItem (row) {
+ this.modelSettingForm.name = row.alias
+ this.modelSettingForm.settingItem = 'Auto Segment Build'
+ this.modelSettingForm.autoSegmentBuild = JSON.parse(JSON.stringify(row.auto_segment_build))
+ this.activeRow = row
+ this.step = 'stepTwo'
+ this.isEdit = true
+ this.editModelSetting = true
+ }
editSparkItem (row, sparkItemKey) {
this.modelSettingForm.name = row.alias
this.modelSettingForm.settingItem = sparkItemKey.substring(24)
@@ -489,6 +547,13 @@ export default class SettingStorage extends Vue {
if (this.modelSettingForm.settingItem === 'Retention Threshold') {
this.activeRow.retention_range = this.modelSettingForm.retentionThreshold
}
+ if (this.modelSettingForm.settingItem === 'Auto Segment Build') {
+ this.activeRow.auto_segment_build = {
+ ...this.modelSettingForm.autoSegmentBuild,
+ enabled: true,
+ logical_date_offset_days: Number(this.modelSettingForm.autoSegmentBuild.logical_date_offset_days)
+ }
+ }
if (this.modelSettingForm.settingItem.indexOf('spark.') !== -1) {
this.activeRow.override_props['kylin.engine.spark-conf.' + this.modelSettingForm.settingItem] = this.modelSettingForm[this.modelSettingForm.settingItem]
}
@@ -544,6 +609,22 @@ export default class SettingStorage extends Vue {
removeCustomSetting (index) {
this.modelSettingForm[this.modelSettingForm.settingItem].splice(index, 1)
}
+ isValidAutoSegmentBuild () {
+ const config = this.modelSettingForm.autoSegmentBuild
+ if (!config || !config.trigger_time || !config.data_range_start_time || !config.data_range_end_time) return false
+ if (!(+config.logical_date_offset_days >= 1)) return false
+ if (config.data_range_end_time === END_OF_DAY) return true
+ return config.data_range_start_time < config.data_range_end_time
+ }
+ formatAutoSegmentBuild (config) {
+ if (!config) return ''
+ return this.$t('autoSegmentBuildSummary', {
+ trigger: config.trigger_time,
+ offset: config.logical_date_offset_days,
+ start: config.data_range_start_time,
+ end: config.data_range_end_time
+ })
+ }
created () {
this.getConfigList()
}
diff --git a/kystudio/src/components/setting/SettingModel/locales.js b/kystudio/src/components/setting/SettingModel/locales.js
index a0c7c22c1e8..d8ac1a2f98d 100644
--- a/kystudio/src/components/setting/SettingModel/locales.js
+++ b/kystudio/src/components/setting/SettingModel/locales.js
@@ -6,6 +6,7 @@ export default {
segmentMerge: 'Segment Merge:',
volatileRange: 'Volatile Range:',
retention: 'Retention Threshold:',
+ autoSegmentBuild: 'Auto Segment Build:',
newSetting: 'Add Parameter Configuration',
editSetting: 'Edit Parameter Configuration',
modelName: 'Model Name',
@@ -34,6 +35,7 @@ export default {
'isDel_kylin.engine.spark-conf.spark.executor.memory': 'Are you sure delete spark.executor.memory item?',
'isDel_kylin.engine.spark-conf.spark.sql.shuffle.partitions': 'Are you sure delete spark.sql.shuffle.partitions item?',
'isDel_kylin.cube.aggrgroup.is-base-cuboid-always-valid': 'Are you sure delete is-base-cuboid-always-valid item?',
+ isDel_auto_segment_build: 'Are you sure delete auto segment build configuration?',
autoMergeTip: 'The system could auto-merge segment fragments over different merging threshold. Auto-merge will optimize storage to enhance query performance.',
volatileTip: '"Auto-Merge" will not merge the latest segments defined in "Volatile Range". The default value is 0.',
retentionThresholdDesc: 'The segments within the retention threshold would be kept. The rest would be removed automatically.',
@@ -41,6 +43,7 @@ export default {
'Auto-merge': 'Auto Merge',
'Volatile Range': 'Volatile Range',
'Retention Threshold': 'Retention Threshold',
+ 'Auto Segment Build': 'Auto Segment Build',
'kylin.engine.spark-conf.spark.executor.cores': 'kylin.engine.spark-conf.spark.executor.cores',
'kylin.engine.spark-conf.spark.executor.instances': 'kylin.engine.spark-conf.spark.executor.instances',
'kylin.engine.spark-conf.spark.executor.memory': 'kylin.engine.spark-conf.spark.executor.memory',
@@ -55,6 +58,12 @@ export default {
customOptions: 'Besides the defined configurations, you can also add some advanced settings.
Note: It\'s highly recommended to use this feature with the support of Kylin 5 Team.',
customSettingKeyPlaceholder: 'Configuration Name',
customSettingValuePlaceholder: 'Value',
- delCustomConfigTip: 'Are you sure you want to delete custom setting item {name}?'
+ delCustomConfigTip: 'Are you sure you want to delete custom setting item {name}?',
+ autoSegmentBuildTriggerTime: 'Trigger Time',
+ autoSegmentBuildLogicalOffset: 'Logical Date Offset (Day)',
+ autoSegmentBuildRangeStart: 'Range Start Time',
+ autoSegmentBuildRangeEnd: 'Range End Time',
+ 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 dc1b57b9ed2..13007960639 100644
--- a/src/core-common/src/main/resources/kylin-defaults0.properties
+++ b/src/core-common/src/main/resources/kylin-defaults0.properties
@@ -449,6 +449,7 @@ 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 * * * * ?
kylin.model.recommendation-page-size=500
kylin.model.dimension-measure-name.max-length=300
diff --git a/src/core-metadata/src/main/java/org/apache/kylin/metadata/model/AutoSegmentBuildConfig.java b/src/core-metadata/src/main/java/org/apache/kylin/metadata/model/AutoSegmentBuildConfig.java
new file mode 100644
index 00000000000..0b6fb728cd6
--- /dev/null
+++ b/src/core-metadata/src/main/java/org/apache/kylin/metadata/model/AutoSegmentBuildConfig.java
@@ -0,0 +1,49 @@
+/*
+ * 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.kylin.metadata.model;
+
+import java.io.Serializable;
+
+import com.fasterxml.jackson.annotation.JsonProperty;
+
+import lombok.Getter;
+import lombok.Setter;
+
+@Getter
+@Setter
+public class AutoSegmentBuildConfig implements Serializable {
+
+ // LocalTime cannot represent 24:00:00; this value denotes midnight at the end of the logical date.
+ public static final String END_OF_DAY = "24:00:00";
+
+ @JsonProperty("enabled")
+ private boolean enabled = false;
+
+ @JsonProperty("trigger_time")
+ private String triggerTime;
+
+ @JsonProperty("logical_date_offset_days")
+ private Integer logicalDateOffsetDays;
+
+ @JsonProperty("data_range_start_time")
+ private String dataRangeStartTime;
+
+ @JsonProperty("data_range_end_time")
+ private String dataRangeEndTime;
+}
diff --git a/src/core-metadata/src/main/java/org/apache/kylin/metadata/model/SegmentConfig.java b/src/core-metadata/src/main/java/org/apache/kylin/metadata/model/SegmentConfig.java
index 9af96a8caef..b8f41c52e4f 100644
--- a/src/core-metadata/src/main/java/org/apache/kylin/metadata/model/SegmentConfig.java
+++ b/src/core-metadata/src/main/java/org/apache/kylin/metadata/model/SegmentConfig.java
@@ -49,6 +49,9 @@ public class SegmentConfig implements Serializable {
@JsonProperty("create_empty_segment_enabled")
private Boolean createEmptySegmentEnabled = false;
+ @JsonProperty("auto_segment_build")
+ private AutoSegmentBuildConfig autoSegmentBuild = new AutoSegmentBuildConfig();
+
public boolean canSkipAutoMerge() {
return !autoMergeEnabled;
}
diff --git a/src/core-metadata/src/main/java/org/apache/kylin/metadata/project/ProjectInstance.java b/src/core-metadata/src/main/java/org/apache/kylin/metadata/project/ProjectInstance.java
index 50c121dc453..57042062e10 100644
--- a/src/core-metadata/src/main/java/org/apache/kylin/metadata/project/ProjectInstance.java
+++ b/src/core-metadata/src/main/java/org/apache/kylin/metadata/project/ProjectInstance.java
@@ -41,6 +41,7 @@
import org.apache.kylin.guava30.shaded.common.collect.Lists;
import org.apache.kylin.metadata.MetadataConstants;
import org.apache.kylin.metadata.model.AutoMergeTimeEnum;
+import org.apache.kylin.metadata.model.AutoSegmentBuildConfig;
import org.apache.kylin.metadata.model.ISourceAware;
import org.apache.kylin.metadata.model.MaintainModelType;
import org.apache.kylin.metadata.model.RetentionRange;
@@ -114,7 +115,7 @@ public class ProjectInstance extends RootPersistentEntity implements ISourceAwar
@Setter
private SegmentConfig segmentConfig = new SegmentConfig(false, Lists.newArrayList(AutoMergeTimeEnum.WEEK,
AutoMergeTimeEnum.MONTH, AutoMergeTimeEnum.QUARTER, AutoMergeTimeEnum.YEAR), new VolatileRange(),
- new RetentionRange(), false);
+ new RetentionRange(), false, new AutoSegmentBuildConfig());
public static ProjectInstance create(String name, String owner, String description,
LinkedHashMap 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 f8d181a84fc..26a5ed736ed 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));
+ new SegmentConfig(false, Lists.newArrayList(AutoMergeTimeEnum.WEEK), null, null, false, 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
new file mode 100644
index 00000000000..f7968b31f7b
--- /dev/null
+++ b/src/data-loading-service/src/main/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentScheduler.java
@@ -0,0 +1,226 @@
+/*
+ * 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.kylin.rest.scheduler;
+
+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.util.List;
+import java.util.Locale;
+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.collect.Lists;
+import org.apache.kylin.job.execution.AbstractExecutable;
+import org.apache.kylin.job.execution.ExecutableManager;
+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.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.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.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+
+import lombok.val;
+import lombok.extern.slf4j.Slf4j;
+
+@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);
+
+ @Autowired
+ @Qualifier("modelBuildService")
+ private ModelBuildService modelBuildService;
+
+ private final AtomicBoolean dispatching = new AtomicBoolean(false);
+ private final AtomicReference lastDispatchTime = new AtomicReference<>();
+
+ @Scheduled(cron = "${kylin.model.auto-segment-build.dispatcher-cron:*/30 * * * * ?}")
+ public void schedulerAutoBuildSegment() {
+ val currentTime = Instant.now();
+ if (!JobContextUtil.getJobContext(KylinConfig.getInstanceFromEnv()).getJobScheduler().isMaster()) {
+ lastDispatchTime.set(currentTime);
+ return;
+ }
+ if (!dispatching.compareAndSet(false, true)) {
+ log.warn("Skip auto build segment dispatch because the previous dispatch is still running");
+ return;
+ }
+
+ val previousTime = getPreviousDispatchTime(currentTime);
+ try {
+ dispatch(previousTime, 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.minus(INITIAL_TRIGGER_LOOKBACK);
+ }
+ return previousTime;
+ }
+
+ @VisibleForTesting
+ void dispatch(Instant previousTime, Instant currentTime) {
+ val systemConfig = KylinConfig.readSystemKylinConfig();
+ val projectManager = NProjectManager.getInstance(systemConfig);
+ for (ProjectInstance project : projectManager.listAllProjects()) {
+ try {
+ dispatchProject(systemConfig, project, previousTime, 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) {
+ val projectName = project.getName();
+ val zoneId = ZoneId.of(project.getConfig().getTimeZone());
+ val dataflowManager = NDataflowManager.getInstance(systemConfig, projectName);
+ for (NDataModel model : dataflowManager.listOnlineDataModels()) {
+ try {
+ dispatchModel(projectName, model, zoneId, previousTime, currentTime);
+ } catch (Exception e) {
+ log.error("Auto build segment dispatch failed, project: {}, model: {}", projectName, model.getUuid(),
+ e);
+ }
+ }
+ }
+
+ private void dispatchModel(String project, NDataModel model, ZoneId zoneId, Instant previousTime,
+ Instant currentTime) {
+ val autoSegmentBuild = getEligibleConfig(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();
+ if (!scheduledInstant.isAfter(previousTime) || scheduledInstant.isAfter(currentTime)) {
+ return;
+ }
+ submitJob(project, model, autoSegmentBuild, scheduledTime);
+ }
+
+ private AutoSegmentBuildConfig getEligibleConfig(NDataModel model) {
+ if (model.isBroken() || model.isStreaming() || model.isMultiPartitionModel()
+ || PartitionDesc.isEmptyPartitionDesc(model.getPartitionDesc()) || model.getSegmentConfig() == null) {
+ return null;
+ }
+ val autoSegmentBuild = model.getSegmentConfig().getAutoSegmentBuild();
+ return autoSegmentBuild != null && autoSegmentBuild.isEnabled() ? autoSegmentBuild : null;
+ }
+
+ @VisibleForTesting
+ ZonedDateTime getLatestScheduledTime(String triggerTime, ZoneId zoneId, Instant currentTime) {
+ val localCurrentTime = currentTime.atZone(zoneId);
+ val parsedTriggerTime = LocalTime.parse(triggerTime, TIME_FORMATTER);
+ ZonedDateTime scheduledTime = ZonedDateTime.of(localCurrentTime.toLocalDate(), parsedTriggerTime, zoneId);
+ if (scheduledTime.toInstant().isAfter(currentTime)) {
+ scheduledTime = scheduledTime.minusDays(1);
+ }
+ return scheduledTime;
+ }
+
+ @VisibleForTesting
+ void submitJob(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);
+ 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 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, "System");
+ 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);
+ }
+ }
+
+ @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);
+ JobTypeEnum[] buildJobTypes = JobTypeEnum.getJobTypeByCategory(JobTypeEnum.Category.BUILD)
+ .toArray(new JobTypeEnum[0]);
+ List jobs = executableManager.listExecByModelAndStatus(modelId,
+ ExecutableState::isProgressing, buildJobTypes);
+ return !jobs.isEmpty();
+ }
+
+ private Pair parseTime(String time, boolean allow24Hour) {
+ if (allow24Hour && StringUtils.equals(time, END_OF_DAY)) {
+ return Pair.newPair(LocalTime.MIDNIGHT, true);
+ }
+ return Pair.newPair(LocalTime.parse(time, TIME_FORMATTER), false);
+ }
+}
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
new file mode 100644
index 00000000000..fbe61850085
--- /dev/null
+++ b/src/data-loading-service/src/test/java/org/apache/kylin/rest/scheduler/AutoBuildSegmentSchedulerTest.java
@@ -0,0 +1,136 @@
+/*
+ * 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.kylin.rest.scheduler;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+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 org.apache.kylin.common.KylinConfig;
+import org.apache.kylin.junit.annotation.MetadataInfo;
+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.project.NProjectManager;
+import org.apache.kylin.metadata.realization.RealizationStatusEnum;
+import org.apache.kylin.rest.service.ModelBuildService;
+import org.apache.kylin.rest.service.params.IncrementBuildSegmentParams;
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+import org.mockito.Mockito;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import lombok.val;
+
+@MetadataInfo
+class AutoBuildSegmentSchedulerTest {
+ private static final String PROJECT = "default";
+ private static final String MODEL_ID = "89af4ee2-2cdb-4b07-b39e-4c29856309aa";
+
+ @Test
+ void testLatestScheduledTimeUsesProjectTimeZone() {
+ val scheduler = new AutoBuildSegmentScheduler();
+ val zoneId = ZoneId.of("Asia/Shanghai");
+ val currentTime = Instant.parse("2026-08-12T16:00:20Z");
+
+ assertEquals(ZonedDateTime.of(2026, 8, 13, 0, 0, 0, 0, zoneId),
+ 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));
+ }
+
+ @Test
+ void testDispatchModelOnlyOnceForTriggerWindow() 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();
+
+ 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());
+
+ val paramsCaptor = ArgumentCaptor.forClass(IncrementBuildSegmentParams.class);
+ Mockito.verify(modelBuildService).incrementBuildSegmentsByScheduler(paramsCaptor.capture(),
+ Mockito.eq("System"));
+ val logicalDate = scheduledTime.toLocalDate().minusDays(1);
+ assertEquals(String.valueOf(LocalDateTime.of(logicalDate, LocalTime.MIDNIGHT).atZone(zoneId).toInstant()
+ .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 {
+ 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();
+
+ scheduler.submitJob(PROJECT, model, config, ZonedDateTime.of(2026, 8, 12, 1, 0, 0, 0, ZoneId.of("UTC")));
+
+ Mockito.verifyNoInteractions(modelBuildService);
+ }
+
+ @Test
+ void testSkipOfflineModel() 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();
+ NDataflowManager.getInstance(KylinConfig.getInstanceFromEnv(), PROJECT).updateDataflowStatus(MODEL_ID,
+ RealizationStatusEnum.OFFLINE);
+
+ 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 -> {
+ AutoSegmentBuildConfig config = new AutoSegmentBuildConfig();
+ config.setEnabled(true);
+ config.setTriggerTime("01:00:00");
+ config.setLogicalDateOffsetDays(1);
+ config.setDataRangeStartTime("00:00:00");
+ config.setDataRangeEndTime("24:00:00");
+ copyForWrite.getSegmentConfig().setAutoSegmentBuild(config);
+ });
+ }
+}
diff --git a/src/metadata-server/src/main/java/org/apache/kylin/rest/controller/open/OpenModelController.java b/src/metadata-server/src/main/java/org/apache/kylin/rest/controller/open/OpenModelController.java
index a397c6a6059..4aadbeca874 100644
--- a/src/metadata-server/src/main/java/org/apache/kylin/rest/controller/open/OpenModelController.java
+++ b/src/metadata-server/src/main/java/org/apache/kylin/rest/controller/open/OpenModelController.java
@@ -619,6 +619,7 @@ private ModelConfigRequest newModelConfigRequestByConfig(ModelConfigResponse mod
modelConfigRequest.setRetentionRange(modelConfig.getRetentionRange());
modelConfigRequest.setVolatileRange(modelConfig.getVolatileRange());
modelConfigRequest.setAutoMergeTimeRanges(modelConfig.getAutoMergeTimeRanges());
+ modelConfigRequest.setAutoSegmentBuild(modelConfig.getAutoSegmentBuild());
modelConfigRequest.setOverrideProps(modelConfig.getOverrideProps());
return modelConfigRequest;
}
diff --git a/src/modeling-service/src/main/java/org/apache/kylin/rest/request/ModelConfigRequest.java b/src/modeling-service/src/main/java/org/apache/kylin/rest/request/ModelConfigRequest.java
index 0b05af4e1f8..e8d9ec34926 100644
--- a/src/modeling-service/src/main/java/org/apache/kylin/rest/request/ModelConfigRequest.java
+++ b/src/modeling-service/src/main/java/org/apache/kylin/rest/request/ModelConfigRequest.java
@@ -24,6 +24,7 @@
import org.apache.kylin.guava30.shaded.common.collect.Maps;
import org.apache.kylin.metadata.insensitive.ProjectInsensitiveRequest;
import org.apache.kylin.metadata.model.AutoMergeTimeEnum;
+import org.apache.kylin.metadata.model.AutoSegmentBuildConfig;
import org.apache.kylin.metadata.model.RetentionRange;
import org.apache.kylin.metadata.model.VolatileRange;
@@ -44,6 +45,8 @@ public class ModelConfigRequest implements ProjectInsensitiveRequest {
private VolatileRange volatileRange;
@JsonProperty("retention_range")
private RetentionRange retentionRange;
+ @JsonProperty("auto_segment_build")
+ private AutoSegmentBuildConfig autoSegmentBuild;
@JsonProperty("override_props")
LinkedHashMap overrideProps = Maps.newLinkedHashMap();
}
diff --git a/src/modeling-service/src/main/java/org/apache/kylin/rest/response/ModelConfigResponse.java b/src/modeling-service/src/main/java/org/apache/kylin/rest/response/ModelConfigResponse.java
index b2e67171e79..4b6137fe781 100644
--- a/src/modeling-service/src/main/java/org/apache/kylin/rest/response/ModelConfigResponse.java
+++ b/src/modeling-service/src/main/java/org/apache/kylin/rest/response/ModelConfigResponse.java
@@ -23,6 +23,7 @@
import org.apache.kylin.guava30.shaded.common.collect.Maps;
import org.apache.kylin.metadata.model.AutoMergeTimeEnum;
+import org.apache.kylin.metadata.model.AutoSegmentBuildConfig;
import org.apache.kylin.metadata.model.RetentionRange;
import org.apache.kylin.metadata.model.VolatileRange;
@@ -48,6 +49,8 @@ public class ModelConfigResponse {
private VolatileRange volatileRange;
@JsonProperty("retention_range")
private RetentionRange retentionRange;
+ @JsonProperty("auto_segment_build")
+ private AutoSegmentBuildConfig autoSegmentBuild;
@JsonProperty("config_last_modifier")
private String configLastModifier;
@JsonProperty("config_last_modified")
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 89e897a8058..923fe9d8518 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
@@ -281,9 +281,21 @@ public JobInfoResponse incrementBuildSegmentsManually(String project, String mod
@Override
public JobInfoResponse incrementBuildSegmentsManually(IncrementBuildSegmentParams params) throws Exception {
+ return incrementBuildSegmentsInternal(params, getUsername(), true);
+ }
+
+ public JobInfoResponse incrementBuildSegmentsByScheduler(IncrementBuildSegmentParams params, String submitter)
+ throws Exception {
+ return incrementBuildSegmentsInternal(params, submitter, false);
+ }
+
+ private JobInfoResponse incrementBuildSegmentsInternal(IncrementBuildSegmentParams params, String submitter,
+ boolean checkPermission) throws Exception {
String project = params.getProject();
- aclEvaluate.checkProjectOperationPermission(project);
- checkModelPermission(project, params.getModelId());
+ if (checkPermission) {
+ aclEvaluate.checkProjectOperationPermission(project);
+ checkModelPermission(project, params.getModelId());
+ }
val modelManager = getManager(NDataModelManager.class, project);
if (PartitionDesc.isEmptyPartitionDesc(params.getPartitionDesc())) {
throw new KylinException(EMPTY_PARTITION_COLUMN, "Partition column is null.'");
@@ -320,7 +332,7 @@ public JobInfoResponse incrementBuildSegmentsManually(IncrementBuildSegmentParam
.withTag(params.getTag());
List jobIds = EnhancedUnitOfWork.doInTransactionWithCheckAndRetry(() -> {
- List paramList = createSegmentsAndJobParams(buildSegmentParams);
+ List paramList = createSegmentsAndJobParams(buildSegmentParams, submitter);
return createJob(paramList);
}, project);
@@ -342,6 +354,11 @@ private List createJob(List jobParamList) {
}
public List createSegmentsAndJobParams(IncrementBuildSegmentParams params) throws IOException {
+ return createSegmentsAndJobParams(params, getUsername());
+ }
+
+ public List createSegmentsAndJobParams(IncrementBuildSegmentParams params, String submitter)
+ throws IOException {
modelService.checkModelAndIndexManually(params);
if (CollectionUtils.isEmpty(params.getSegmentHoles())) {
params.setSegmentHoles(Lists.newArrayList());
@@ -374,7 +391,7 @@ public List createSegmentsAndJobParams(IncrementBuildSegmentParams par
.withBatchIndexIds(params.getBatchIndexIds()).withYarnQueue(params.getYarnQueue())
.withTag(params.getTag());
NDataSegment segment = createSegment(relParams);
- res.add(createJobParam(relParams, segment));
+ res.add(createJobParam(relParams, segment, submitter));
}
IncrementBuildSegmentParams relParams = new IncrementBuildSegmentParams(params.getProject(),
params.getModelId(), params.getStart(), params.getEnd(), params.getPartitionColFormat(),
@@ -386,18 +403,22 @@ public List createSegmentsAndJobParams(IncrementBuildSegmentParams par
.withBatchIndexIds(params.getBatchIndexIds()).withYarnQueue(params.getYarnQueue())
.withTag(params.getTag());
NDataSegment segment = createSegment(relParams);
- res.add(createJobParam(relParams, segment));
+ res.add(createJobParam(relParams, segment, submitter));
return res;
}
public JobParam createJobParam(IncrementBuildSegmentParams params, NDataSegment segment) {
+ return createJobParam(params, segment, getUsername());
+ }
+
+ public JobParam createJobParam(IncrementBuildSegmentParams params, NDataSegment segment, String submitter) {
if (!params.isNeedBuild()) {
return null;
}
String project = params.getProject();
String modelId = params.getModelId();
NDataModel dataModel = getManager(NDataModelManager.class, project).getDataModelDesc(modelId);
- JobParam jobParam = new JobParam(segment, modelId, getUsername())
+ JobParam jobParam = new JobParam(segment, modelId, submitter)
.withIgnoredSnapshotTables(params.getIgnoredSnapshotTables()).withPriority(params.getPriority())
.withYarnQueue(params.getYarnQueue()).withTag(params.getTag()).withProject(project);
addJobParamExtParams(jobParam, params);
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 33f9c643164..3388f4fa7a9 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
@@ -63,6 +63,7 @@
import static org.apache.kylin.common.exception.code.ErrorCodeServer.SEGMENT_NOT_EXIST_ID;
import static org.apache.kylin.common.exception.code.ErrorCodeServer.SEGMENT_NOT_EXIST_NAME;
import static org.apache.kylin.common.exception.code.ErrorCodeServer.SEGMENT_STATUS;
+import static org.apache.kylin.metadata.model.AutoSegmentBuildConfig.END_OF_DAY;
import static org.apache.kylin.metadata.model.FunctionDesc.PARAMETER_TYPE_COLUMN;
import java.io.IOException;
@@ -70,6 +71,9 @@
import java.sql.SQLException;
import java.text.MessageFormat;
import java.text.SimpleDateFormat;
+import java.time.LocalTime;
+import java.time.format.DateTimeFormatter;
+import java.time.format.DateTimeParseException;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
@@ -158,6 +162,7 @@
import org.apache.kylin.metadata.cube.model.RuleBasedIndex;
import org.apache.kylin.metadata.model.AntiFlatChecker;
import org.apache.kylin.metadata.model.AutoMergeTimeEnum;
+import org.apache.kylin.metadata.model.AutoSegmentBuildConfig;
import org.apache.kylin.metadata.model.ColumnDesc;
import org.apache.kylin.metadata.model.ComputedColumnDesc;
import org.apache.kylin.metadata.model.DataCheckDesc;
@@ -285,7 +290,8 @@ 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);
//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", //
@@ -851,13 +857,14 @@ private ModelStatusToDisplayEnum convertFusionModelStatusToDisplay(NDataModel mo
inconsistentBatchSegmentCount);
if (!batchModel.isBroken() && !modelDesc.isBroken()) {
switch (modelResponseStatus) {
- case ONLINE:
- return (batchModelResponseStatus == ModelStatusToDisplayEnum.WARNING ? ModelStatusToDisplayEnum.WARNING
- : modelResponseStatus);
- case OFFLINE:
- return batchModelResponseStatus;
- default:
- return modelResponseStatus;
+ case ONLINE:
+ return (batchModelResponseStatus == ModelStatusToDisplayEnum.WARNING
+ ? ModelStatusToDisplayEnum.WARNING
+ : modelResponseStatus);
+ case OFFLINE:
+ return batchModelResponseStatus;
+ default:
+ return modelResponseStatus;
}
} else {
return modelResponseStatus;
@@ -3261,6 +3268,7 @@ public List getModelConfig(String project, String modelName
response.setAutoMergeTimeRanges(segmentConfig.getAutoMergeTimeRanges());
response.setVolatileRange(segmentConfig.getVolatileRange());
response.setRetentionRange(segmentConfig.getRetentionRange());
+ response.setAutoSegmentBuild(segmentConfig.getAutoSegmentBuild());
response.setConfigLastModified(dataModel.getConfigLastModified());
response.setConfigLastModifier(dataModel.getConfigLastModifier());
val indexPlan = getIndexPlan(dataModel.getUuid(), project);
@@ -3277,7 +3285,7 @@ public List getModelConfig(String project, String modelName
@Transaction(project = 0)
public void updateModelConfig(String project, String modelId, ModelConfigRequest request) {
aclEvaluate.checkProjectWritePermission(project);
- checkModelConfigParameters(request);
+ checkModelConfigParameters(project, modelId, request);
val dataModelManager = getManager(NDataModelManager.class, project);
dataModelManager.updateDataModel(modelId, copyForWrite -> {
val segmentConfig = copyForWrite.getSegmentConfig();
@@ -3285,6 +3293,7 @@ public void updateModelConfig(String project, String modelId, ModelConfigRequest
segmentConfig.setAutoMergeTimeRanges(request.getAutoMergeTimeRanges());
segmentConfig.setVolatileRange(request.getVolatileRange());
segmentConfig.setRetentionRange(request.getRetentionRange());
+ segmentConfig.setAutoSegmentBuild(request.getAutoSegmentBuild());
copyForWrite.setConfigLastModified(System.currentTimeMillis());
copyForWrite.setConfigLastModifier(BasicService.getUsername());
@@ -3352,10 +3361,16 @@ private void checkPropParameter(ModelConfigRequest request) {
@VisibleForTesting
public void checkModelConfigParameters(ModelConfigRequest request) {
+ checkModelConfigParameters(request.getProject(), null, request);
+ }
+
+ @VisibleForTesting
+ public void checkModelConfigParameters(String project, String modelId, ModelConfigRequest request) {
Boolean autoMergeEnabled = request.getAutoMergeEnabled();
List timeRanges = request.getAutoMergeTimeRanges();
VolatileRange volatileRange = request.getVolatileRange();
RetentionRange retentionRange = request.getRetentionRange();
+ AutoSegmentBuildConfig autoSegmentBuild = request.getAutoSegmentBuild();
if (Boolean.TRUE.equals(autoMergeEnabled) && (null == timeRanges || timeRanges.isEmpty())) {
throw new KylinException(INVALID_PARAMETER, MsgPicker.getMsg().getInvalidAutoMergeConfig());
@@ -3368,9 +3383,62 @@ public void checkModelConfigParameters(ModelConfigRequest request) {
&& retentionRange.getRetentionRangeNumber() < 0) {
throw new KylinException(INVALID_PARAMETER, MsgPicker.getMsg().getInvalidRetentionRangeConfig());
}
+ if (null != autoSegmentBuild && autoSegmentBuild.isEnabled()) {
+ checkAutoSegmentBuildConfig(project, modelId, autoSegmentBuild);
+ }
checkPropParameter(request);
}
+ private void checkAutoSegmentBuildConfig(String project, String modelId, AutoSegmentBuildConfig autoSegmentBuild) {
+ if (StringUtils.isBlank(autoSegmentBuild.getTriggerTime())
+ || StringUtils.isBlank(autoSegmentBuild.getDataRangeStartTime())
+ || StringUtils.isBlank(autoSegmentBuild.getDataRangeEndTime())
+ || autoSegmentBuild.getLogicalDateOffsetDays() == null) {
+ throw new KylinException(INVALID_PARAMETER, "Invalid auto_segment_build config.");
+ }
+ if (autoSegmentBuild.getLogicalDateOffsetDays() < 1) {
+ throw new KylinException(INVALID_PARAMETER, "logical_date_offset_days must be >= 1.");
+ }
+ parseSegmentBuildTime(autoSegmentBuild.getTriggerTime(), "trigger_time", false);
+ LocalTime startTime = parseSegmentBuildTime(autoSegmentBuild.getDataRangeStartTime(), "data_range_start_time",
+ false);
+ boolean endOfNextDay = END_OF_DAY.equals(autoSegmentBuild.getDataRangeEndTime());
+ LocalTime endTime = parseSegmentBuildTime(autoSegmentBuild.getDataRangeEndTime(), "data_range_end_time", true);
+ if (!endOfNextDay && !startTime.isBefore(endTime)) {
+ throw new KylinException(INVALID_PARAMETER,
+ "data_range_start_time must be earlier than data_range_end_time.");
+ }
+ if (StringUtils.isBlank(project) || StringUtils.isBlank(modelId)) {
+ return;
+ }
+ val model = getManager(NDataModelManager.class, project).getDataModelDesc(modelId);
+ if (model == null) {
+ throw new KylinException(MODEL_ID_NOT_EXIST, modelId);
+ }
+ if (model.isStreaming()) {
+ throw new KylinException(INVALID_PARAMETER, "auto_segment_build is not supported for streaming model.");
+ }
+ if (model.isMultiPartitionModel()) {
+ throw new KylinException(INVALID_PARAMETER,
+ "auto_segment_build is not supported for multi partition model.");
+ }
+ if (PartitionDesc.isEmptyPartitionDesc(model.getPartitionDesc())) {
+ throw new KylinException(INVALID_PARAMETER, "auto_segment_build requires a valid partition column.");
+ }
+ }
+
+ private LocalTime parseSegmentBuildTime(String value, String fieldName, boolean allow24Hour) {
+ if (allow24Hour && StringUtils.equals(value, END_OF_DAY)) {
+ return LocalTime.MIDNIGHT;
+ }
+ try {
+ return LocalTime.parse(value, SEGMENT_AUTO_BUILD_TIME_FORMATTER);
+ } catch (DateTimeParseException e) {
+ throw new KylinException(INVALID_PARAMETER,
+ String.format(Locale.ROOT, "Invalid %s, expected HH:mm:ss.", fieldName));
+ }
+ }
+
public ExistedDataRangeResponse getLatestDataRange(String project, String modelId, PartitionDesc desc) {
aclEvaluate.checkProjectReadPermission(project);
Preconditions.checkNotNull(modelId);
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 55e19c629eb..6185d604a8c 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
@@ -125,6 +125,7 @@
import org.apache.kylin.metadata.cube.model.RuleBasedIndex;
import org.apache.kylin.metadata.cube.optimization.FrequencyMap;
import org.apache.kylin.metadata.model.AutoMergeTimeEnum;
+import org.apache.kylin.metadata.model.AutoSegmentBuildConfig;
import org.apache.kylin.metadata.model.BadModelException;
import org.apache.kylin.metadata.model.BadModelException.CauseType;
import org.apache.kylin.metadata.model.ColumnDesc;
@@ -2788,6 +2789,13 @@ public void testUpdateAndGetModelConfig() {
modelConfigRequest.setProject(project);
modelConfigRequest.setAutoMergeEnabled(false);
modelConfigRequest.setAutoMergeTimeRanges(Lists.newArrayList(AutoMergeTimeEnum.WEEK));
+ AutoSegmentBuildConfig autoSegmentBuildConfig = new AutoSegmentBuildConfig();
+ autoSegmentBuildConfig.setEnabled(true);
+ autoSegmentBuildConfig.setTriggerTime("01:00:00");
+ autoSegmentBuildConfig.setLogicalDateOffsetDays(1);
+ autoSegmentBuildConfig.setDataRangeStartTime("00:00:00");
+ autoSegmentBuildConfig.setDataRangeEndTime("24:00:00");
+ modelConfigRequest.setAutoSegmentBuild(autoSegmentBuildConfig);
modelService.updateModelConfig(project, model, modelConfigRequest);
var modelConfigResponses = modelService.getModelConfig(project, null);
@@ -2795,6 +2803,8 @@ public void testUpdateAndGetModelConfig() {
if (modelConfigResponse.getModel().equals(model)) {
Assert.assertEquals(false, modelConfigResponse.getAutoMergeEnabled());
Assert.assertEquals(1, modelConfigResponse.getAutoMergeTimeRanges().size());
+ Assert.assertNotNull(modelConfigResponse.getAutoSegmentBuild());
+ Assert.assertTrue(modelConfigResponse.getAutoSegmentBuild().isEnabled());
}
});
@@ -4330,6 +4340,81 @@ public void testCheckModelConfigParameters() {
checkPropParameter(request);
}
+ @Test
+ public void testCheckModelConfigParameters_AutoSegmentBuildInvalidConfig() {
+ ModelConfigRequest request = new ModelConfigRequest();
+ AutoSegmentBuildConfig config = new AutoSegmentBuildConfig();
+ config.setEnabled(true);
+ config.setTriggerTime("01:00:00");
+ config.setLogicalDateOffsetDays(0);
+ config.setDataRangeStartTime("00:00:00");
+ config.setDataRangeEndTime("24:00:00");
+ request.setAutoSegmentBuild(config);
+ try {
+ modelService.checkModelConfigParameters(request);
+ Assert.fail();
+ } catch (Exception e) {
+ Assert.assertTrue(e instanceof KylinException);
+ Assert.assertTrue(e.getMessage().contains("logical_date_offset_days"));
+ }
+
+ config.setLogicalDateOffsetDays(1);
+ config.setTriggerTime("25: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("10:00:00");
+ config.setDataRangeEndTime("09: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"));
+ }
+ }
+
+ @Test
+ public void testCheckModelConfigParameters_AutoSegmentBuildModelConstraints() {
+ ModelConfigRequest request = new ModelConfigRequest();
+ AutoSegmentBuildConfig config = new AutoSegmentBuildConfig();
+ config.setEnabled(true);
+ config.setTriggerTime("01:00:00");
+ config.setLogicalDateOffsetDays(1);
+ config.setDataRangeStartTime("00:00:00");
+ config.setDataRangeEndTime("24:00:00");
+ request.setAutoSegmentBuild(config);
+
+ NDataModel streamingModel = NDataModelManager.getInstance(getTestConfig(), "streaming_test").listAllModels()
+ .stream().filter(NDataModel::isStreaming).findFirst().orElse(null);
+ Assert.assertNotNull(streamingModel);
+ try {
+ modelService.checkModelConfigParameters("streaming_test", streamingModel.getId(), request);
+ Assert.fail();
+ } catch (Exception e) {
+ Assert.assertTrue(e instanceof KylinException);
+ Assert.assertTrue(e.getMessage().contains("streaming"));
+ }
+
+ NDataModel noPartitionModel = NDataModelManager.getInstance(getTestConfig(), "default").listAllModels()
+ .stream().filter(model -> PartitionDesc.isEmptyPartitionDesc(model.getPartitionDesc())).findFirst()
+ .orElse(null);
+ Assert.assertNotNull(noPartitionModel);
+ try {
+ modelService.checkModelConfigParameters("default", noPartitionModel.getId(), request);
+ Assert.fail();
+ } catch (Exception e) {
+ Assert.assertTrue(e instanceof KylinException);
+ Assert.assertTrue(e.getMessage().contains("partition"));
+ }
+ }
+
@Test
public void testBatchUpdateMultiPartition() {
val modelId = "b780e4e4-69af-449e-b09f-05c90dfa04b6";