This is an automated email from the ASF dual-hosted git repository. SvenO3 pushed a commit to branch 4430-chart-settings-not-correctly-migrated-after-event-schema-changes in repository https://gitbox.apache.org/repos/asf/streampipes.git
commit 3bc596a9000c08e93fe995f36a258815af23732f Author: Sven Oehler <[email protected]> AuthorDate: Tue May 12 13:53:46 2026 +0200 Add chart infos and health status for schema updates --- ...eUpdateInfo.java => ChartSchemaUpdateInfo.java} | 58 +++--- .../model/connect/adapter/PipelineUpdateInfo.java | 11 + .../datalake/DataExplorerWidgetHealthStatus.java | 24 +++ .../model/datalake/DataExplorerWidgetModel.java | 11 + .../update/ChartSchemaUpdateCoordinator.java | 222 +++++++++++++++++++++ .../pipeline/update/PipelineUpdateCoordinator.java | 19 +- .../src/lib/model/gen/streampipes-model.ts | 34 ++++ 7 files changed, 345 insertions(+), 34 deletions(-) diff --git a/streampipes-model/src/main/java/org/apache/streampipes/model/connect/adapter/PipelineUpdateInfo.java b/streampipes-model/src/main/java/org/apache/streampipes/model/connect/adapter/ChartSchemaUpdateInfo.java similarity index 52% copy from streampipes-model/src/main/java/org/apache/streampipes/model/connect/adapter/PipelineUpdateInfo.java copy to streampipes-model/src/main/java/org/apache/streampipes/model/connect/adapter/ChartSchemaUpdateInfo.java index 0d9f79f47b..6c854e69d0 100644 --- a/streampipes-model/src/main/java/org/apache/streampipes/model/connect/adapter/PipelineUpdateInfo.java +++ b/streampipes-model/src/main/java/org/apache/streampipes/model/connect/adapter/ChartSchemaUpdateInfo.java @@ -18,63 +18,61 @@ package org.apache.streampipes.model.connect.adapter; -import org.apache.streampipes.model.pipeline.PipelineElementValidationInfo; import org.apache.streampipes.model.shared.annotation.TsModel; -import java.util.HashMap; +import java.util.ArrayList; import java.util.List; -import java.util.Map; @TsModel -public class PipelineUpdateInfo { +public class ChartSchemaUpdateInfo { - private String pipelineId; - private String pipelineName; + private String chartId; + private String chartTitle; + private String measureName; private boolean canAutoMigrate; - private String migrationInfo; - private Map<String, List<PipelineElementValidationInfo>> validationInfos; + private List<String> affectedFields; - public PipelineUpdateInfo() { - this.validationInfos = new HashMap<>(); + public ChartSchemaUpdateInfo() { + this.affectedFields = new ArrayList<>(); } - public String getPipelineId() { - return pipelineId; + public String getChartId() { + return chartId; } - public void setPipelineId(String pipelineId) { - this.pipelineId = pipelineId; + public void setChartId(String chartId) { + this.chartId = chartId; } - public String getPipelineName() { - return pipelineName; + public String getChartTitle() { + return chartTitle; } - public void setPipelineName(String pipelineName) { - this.pipelineName = pipelineName; + public void setChartTitle(String chartTitle) { + this.chartTitle = chartTitle; } - public boolean isCanAutoMigrate() { - return canAutoMigrate; + public String getMeasureName() { + return measureName; } - public void setCanAutoMigrate(boolean canAutoMigrate) { - this.canAutoMigrate = canAutoMigrate; + public void setMeasureName(String measureName) { + this.measureName = measureName; } - public String getMigrationInfo() { - return migrationInfo; + public boolean isCanAutoMigrate() { + return canAutoMigrate; } - public void setMigrationInfo(String migrationInfo) { - this.migrationInfo = migrationInfo; + public void setCanAutoMigrate(boolean canAutoMigrate) { + this.canAutoMigrate = canAutoMigrate; } - public Map<String, List<PipelineElementValidationInfo>> getValidationInfos() { - return validationInfos; + public List<String> getAffectedFields() { + return affectedFields; } - public void setValidationInfos(Map<String, List<PipelineElementValidationInfo>> validationInfos) { - this.validationInfos = validationInfos; + public void setAffectedFields(List<String> affectedFields) { + this.affectedFields = affectedFields; } } diff --git a/streampipes-model/src/main/java/org/apache/streampipes/model/connect/adapter/PipelineUpdateInfo.java b/streampipes-model/src/main/java/org/apache/streampipes/model/connect/adapter/PipelineUpdateInfo.java index 0d9f79f47b..8ee7184720 100644 --- a/streampipes-model/src/main/java/org/apache/streampipes/model/connect/adapter/PipelineUpdateInfo.java +++ b/streampipes-model/src/main/java/org/apache/streampipes/model/connect/adapter/PipelineUpdateInfo.java @@ -21,6 +21,7 @@ package org.apache.streampipes.model.connect.adapter; import org.apache.streampipes.model.pipeline.PipelineElementValidationInfo; import org.apache.streampipes.model.shared.annotation.TsModel; +import java.util.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -33,9 +34,11 @@ public class PipelineUpdateInfo { private boolean canAutoMigrate; private String migrationInfo; private Map<String, List<PipelineElementValidationInfo>> validationInfos; + private List<ChartSchemaUpdateInfo> chartSchemaUpdateInfos; public PipelineUpdateInfo() { this.validationInfos = new HashMap<>(); + this.chartSchemaUpdateInfos = new ArrayList<>(); } public String getPipelineId() { @@ -77,4 +80,12 @@ public class PipelineUpdateInfo { public void setValidationInfos(Map<String, List<PipelineElementValidationInfo>> validationInfos) { this.validationInfos = validationInfos; } + + public List<ChartSchemaUpdateInfo> getChartSchemaUpdateInfos() { + return chartSchemaUpdateInfos; + } + + public void setChartSchemaUpdateInfos(List<ChartSchemaUpdateInfo> chartSchemaUpdateInfos) { + this.chartSchemaUpdateInfos = chartSchemaUpdateInfos; + } } diff --git a/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/DataExplorerWidgetHealthStatus.java b/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/DataExplorerWidgetHealthStatus.java new file mode 100644 index 0000000000..0f2c2faa42 --- /dev/null +++ b/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/DataExplorerWidgetHealthStatus.java @@ -0,0 +1,24 @@ +/* + * 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.streampipes.model.datalake; + +public enum DataExplorerWidgetHealthStatus { + OK, + REQUIRES_ATTENTION +} diff --git a/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/DataExplorerWidgetModel.java b/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/DataExplorerWidgetModel.java index af5d2739c9..c72a5a6d16 100644 --- a/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/DataExplorerWidgetModel.java +++ b/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/DataExplorerWidgetModel.java @@ -45,12 +45,15 @@ public class DataExplorerWidgetModel extends DashboardEntity { @JsonSerialize(using = CustomMapSerializer.class, as = Map.class) private Map<String, Object> timeSettings; + private DataExplorerWidgetHealthStatus healthStatus; + public DataExplorerWidgetModel() { super(); this.baseAppearanceConfig = new HashMap<>(); this.visualizationConfig = new HashMap<>(); this.dataConfig = new HashMap<>(); this.timeSettings = new HashMap<>(); + this.healthStatus = DataExplorerWidgetHealthStatus.OK; } public String getWidgetId() { @@ -101,4 +104,12 @@ public class DataExplorerWidgetModel extends DashboardEntity { this.timeSettings = timeSettings; } + public DataExplorerWidgetHealthStatus getHealthStatus() { + return healthStatus; + } + + public void setHealthStatus(DataExplorerWidgetHealthStatus healthStatus) { + this.healthStatus = healthStatus; + } + } diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/ChartSchemaUpdateCoordinator.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/ChartSchemaUpdateCoordinator.java new file mode 100644 index 0000000000..471bd2eafb --- /dev/null +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/ChartSchemaUpdateCoordinator.java @@ -0,0 +1,222 @@ +/* + * 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.streampipes.manager.pipeline.update; + +import org.apache.streampipes.model.connect.adapter.ChartSchemaUpdateInfo; +import org.apache.streampipes.model.datalake.DataExplorerWidgetHealthStatus; +import org.apache.streampipes.model.datalake.DataExplorerWidgetModel; +import org.apache.streampipes.model.graph.DataSinkInvocation; +import org.apache.streampipes.model.pipeline.Pipeline; +import org.apache.streampipes.model.schema.EventProperty; +import org.apache.streampipes.model.schema.EventSchema; +import org.apache.streampipes.model.staticproperty.FreeTextStaticProperty; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; +import org.apache.streampipes.storage.management.StorageDispatcher; + +import java.util.LinkedHashSet; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Optional; +import java.util.Set; +import java.util.stream.Collectors; + +public class ChartSchemaUpdateCoordinator { + + private static final String DATA_LAKE_SINK_APP_ID = "org.apache.streampipes.sinks.internal.jvm.datalake"; + private static final String DATA_LAKE_MEASUREMENT_FIELD = "db_measurement"; + private static final String SOURCE_CONFIGS = "sourceConfigs"; + private static final String MEASURE_NAME = "measureName"; + private static final String QUERY_CONFIG = "queryConfig"; + private static final String FIELDS = "fields"; + private static final String RUNTIME_NAME = "runtimeName"; + private static final String SELECTED = "selected"; + private static final String WIDGET_TITLE = "widgetTitle"; + + private final IDataExplorerWidgetStorage widgetStorage; + + public ChartSchemaUpdateCoordinator() { + this(StorageDispatcher.INSTANCE.getNoSqlStore().getDataExplorerWidgetStorage()); + } + + ChartSchemaUpdateCoordinator(IDataExplorerWidgetStorage widgetStorage) { + this.widgetStorage = widgetStorage; + } + + public List<ChartSchemaUpdateInfo> checkChartMigrations(Pipeline pipeline, + EventSchema updatedSchema) { + var measureNames = extractMeasureNames(pipeline); + return widgetStorage + .findAll() + .stream() + .map(widget -> makeUpdateInfo(widget, measureNames, updatedSchema)) + .flatMap(Optional::stream) + .toList(); + } + + public void updateCharts(Pipeline pipeline, + EventSchema updatedSchema) { + var measureNames = extractMeasureNames(pipeline); + widgetStorage + .findAll() + .stream() + .filter(widget -> makeUpdateInfo(widget, measureNames, updatedSchema).isPresent()) + .forEach(widget -> { + widget.setHealthStatus(DataExplorerWidgetHealthStatus.REQUIRES_ATTENTION); + widgetStorage.updateElement(widget); + }); + } + + Optional<ChartSchemaUpdateInfo> makeUpdateInfo(DataExplorerWidgetModel widget, + Set<String> measureNames, + EventSchema updatedSchema) { + var matchingSourceConfigs = getSourceConfigs(widget) + .stream() + .filter(sourceConfig -> sourceConfig.get(MEASURE_NAME) instanceof String measureName + && measureNames.contains(measureName)) + .toList(); + + var affectedFields = matchingSourceConfigs + .stream() + .flatMap(sourceConfig -> findMissingSelectedFields(sourceConfig, updatedSchema).stream()) + .collect(Collectors.toCollection(LinkedHashSet::new)); + + if (affectedFields.isEmpty()) { + return Optional.empty(); + } else { + var info = new ChartSchemaUpdateInfo(); + info.setChartId(widget.getElementId()); + info.setChartTitle(getChartTitle(widget)); + info.setMeasureName(matchingSourceConfigs.stream() + .map(sourceConfig -> sourceConfig.get(MEASURE_NAME)) + .filter(String.class::isInstance) + .map(String.class::cast) + .findFirst() + .orElse(null)); + info.setCanAutoMigrate(false); + info.setAffectedFields(affectedFields + .stream() + .map("Referenced field '%s' no longer exists."::formatted) + .toList()); + return Optional.of(info); + } + } + + private Set<String> findMissingSelectedFields(Map<String, Object> sourceConfig, + EventSchema updatedSchema) { + var updatedFieldNames = extractFieldNames(updatedSchema); + return collectSelectedFields(sourceConfig) + .stream() + .filter(fieldName -> !updatedFieldNames.contains(fieldName)) + .collect(Collectors.toCollection(LinkedHashSet::new)); + } + + private Set<String> extractFieldNames(EventSchema schema) { + if (schema == null || schema.getEventProperties() == null) { + return Set.of(); + } + + return schema + .getEventProperties() + .stream() + .map(EventProperty::getRuntimeName) + .filter(Objects::nonNull) + .collect(Collectors.toSet()); + } + + @SuppressWarnings("unchecked") + private List<Map<String, Object>> getSourceConfigs(DataExplorerWidgetModel widget) { + var sourceConfigs = widget.getDataConfig().get(SOURCE_CONFIGS); + if (sourceConfigs instanceof List<?> configs) { + return configs + .stream() + .filter(Map.class::isInstance) + .map(config -> (Map<String, Object>) config) + .toList(); + } else { + return List.of(); + } + } + + private Set<String> collectSelectedFields(Map<String, Object> sourceConfig) { + var selectedFields = new LinkedHashSet<String>(); + var queryConfig = sourceConfig.get(QUERY_CONFIG); + if (queryConfig instanceof Map<?, ?> queryConfigMap) { + collectFieldConfigs(queryConfigMap.get(FIELDS), selectedFields); + } + return selectedFields; + } + + private void collectFieldConfigs(Object fieldConfigs, + Set<String> selectedFields) { + if (fieldConfigs instanceof List<?> fields) { + fields + .stream() + .filter(Map.class::isInstance) + .map(Map.class::cast) + .filter(field -> Boolean.TRUE.equals(field.get(SELECTED))) + .map(field -> field.get(RUNTIME_NAME)) + .filter(String.class::isInstance) + .map(String.class::cast) + .forEach(selectedFields::add); + } + } + + private Set<String> extractMeasureNames(Pipeline pipeline) { + if (pipeline.getActions() == null) { + return Set.of(); + } + + return pipeline + .getActions() + .stream() + .filter(ChartSchemaUpdateCoordinator::isDataLakeSink) + .map(this::extractMeasureName) + .flatMap(Optional::stream) + .collect(Collectors.toSet()); + } + + private Optional<String> extractMeasureName(DataSinkInvocation sink) { + return Optional + .ofNullable(sink.getStaticProperties()) + .stream() + .flatMap(List::stream) + .filter(property -> DATA_LAKE_MEASUREMENT_FIELD.equals(property.getInternalName())) + .filter(FreeTextStaticProperty.class::isInstance) + .map(FreeTextStaticProperty.class::cast) + .map(FreeTextStaticProperty::getValue) + .filter(Objects::nonNull) + .findFirst(); + } + + private String getChartTitle(DataExplorerWidgetModel widget) { + var baseAppearanceConfig = widget.getBaseAppearanceConfig(); + if (baseAppearanceConfig != null && baseAppearanceConfig.get(WIDGET_TITLE) instanceof String widgetTitle) { + return widgetTitle; + } else if (widget.getWidgetId() != null) { + return widget.getWidgetId(); + } else { + return widget.getElementId(); + } + } + + private static boolean isDataLakeSink(DataSinkInvocation dataSink) { + return DATA_LAKE_SINK_APP_ID.equals(dataSink.getAppId()); + } +} diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/PipelineUpdateCoordinator.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/PipelineUpdateCoordinator.java index 228477dd2d..3a3b290018 100644 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/PipelineUpdateCoordinator.java +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/PipelineUpdateCoordinator.java @@ -50,9 +50,16 @@ public class PipelineUpdateCoordinator { private static final PipelinesStats PIPELINES_STATS = new PipelinesStats(); private final ExtensionServiceRequestManager requestManager; + private final ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator; public PipelineUpdateCoordinator(ExtensionServiceRequestManager requestManager) { + this(requestManager, new ChartSchemaUpdateCoordinator()); + } + + PipelineUpdateCoordinator(ExtensionServiceRequestManager requestManager, + ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator) { this.requestManager = requestManager; + this.chartSchemaUpdateCoordinator = chartSchemaUpdateCoordinator; } public void updatePipelines(SpDataStream dataStream) { @@ -127,6 +134,7 @@ public class PipelineUpdateCoordinator { } StorageDispatcher.INSTANCE.getNoSqlStore().getPipelineStorageAPI().updateElement(modifiedPipeline); + chartSchemaUpdateCoordinator.updateCharts(modifiedPipeline, updatedEventSchema); if (shouldRestartPipeline && canAutoMigrate) { new PipelineExecutor(PipelineManager.getPipeline(pipeline.getPipelineId()), requestManager).startPipeline(); @@ -141,20 +149,23 @@ public class PipelineUpdateCoordinator { String updatedStreamName, EventSchema updatedEventSchema) { var affectedPipelines = PipelineManager.getPipelinesContainingElements(affectedElementId); - var updateInfos = new ArrayList<PipelineUpdateInfo>(); + var pipelineUpdateInfos = new ArrayList<PipelineUpdateInfo>(); affectedPipelines.forEach(pipeline -> { var updatedPipeline = updatePipeline(pipeline, affectedElementId, updatedStreamName, updatedEventSchema); try { - var modificationMessage = new PipelineVerificationHandlerV2(updatedPipeline, requestManager).verifyPipeline(); + var verificationHandler = new PipelineVerificationHandlerV2(updatedPipeline, requestManager); + var modificationMessage = verificationHandler.verifyPipeline(); var updateInfo = makeUpdateInfo(modificationMessage, updatedPipeline); - updateInfos.add(updateInfo); + updateInfo.setChartSchemaUpdateInfos( + chartSchemaUpdateCoordinator.checkChartMigrations(updatedPipeline, updatedEventSchema)); + pipelineUpdateInfos.add(updateInfo); } catch (Exception e) { throw new RuntimeException(e); } }); - return updateInfos; + return pipelineUpdateInfos; } private Pipeline updatePipeline(Pipeline pipeline, diff --git a/ui/projects/streampipes/platform-services/src/lib/model/gen/streampipes-model.ts b/ui/projects/streampipes/platform-services/src/lib/model/gen/streampipes-model.ts index 1c09e3e189..b1bff83646 100644 --- a/ui/projects/streampipes/platform-services/src/lib/model/gen/streampipes-model.ts +++ b/ui/projects/streampipes/platform-services/src/lib/model/gen/streampipes-model.ts @@ -1302,6 +1302,7 @@ export class DashboardModel implements Storable, SpResource { export class DataExplorerWidgetModel extends DashboardEntity { baseAppearanceConfig: { [index: string]: any }; dataConfig: { [index: string]: any }; + healthStatus: DataExplorerWidgetHealthStatus; timeSettings: { [index: string]: any }; visualizationConfig: { [index: string]: any }; widgetId: string; @@ -1322,6 +1323,7 @@ export class DataExplorerWidgetModel extends DashboardEntity { instance.dataConfig = __getCopyObjectFn(__identity<any>())( data.dataConfig, ); + instance.healthStatus = data.healthStatus; instance.timeSettings = __getCopyObjectFn(__identity<any>())( data.timeSettings, ); @@ -3286,6 +3288,7 @@ export class PipelineTemplateGenerationRequest { export class PipelineUpdateInfo { canAutoMigrate: boolean; + chartSchemaUpdateInfos: ChartSchemaUpdateInfo[]; migrationInfo: string; pipelineId: string; pipelineName: string; @@ -3300,6 +3303,9 @@ export class PipelineUpdateInfo { } const instance = target || new PipelineUpdateInfo(); instance.canAutoMigrate = data.canAutoMigrate; + instance.chartSchemaUpdateInfos = __getCopyArrayFn( + ChartSchemaUpdateInfo.fromData, + )(data.chartSchemaUpdateInfos); instance.migrationInfo = data.migrationInfo; instance.pipelineId = data.pipelineId; instance.pipelineName = data.pipelineName; @@ -3310,6 +3316,32 @@ export class PipelineUpdateInfo { } } +export class ChartSchemaUpdateInfo { + affectedFields: string[]; + canAutoMigrate: boolean; + chartId: string; + chartTitle: string; + measureName: string; + + static fromData( + data: ChartSchemaUpdateInfo, + target?: ChartSchemaUpdateInfo, + ): ChartSchemaUpdateInfo { + if (!data) { + return data; + } + const instance = target || new ChartSchemaUpdateInfo(); + instance.affectedFields = __getCopyArrayFn(__identity<string>())( + data.affectedFields, + ); + instance.canAutoMigrate = data.canAutoMigrate; + instance.chartId = data.chartId; + instance.chartTitle = data.chartTitle; + instance.measureName = data.measureName; + return instance; + } +} + export class ProducedMessagesInfo extends MessagesInfo { totalProducedMessages: number; totalProducedMessagesSincePipelineStart: number; @@ -4640,6 +4672,8 @@ export type PipelineHealthStatus = export type MeasurementUpdateAction = 'edit-pipeline' | 'manage-datasets'; +export type DataExplorerWidgetHealthStatus = 'OK' | 'REQUIRES_ATTENTION'; + export type PropertyScope = | 'HEADER_PROPERTY' | 'DIMENSION_PROPERTY'
