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


The following commit(s) were added to 
refs/heads/4430-chart-settings-not-correctly-migrated-after-event-schema-changes
 by this push:
     new 7da1324f5c Update charts when pipeline is started
7da1324f5c is described below

commit 7da1324f5c560db14acc7e01078c44c59b0cc3a0
Author: Sven Oehler <[email protected]>
AuthorDate: Wed May 13 18:03:46 2026 +0200

    Update charts when pipeline is started
---
 .../dataexplorer/DataExplorerSchemaManagement.java | 11 +++
 .../update/ChartSchemaUpdateCoordinator.java       | 98 ++++++++++++++++++----
 .../pipeline/update/PipelineUpdateCoordinator.java |  1 -
 .../components/chart-view/chart-view.component.ts  | 38 ---------
 4 files changed, 95 insertions(+), 53 deletions(-)

diff --git 
a/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/DataExplorerSchemaManagement.java
 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/DataExplorerSchemaManagement.java
index 7a2076a879..a833e41860 100644
--- 
a/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/DataExplorerSchemaManagement.java
+++ 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/DataExplorerSchemaManagement.java
@@ -21,6 +21,7 @@ package org.apache.streampipes.dataexplorer;
 import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement;
 import 
org.apache.streampipes.manager.matching.v2.pipeline.MeasurementChangeDetector;
 import org.apache.streampipes.manager.permission.DataLakePermissionManager;
+import 
org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator;
 import org.apache.streampipes.model.datalake.DataLakeMeasure;
 import 
org.apache.streampipes.model.datalake.DataLakeMeasureSchemaUpdateStrategy;
 import org.apache.streampipes.model.schema.EventProperty;
@@ -30,6 +31,7 @@ import org.apache.streampipes.storage.api.core.CRUDStorage;
 import java.util.ArrayList;
 import java.util.List;
 import java.util.Optional;
+import java.util.Set;
 import java.util.UUID;
 import java.util.function.Function;
 import java.util.stream.Collectors;
@@ -39,11 +41,19 @@ public class DataExplorerSchemaManagement implements 
IDataExplorerSchemaManageme
 
   CRUDStorage<DataLakeMeasure> dataLakeStorage;
   private final DataLakePermissionManager permissionManager;
+  private final ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator;
 
   public DataExplorerSchemaManagement(CRUDStorage<DataLakeMeasure> 
dataLakeStorage,
                                       DataLakePermissionManager 
permissionManager) {
+    this(dataLakeStorage, permissionManager, new 
ChartSchemaUpdateCoordinator());
+  }
+
+  DataExplorerSchemaManagement(CRUDStorage<DataLakeMeasure> dataLakeStorage,
+                               DataLakePermissionManager permissionManager,
+                               ChartSchemaUpdateCoordinator 
chartSchemaUpdateCoordinator) {
     this.dataLakeStorage = dataLakeStorage;
     this.permissionManager = permissionManager;
+    this.chartSchemaUpdateCoordinator = chartSchemaUpdateCoordinator;
   }
 
   @Override
@@ -97,6 +107,7 @@ public class DataExplorerSchemaManagement implements 
IDataExplorerSchemaManageme
       // one
       unifyEventSchemaAndUpdateMeasure(measure, existingMeasure);
     }
