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;
}