This is an automated email from the ASF dual-hosted git repository. dominikriemer pushed a commit to branch fix-data-lake-sink-retention in repository https://gitbox.apache.org/repos/asf/streampipes.git
commit a1ee9022630c088935591571229d45493c810a30 Author: Dominik Riemer <[email protected]> AuthorDate: Tue May 12 21:45:46 2026 +0200 fix: Fetch retention config through client --- .../streampipes/sinks/internal/jvm/datalake/DataLakeSink.java | 9 +++------ 1 file changed, 3 insertions(+), 6 deletions(-) diff --git a/streampipes-extensions/streampipes-sinks-internal-jvm/src/main/java/org/apache/streampipes/sinks/internal/jvm/datalake/DataLakeSink.java b/streampipes-extensions/streampipes-sinks-internal-jvm/src/main/java/org/apache/streampipes/sinks/internal/jvm/datalake/DataLakeSink.java index 40230d9e5a..10c139c48b 100644 --- a/streampipes-extensions/streampipes-sinks-internal-jvm/src/main/java/org/apache/streampipes/sinks/internal/jvm/datalake/DataLakeSink.java +++ b/streampipes-extensions/streampipes-sinks-internal-jvm/src/main/java/org/apache/streampipes/sinks/internal/jvm/datalake/DataLakeSink.java @@ -22,7 +22,6 @@ import org.apache.streampipes.client.api.IStreamPipesClient; import org.apache.streampipes.commons.environment.Environments; import org.apache.streampipes.commons.exceptions.SpRuntimeException; import org.apache.streampipes.dataexplorer.TimeSeriesStore; -import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement; import org.apache.streampipes.dataexplorer.management.DataExplorerDispatcher; import org.apache.streampipes.extensions.api.extractor.IStaticPropertyExtractor; import org.apache.streampipes.extensions.api.pe.IStreamPipesDataSink; @@ -164,12 +163,10 @@ public class DataLakeSink implements IStreamPipesDataSink, SupportsRuntimeConfig return staticProperty; } - private RetentionTimeConfig getRetentionTime(String measureName, IStreamPipesClient client){ + private RetentionTimeConfig getRetentionTime(String measureName, + IStreamPipesClient client){ - IDataExplorerSchemaManagement dataExplorerSchemaManagement = new DataExplorerDispatcher().getDataExplorerManager() - .getSchemaManagement(); - - var originalMeasure = dataExplorerSchemaManagement.getExistingMeasureByName(measureName); + var originalMeasure = client.dataLakeMeasureApi().getByDatasetName(measureName); RetentionTimeConfig retentionTime = null;