+    
chartSchemaUpdateCoordinator.updateCharts(Set.of(measure.getMeasureName()), 
measure.getEventSchema());
   }
 
   /**
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
index cbe2fd013c..38015777f3 100644
--- 
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
@@ -21,14 +21,19 @@ 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.datalake.DataLakeMeasure;
 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.serializers.json.JacksonSerializer;
 import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage;
 import org.apache.streampipes.storage.management.StorageDispatcher;
 
+import com.fasterxml.jackson.core.type.TypeReference;
+import com.fasterxml.jackson.databind.ObjectMapper;
+
 import java.util.LinkedHashSet;
 import java.util.List;
 import java.util.Map;
@@ -43,13 +48,17 @@ public class ChartSchemaUpdateCoordinator {
   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 MEASURE = "measure";
   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 static final TypeReference<Map<String, Object>> MAP_TYPE = new 
TypeReference<>() {
+  };
 
   private final IDataExplorerWidgetStorage widgetStorage;
+  private final ObjectMapper objectMapper;
 
   public ChartSchemaUpdateCoordinator() {
     
this(StorageDispatcher.INSTANCE.getNoSqlStore().getDataExplorerWidgetStorage());
@@ -57,6 +66,7 @@ public class ChartSchemaUpdateCoordinator {
 
   ChartSchemaUpdateCoordinator(IDataExplorerWidgetStorage widgetStorage) {
     this.widgetStorage = widgetStorage;
+    this.objectMapper = JacksonSerializer.getObjectMapper();
   }
 
   public List<ChartSchemaUpdateInfo> checkChartMigrations(Pipeline pipeline,
@@ -73,27 +83,39 @@ public class ChartSchemaUpdateCoordinator {
   public void updateCharts(Pipeline pipeline,
                            EventSchema updatedSchema) {
     var measureNames = extractMeasureNames(pipeline);
+    updateCharts(measureNames, updatedSchema);
+  }
+
+  public void updateCharts(Set<String> measureNames, EventSchema 
updatedSchema) {
     widgetStorage
         .findAll()
-        .stream()
-        .map(widget -> makeUpdateInfo(widget, measureNames, updatedSchema)
-            .map(updateInfo -> Map.entry(widget, updateInfo)))
-        .flatMap(Optional::stream)
-        .forEach(entry -> {
-          
entry.getKey().setHealthStatus(DataExplorerWidgetHealthStatus.REQUIRES_ATTENTION);
-          
entry.getKey().setAffectedSchemaUpdateFields(entry.getValue().getAffectedFields());
-          widgetStorage.updateElement(entry.getKey());
-        });
+        .forEach(widget -> updateChart(widget, measureNames, updatedSchema));
+  }
+
+  private void updateChart(DataExplorerWidgetModel widget,
+                           Set<String> measureNames,
+                           EventSchema updatedSchema) {
+    var matchingSourceConfigs = getMatchingSourceConfigs(widget, measureNames);
+    if (matchingSourceConfigs.isEmpty()) {
+      return;
+    }
+
+    var updateInfo = makeUpdateInfo(widget, measureNames, updatedSchema);
+    var measureSchemaUpdated = 
updateSourceConfigMeasures(matchingSourceConfigs, updatedSchema);
+    updateInfo.ifPresent(info -> {
+      
widget.setHealthStatus(DataExplorerWidgetHealthStatus.REQUIRES_ATTENTION);
+      widget.setAffectedSchemaUpdateFields(info.getAffectedFields());
+    });
+
+    if (measureSchemaUpdated || updateInfo.isPresent()) {
+      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 matchingSourceConfigs = getMatchingSourceConfigs(widget, measureNames);
 
     var affectedFields = matchingSourceConfigs
         .stream()
@@ -127,6 +149,54 @@ public class ChartSchemaUpdateCoordinator {
         .collect(Collectors.toCollection(LinkedHashSet::new));
   }
 
+  private List<Map<String, Object>> 
getMatchingSourceConfigs(DataExplorerWidgetModel widget,
+                                                             Set<String> 
measureNames) {
+    return getSourceConfigs(widget)
+        .stream()
+        .filter(sourceConfig -> sourceConfig.get(MEASURE_NAME) instanceof 
String measureName
+            && measureNames.contains(measureName))
+        .toList();
+  }
+
+  private boolean updateSourceConfigMeasures(List<Map<String, Object>> 
sourceConfigs,
+                                             EventSchema updatedSchema) {
+    return sourceConfigs
+        .stream()
+        .map(sourceConfig -> updateSourceConfigMeasure(sourceConfig, 
updatedSchema))
+        .reduce(false, Boolean::logicalOr);
+  }
+
+  private boolean updateSourceConfigMeasure(Map<String, Object> sourceConfig,
+                                            EventSchema updatedSchema) {
+    if (sourceConfig.get(MEASURE_NAME) instanceof String measureName) {
+      var measure = parseMeasure(sourceConfig.get(MEASURE), measureName);
+      measure.setEventSchema(updatedSchema);
+      sourceConfig.put(MEASURE, serializeMeasure(measure));
+      return true;
+    } else {
+      return false;
+    }
+  }
+
+  private DataLakeMeasure parseMeasure(Object measure,
+                                       String measureName) {
+    var dataLakeMeasure = objectMapper.convertValue(measure, 
DataLakeMeasure.class);
+    if (dataLakeMeasure == null) {
+      dataLakeMeasure = new DataLakeMeasure();
+    }
+    if (dataLakeMeasure.getMeasureName() == null) {
+      dataLakeMeasure.setMeasureName(measureName);
+    }
+    if (dataLakeMeasure.getSchemaVersion() == null) {
+      dataLakeMeasure.setSchemaVersion(DataLakeMeasure.CURRENT_SCHEMA_VERSION);
+    }
+    return dataLakeMeasure;
+  }
+
+  private Map<String, Object> serializeMeasure(DataLakeMeasure measure) {
+    return objectMapper.convertValue(measure, MAP_TYPE);
+  }
+
   private Set<String> extractFieldNames(EventSchema schema) {
     if (schema == null || schema.getEventProperties() == null) {
       return Set.of();
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 3a3b290018..c71d241607 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
@@ -134,7 +134,6 @@ public class PipelineUpdateCoordinator {
         }
 
         
StorageDispatcher.INSTANCE.getNoSqlStore().getPipelineStorageAPI().updateElement(modifiedPipeline);
-        chartSchemaUpdateCoordinator.updateCharts(modifiedPipeline, 
updatedEventSchema);
 
         if (shouldRestartPipeline && canAutoMigrate) {
           new 
PipelineExecutor(PipelineManager.getPipeline(pipeline.getPipelineId()), 
requestManager).startPipeline();
diff --git a/ui/src/app/chart/components/chart-view/chart-view.component.ts 
b/ui/src/app/chart/components/chart-view/chart-view.component.ts
index 9a826876f8..6c18303e2f 100644
--- a/ui/src/app/chart/components/chart-view/chart-view.component.ts
+++ b/ui/src/app/chart/components/chart-view/chart-view.component.ts
@@ -250,9 +250,6 @@ export class ChartViewComponent
                     this.chartNotFound = true;
                     return of(null);
                 }),
-                switchMap(res =>
-                    res ? this.refreshDataViewMeasureSchemas(res) : of(null),
-                ),
             )
             .subscribe(res => {
                 if (!res) {
@@ -627,41 +624,6 @@ export class ChartViewComponent
         this.queryParams$?.unsubscribe();
     }
 
-    private refreshDataViewMeasureSchemas(
-        dataView: DataExplorerWidgetModel,
-    ): Observable<DataExplorerWidgetModel> {
-        const sourceConfigs = this.getSourceConfigs(dataView);
-        if (sourceConfigs.length === 0) {
-            return of(dataView);
-        }
-
-        return this.datalakeRestService.getAllMeasurementSeries().pipe(
-            map(measures => {
-                const measuresByName = new Map(
-                    measures.map(measure => [measure.measureName, measure]),
-                );
-
-                sourceConfigs.forEach(sourceConfig => {
-                    const latestMeasure = measuresByName.get(
-                        sourceConfig.measureName,
-                    );
-                    if (latestMeasure) {
-                        sourceConfig.measure = latestMeasure;
-                    }
-                });
-
-                return dataView;
-            }),
-            catchError(() => of(dataView)),
-        );
-    }
-
-    private getSourceConfigs(
-        dataView: DataExplorerWidgetModel,
-    ): SourceConfig[] {
-        return dataView?.dataConfig?.sourceConfigs ?? [];
-    }
-
     private hasMultipleSourceConfigs(widget: DataExplorerWidgetModel): boolean 
{
         return (widget?.dataConfig?.sourceConfigs?.length ?? 0) > 1;
     }

Reply via email to