This is an automated email from the ASF dual-hosted git repository.

bossenti pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git


The following commit(s) were added to refs/heads/dev by this push:
     new b75cc2b9c8 refactor: abstract data explorer modules from 
Influx-related code (#2803)
b75cc2b9c8 is described below

commit b75cc2b9c812829de0c829f581f9eca86dcc225a
Author: Tim <[email protected]>
AuthorDate: Fri May 3 15:27:07 2024 +0200

    refactor: abstract data explorer modules from Influx-related code (#2803)
    
    * refactor: introduce new env variable to determine time series storage
    
    * refactor: abstract sanitation
    
    * refactor: abstract query management
    
    * refactor: abstract counter
    
    * refactor: abstract query execution from Influx
    
    * refactor: abstract time series store from Influx
    
    * refactor: several adaptions
    
    * refactor: adapt usages to new APIs
    
    * feat: data storage management
    
    * refactor: small improvements
    
    * style: fix indentation & import order
    
    * style: fix import order
    
    * style: fix import order
    
    * adapt spelling of influx tag
    
    * change data explorer storage dispatcher to regular class
    
    * adapt to new syntax
    
    * refactor: simplify implementation
---
 .../apache/streampipes/commons/constants/Envs.java |   1 +
 .../commons/environment/DefaultEnvironment.java    |   5 +
 .../commons/environment/Environment.java           |   2 +
 streampipes-data-explorer-api/pom.xml              |   5 +
 .../dataexplorer/api/IDataExplorerManager.java     |  48 ++++
 .../api}/IDataExplorerQueryManagement.java         |   7 +-
 .../api/IDataLakeMeasurementCounter.java           |  20 +-
 .../api/IDataLakeMeasurementSanitizer.java         |  48 ++++
 .../dataexplorer/api/ITimeSeriesStorage.java       |  18 +-
 .../influx/DataExplorerInfluxQueryExecutor.java    |  19 +-
 .../influx/DataExplorerManagerInflux.java          |  71 ++++++
 ...java => DataExplorerQueryManagementInflux.java} |  22 +-
 ....java => DataLakeMeasurementCounterInflux.java} |  14 +-
 .../influx/{migrate => }/DeleteDataQuery.java      |   2 +-
 .../dataexplorer/influx/InfluxStore.java           | 244 ---------------------
 .../influx/TimeSeriesStorageInflux.java            | 132 +++++++++++
 .../influx/migrate/DataExplorerUtils.java          |  86 --------
 .../DataLakeMeasurementSanitizerInflux.java        |  44 ++++
 ...eTest.java => TimeSeriesStorageInfluxTest.java} |  10 +-
 .../DataLakeMeasurementSanitizerInfluxTest.java    |  93 ++++++++
 streampipes-data-explorer-management/pom.xml       |   5 +
 .../management/DataExplorerDispatcher.java         |  27 +--
 .../management/SupportedDataExplorerStorages.java  |  16 +-
 .../dataexplorer}/AutoAggregationHandler.java      |  23 +-
 .../dataexplorer/DataLakeMeasurementSanitizer.java | 118 ++++++++++
 .../streampipes/dataexplorer}/ImageStore.java      |  10 +-
 .../streampipes/dataexplorer}/ImageStoreUtils.java |   5 +-
 .../dataexplorer}/QueryResultProvider.java         |  28 ++-
 .../dataexplorer}/StreamedQueryResultProvider.java |  12 +-
 .../dataexplorer/TimeSeriesStorage.java            | 145 ++++++++++++
 .../streampipes/dataexplorer}/TimeSeriesStore.java |  28 +--
 .../query/DataExplorerQueryExecutor.java           |  39 ++--
 .../streampipes-sinks-internal-jvm/pom.xml         |  12 +-
 .../sinks/internal/jvm/datalake/DataLakeSink.java  |  51 +++--
 streampipes-platform-services/pom.xml              |   5 +
 .../streampipes/ps/DataLakeMeasureResourceV4.java  |  17 +-
 .../apache/streampipes/ps/DataLakeResourceV4.java  |  55 ++---
 streampipes-rest/pom.xml                           |   5 +
 .../apache/streampipes/rest/ResetManagement.java   |  15 +-
 39 files changed, 958 insertions(+), 549 deletions(-)

diff --git 
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/constants/Envs.java
 
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/constants/Envs.java
index 5f8de15ef6..edcbb32999 100644
--- 
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/constants/Envs.java
+++ 
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/constants/Envs.java
@@ -54,6 +54,7 @@ public enum Envs {
 
 
   // Time Series Storage
+  SP_TS_STORAGE("SP_TS_STORAGE", "influxdb"),
   SP_TS_STORAGE_PROTOCOL("SP_TS_STORAGE_PROTOCOL", "http"),
   SP_TS_STORAGE_HOST("SP_TS_STORAGE_HOST", "influxdb", 
DefaultEnvValues.LOCALHOST),
   SP_TS_STORAGE_PORT("SP_TS_STORAGE_PORT", "8086"),
diff --git 
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/DefaultEnvironment.java
 
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/DefaultEnvironment.java
index 3fef8c55f7..69b4033551 100644
--- 
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/DefaultEnvironment.java
+++ 
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/DefaultEnvironment.java
@@ -55,6 +55,11 @@ public class DefaultEnvironment implements Environment {
     return new IntEnvironmentVariable(Envs.SP_CORE_PORT);
   }
 
+  @Override
+  public StringEnvironmentVariable getTsStorage() {
+    return new StringEnvironmentVariable(Envs.SP_TS_STORAGE);
+  }
+
   @Override
   public StringEnvironmentVariable getTsStorageProtocol() {
     return new StringEnvironmentVariable(Envs.SP_TS_STORAGE_PROTOCOL);
diff --git 
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/Environment.java
 
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/Environment.java
index ceb2d2d1f2..e190402378 100644
--- 
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/Environment.java
+++ 
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/Environment.java
@@ -38,6 +38,8 @@ public interface Environment {
 
   // Time series storage env variables
 
+  StringEnvironmentVariable getTsStorage();
+
   StringEnvironmentVariable getTsStorageProtocol();
 
   StringEnvironmentVariable getTsStorageHost();
diff --git a/streampipes-data-explorer-api/pom.xml 
b/streampipes-data-explorer-api/pom.xml
index a433295db7..c8c573151d 100644
--- a/streampipes-data-explorer-api/pom.xml
+++ b/streampipes-data-explorer-api/pom.xml
@@ -34,6 +34,11 @@
 
     <dependencies>
         <!-- StreamPipes dependencies -->
+        <dependency>
+            <groupId>org.apache.streampipes</groupId>
+            <artifactId>streampipes-client-api</artifactId>
+            <version>0.95.0-SNAPSHOT</version>
+        </dependency>
         <dependency>
             <groupId>org.apache.streampipes</groupId>
             <artifactId>streampipes-data-explorer-export</artifactId>
diff --git 
a/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerManager.java
 
b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerManager.java
new file mode 100644
index 0000000000..ef1b68a0ca
--- /dev/null
+++ 
b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerManager.java
@@ -0,0 +1,48 @@
+/*
+ * 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.dataexplorer.api;
+
+import org.apache.streampipes.client.api.IStreamPipesClient;
+import org.apache.streampipes.model.datalake.DataLakeMeasure;
+
+import java.util.List;
+
+public interface IDataExplorerManager {
+
+  /**
+   * Provide an instance of {@link IDataLakeMeasurementCounter} for counting 
the sizes of measurements within a data
+   * lake.
+   *
+   * @param allMeasurements     A list of {@link DataLakeMeasure} objects 
representing all measurements in the data lake.
+   * @param measurementsToCount A list of measurement names for which the 
sizes should be counted.
+   * @return An instance of {@link IDataLakeMeasurementCounter} configured to 
count the sizes of the specified measurements.
+   */
+  IDataLakeMeasurementCounter getMeasurementCounter(
+    List<DataLakeMeasure> allMeasurements,
+    List<String> measurementsToCount
+  );
+
+  IDataExplorerQueryManagement 
getQueryManagement(IDataExplorerSchemaManagement dataExplorerSchemaManagement);
+
+  IDataExplorerSchemaManagement getSchemaManagement();
+
+  ITimeSeriesStorage getTimeseriesStorage(DataLakeMeasure measure);
+
+  IDataLakeMeasurementSanitizer getMeasurementSanitizer(IStreamPipesClient 
client, DataLakeMeasure measure);
+}
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/IDataExplorerQueryManagement.java
 
b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerQueryManagement.java
similarity index 88%
rename from 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/IDataExplorerQueryManagement.java
rename to 
streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerQueryManagement.java
index 37e43002fc..1e3b500925 100644
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/IDataExplorerQueryManagement.java
+++ 
b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerQueryManagement.java
@@ -16,7 +16,7 @@
  *
  */
 
-package org.apache.streampipes.dataexplorer.influx.migrate;
+package org.apache.streampipes.dataexplorer.api;
 
 import org.apache.streampipes.dataexplorer.export.OutputFormat;
 import org.apache.streampipes.model.datalake.SpQueryResult;
@@ -28,8 +28,9 @@ import java.util.Map;
 
 public interface IDataExplorerQueryManagement {
 
-  SpQueryResult getData(ProvidedRestQueryParams queryParams,
-                        boolean ignoreMissingData) throws 
IllegalArgumentException;
+  SpQueryResult getData(
+    ProvidedRestQueryParams queryParams,
+    boolean ignoreMissingData) throws IllegalArgumentException;
 
   void getDataAsStream(ProvidedRestQueryParams params,
                        OutputFormat format,
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/package-info.java
 
b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataLakeMeasurementCounter.java
similarity index 60%
copy from 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/package-info.java
copy to 
streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataLakeMeasurementCounter.java
index c96a83d257..a5527f996a 100644
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/package-info.java
+++ 
b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataLakeMeasurementCounter.java
@@ -16,10 +16,20 @@
  *
  */
 
+package org.apache.streampipes.dataexplorer.api;
+
+import java.util.Map;
+
 /**
- * This package contains classes that have direct dependencies on 
InfluxDB-specific implementations.
- * To enhance code organization and promote agnostic design, these classes are 
slated for migration to other modules,
- * such as 'streampipes-data-explorer' or 'streampipes-data-explorer-api'.
- * Pending migration, they are temporarily housed within this package.
+ * Interface for counting the number of events per measurement within the 
StreamPipes data storage.
  */
-package org.apache.streampipes.dataexplorer.influx.migrate;
\ No newline at end of file
+public interface IDataLakeMeasurementCounter {
+
+  /**
+   * Counts the sizes of measurements within the StreamPipes data storage.
+   *
+   * @return A map where each key represents a measurement name and its 
corresponding value represents
+   *         the number of events contained by that measurement.
+   */
+  Map<String, Integer> countMeasurementSizes();
+}
\ No newline at end of file
diff --git 
a/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataLakeMeasurementSanitizer.java
 
b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataLakeMeasurementSanitizer.java
new file mode 100644
index 0000000000..706f66ce38
--- /dev/null
+++ 
b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataLakeMeasurementSanitizer.java
@@ -0,0 +1,48 @@
+/*
+ * 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.dataexplorer.api;
+
+import org.apache.streampipes.model.datalake.DataLakeMeasure;
+
+/**
+ * The IDataLakeMeasurementSanitizer interface defines methods for sanitizing 
and registering or
+ * updating data lake measures.
+ * Implementations of this interface provide functionality to ensure that the 
measurement complies to
+ * the requirements of the underlying time series storage, e.g., to not 
contain any reserved symbols.
+ */
+public interface IDataLakeMeasurementSanitizer {
+
+  /**
+   * Sanitizes and registers a data lake measure.
+   * This method should perform any necessary data validation and cleanup 
operations
+   * before registering the measure in the data lake.
+   *
+   * @return The sanitized and registered data lake measure.
+   */
+  DataLakeMeasure sanitizeAndRegister();
+
+  /**
+   * Sanitizes and updates a data lake measure.
+   * This method should perform any necessary data validation and cleanup 
operations
+   * before updating the measure in the data lake.
+   *
+   * @return The sanitized and updated data lake measure.
+   */
+  DataLakeMeasure sanitizeAndUpdate();
+
+}
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/package-info.java
 
b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/ITimeSeriesStorage.java
similarity index 65%
copy from 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/package-info.java
copy to 
streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/ITimeSeriesStorage.java
index c96a83d257..f3299db16a 100644
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/package-info.java
+++ 
b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/ITimeSeriesStorage.java
@@ -16,10 +16,14 @@
  *
  */
 
-/**
- * This package contains classes that have direct dependencies on 
InfluxDB-specific implementations.
- * To enhance code organization and promote agnostic design, these classes are 
slated for migration to other modules,
- * such as 'streampipes-data-explorer' or 'streampipes-data-explorer-api'.
- * Pending migration, they are temporarily housed within this package.
- */
-package org.apache.streampipes.dataexplorer.influx.migrate;
\ No newline at end of file
+package org.apache.streampipes.dataexplorer.api;
+
+import org.apache.streampipes.commons.exceptions.SpRuntimeException;
+import org.apache.streampipes.model.runtime.Event;
+
+public interface ITimeSeriesStorage {
+
+  void onEvent(Event event) throws SpRuntimeException;
+
+  void close() throws SpRuntimeException;
+}
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerInfluxQueryExecutor.java
 
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerInfluxQueryExecutor.java
index a500a4bd11..b24d2e4a88 100644
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerInfluxQueryExecutor.java
+++ 
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerInfluxQueryExecutor.java
@@ -19,7 +19,6 @@
 package org.apache.streampipes.dataexplorer.influx;
 
 import org.apache.streampipes.dataexplorer.influx.client.InfluxClientProvider;
-import org.apache.streampipes.dataexplorer.influx.migrate.DeleteDataQuery;
 import org.apache.streampipes.dataexplorer.param.DeleteQueryParams;
 import org.apache.streampipes.dataexplorer.param.SelectQueryParams;
 import org.apache.streampipes.dataexplorer.api.IDataLakeQueryBuilder;
@@ -37,6 +36,7 @@ import java.util.Arrays;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.Optional;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicLong;
 import java.util.stream.Collectors;
@@ -45,18 +45,6 @@ import static 
org.apache.streampipes.commons.environment.Environments.getEnviron
 
 public class DataExplorerInfluxQueryExecutor extends 
DataExplorerQueryExecutor<Query, QueryResult> {
 
-  public DataExplorerInfluxQueryExecutor() {
-    super();
-  }
-
-  public DataExplorerInfluxQueryExecutor(String forId) {
-    super(forId);
-  }
-
-  public DataExplorerInfluxQueryExecutor(int maximumAmountOfEvents) {
-    super(maximumAmountOfEvents);
-  }
-
   protected DataSeries convertResult(QueryResult.Series series,
                                      boolean ignoreMissingValues) {
     List<String> columns = series.getColumns();
@@ -79,6 +67,7 @@ public class DataExplorerInfluxQueryExecutor extends 
DataExplorerQueryExecutor<Q
   }
 
   protected SpQueryResult postQuery(QueryResult queryResult,
+                                    Optional<String> forIdOpt,
                                     boolean ignoreMissingValues) throws 
RuntimeException {
     SpQueryResult result = new SpQueryResult();
     AtomicLong lastTimestamp = new AtomicLong();
@@ -96,9 +85,7 @@ public class DataExplorerInfluxQueryExecutor extends 
DataExplorerQueryExecutor<Q
       result.setLastTimestamp(lastTimestamp.get());
     }
 
-    if (this.appendId) {
-      result.setForId(this.forId);
-    }
+    forIdOpt.ifPresent(result::setForId);
 
     return result;
   }
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerManagerInflux.java
 
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerManagerInflux.java
new file mode 100644
index 0000000000..e79cb8d2bc
--- /dev/null
+++ 
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerManagerInflux.java
@@ -0,0 +1,71 @@
+/*
+ * 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.dataexplorer.influx;
+
+import org.apache.streampipes.client.api.IStreamPipesClient;
+import org.apache.streampipes.commons.environment.Environments;
+import org.apache.streampipes.dataexplorer.DataExplorerSchemaManagement;
+import org.apache.streampipes.dataexplorer.api.IDataExplorerQueryManagement;
+import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement;
+import org.apache.streampipes.dataexplorer.api.IDataExplorerManager;
+import org.apache.streampipes.dataexplorer.api.IDataLakeMeasurementCounter;
+import org.apache.streampipes.dataexplorer.api.IDataLakeMeasurementSanitizer;
+import org.apache.streampipes.dataexplorer.api.ITimeSeriesStorage;
+import org.apache.streampipes.dataexplorer.influx.client.InfluxClientProvider;
+import 
org.apache.streampipes.dataexplorer.influx.sanitize.DataLakeMeasurementSanitizerInflux;
+import org.apache.streampipes.model.datalake.DataLakeMeasure;
+import org.apache.streampipes.storage.management.StorageDispatcher;
+
+import java.util.List;
+
+public enum DataExplorerManagerInflux implements IDataExplorerManager {
+
+  INSTANCE;
+
+  @Override
+  public IDataLakeMeasurementCounter getMeasurementCounter(
+    List<DataLakeMeasure> allMeasurements,
+    List<String> measurementsToCount) {
+    return new DataLakeMeasurementCounterInflux(allMeasurements, 
measurementsToCount);
+  }
+
+  @Override
+  public IDataExplorerQueryManagement getQueryManagement(
+    IDataExplorerSchemaManagement dataExplorerSchemaManagement
+    ) {
+    return new DataExplorerQueryManagementInflux(dataExplorerSchemaManagement);
+  }
+
+  @Override
+  public IDataExplorerSchemaManagement getSchemaManagement() {
+    return new DataExplorerSchemaManagement(StorageDispatcher.INSTANCE
+                                                             .getNoSqlStore()
+                                                             
.getDataLakeStorage());
+  }
+
+  @Override
+  public ITimeSeriesStorage getTimeseriesStorage(DataLakeMeasure measure) {
+    return new TimeSeriesStorageInflux(measure, Environments.getEnvironment(), 
new InfluxClientProvider());
+  }
+
+  @Override
+  public IDataLakeMeasurementSanitizer 
getMeasurementSanitizer(IStreamPipesClient client, DataLakeMeasure measure) {
+    return new DataLakeMeasurementSanitizerInflux(client, measure);
+  }
+}
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/DataExplorerQueryManagement.java
 
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerQueryManagementInflux.java
similarity index 78%
rename from 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/DataExplorerQueryManagement.java
rename to 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerQueryManagementInflux.java
index d9dcc5c284..44b9374b80 100644
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/DataExplorerQueryManagement.java
+++ 
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerQueryManagementInflux.java
@@ -16,11 +16,13 @@
  *
  */
 
-package org.apache.streampipes.dataexplorer.influx.migrate;
+package org.apache.streampipes.dataexplorer.influx;
 
+import org.apache.streampipes.dataexplorer.api.IDataExplorerQueryManagement;
 import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement;
 import org.apache.streampipes.dataexplorer.export.OutputFormat;
-import 
org.apache.streampipes.dataexplorer.influx.DataExplorerInfluxQueryExecutor;
+import org.apache.streampipes.dataexplorer.QueryResultProvider;
+import org.apache.streampipes.dataexplorer.StreamedQueryResultProvider;
 import org.apache.streampipes.dataexplorer.param.DeleteQueryParams;
 import 
org.apache.streampipes.dataexplorer.param.ProvidedRestQueryParamConverter;
 import org.apache.streampipes.model.datalake.param.ProvidedRestQueryParams;
@@ -32,18 +34,22 @@ import java.io.OutputStream;
 import java.util.List;
 import java.util.Map;
 
-public class DataExplorerQueryManagement implements 
IDataExplorerQueryManagement {
+public class DataExplorerQueryManagementInflux implements 
IDataExplorerQueryManagement {
 
   private final IDataExplorerSchemaManagement dataExplorerSchemaManagement;
 
-  public DataExplorerQueryManagement(IDataExplorerSchemaManagement 
dataExplorerSchemaManagement) {
+  public DataExplorerQueryManagementInflux(IDataExplorerSchemaManagement 
dataExplorerSchemaManagement) {
     this.dataExplorerSchemaManagement = dataExplorerSchemaManagement;
   }
 
   @Override
   public SpQueryResult getData(ProvidedRestQueryParams queryParams,
                                boolean ignoreMissingData) throws 
IllegalArgumentException {
-    return new QueryResultProvider(queryParams, ignoreMissingData).getData();
+    return new QueryResultProvider(queryParams,
+                                   this,
+                                   new DataExplorerInfluxQueryExecutor(),
+                                   ignoreMissingData
+    ).getData();
   }
 
   @Override
@@ -52,7 +58,11 @@ public class DataExplorerQueryManagement implements 
IDataExplorerQueryManagement
                               boolean ignoreMissingValues,
                               OutputStream outputStream) throws IOException {
 
-    new StreamedQueryResultProvider(params, format, 
ignoreMissingValues).getDataAsStream(outputStream);
+    new StreamedQueryResultProvider(params, format,
+                                    this,
+                                    new DataExplorerInfluxQueryExecutor(),
+                                    ignoreMissingValues
+    ).getDataAsStream(outputStream);
   }
 
   @Override
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/DataLakeMeasurementCount.java
 
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataLakeMeasurementCounterInflux.java
similarity index 88%
rename from 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/DataLakeMeasurementCount.java
rename to 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataLakeMeasurementCounterInflux.java
index aa21d051fa..4f516f52c1 100644
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/DataLakeMeasurementCount.java
+++ 
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataLakeMeasurementCounterInflux.java
@@ -16,10 +16,9 @@
  *
  */
 
-package org.apache.streampipes.dataexplorer.influx.migrate;
+package org.apache.streampipes.dataexplorer.influx;
 
-import 
org.apache.streampipes.dataexplorer.influx.DataExplorerInfluxQueryExecutor;
-import org.apache.streampipes.dataexplorer.influx.DataLakeInfluxQueryBuilder;
+import org.apache.streampipes.dataexplorer.api.IDataLakeMeasurementCounter;
 import org.apache.streampipes.model.datalake.AggregationFunction;
 import org.apache.streampipes.model.datalake.DataLakeMeasure;
 import org.apache.streampipes.model.datalake.SpQueryResult;
@@ -29,19 +28,20 @@ import org.apache.streampipes.model.schema.PropertyScope;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.Optional;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.ExecutionException;
 import java.util.stream.Collectors;
 
-public class DataLakeMeasurementCount {
+public class DataLakeMeasurementCounterInflux implements 
IDataLakeMeasurementCounter {
 
   private final List<DataLakeMeasure> allMeasurements;
   private final List<String> measurementNames;
 
   private static final String COUNT_FIELD = "count";
 
-  public DataLakeMeasurementCount(List<DataLakeMeasure> allMeasurements,
-                                  List<String> measurementNames) {
+  public DataLakeMeasurementCounterInflux(List<DataLakeMeasure> 
allMeasurements,
+                                          List<String> measurementNames) {
     this.allMeasurements = allMeasurements;
     this.measurementNames = measurementNames;
   }
@@ -55,7 +55,7 @@ public class DataLakeMeasurementCount {
           var builder = DataLakeInfluxQueryBuilder
               
.create(m.getMeasureName()).withEndTime(System.currentTimeMillis())
               .withAggregatedColumn(firstColumn, AggregationFunction.COUNT);
-          var queryResult = new 
DataExplorerInfluxQueryExecutor().executeQuery(builder.build(), true);
+          var queryResult = new 
DataExplorerInfluxQueryExecutor().executeQuery(builder.build(), 
Optional.empty(), true);
           if (queryResult.getTotal() > 0) {
             var headers = queryResult.getHeaders();
             return extractResult(queryResult, headers);
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/DeleteDataQuery.java
 
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DeleteDataQuery.java
similarity index 96%
rename from 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/DeleteDataQuery.java
rename to 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DeleteDataQuery.java
index 4885cd1ed9..63eb528d25 100644
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/DeleteDataQuery.java
+++ 
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DeleteDataQuery.java
@@ -15,7 +15,7 @@
  * limitations under the License.
  *
  */
-package org.apache.streampipes.dataexplorer.influx.migrate;
+package org.apache.streampipes.dataexplorer.influx;
 
 import org.apache.streampipes.commons.environment.Environment;
 import org.apache.streampipes.commons.environment.Environments;
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/InfluxStore.java
 
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/InfluxStore.java
deleted file mode 100644
index 0d10ac147f..0000000000
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/InfluxStore.java
+++ /dev/null
@@ -1,244 +0,0 @@
-/*
- * 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.dataexplorer.influx;
-
-import org.apache.streampipes.commons.environment.Environment;
-import org.apache.streampipes.commons.exceptions.SpRuntimeException;
-import org.apache.streampipes.dataexplorer.influx.client.InfluxClientProvider;
-import org.apache.streampipes.dataexplorer.influx.sanitize.InfluxNameSanitizer;
-import org.apache.streampipes.model.datalake.DataLakeMeasure;
-import org.apache.streampipes.model.runtime.Event;
-import org.apache.streampipes.model.schema.EventProperty;
-import org.apache.streampipes.model.schema.EventPropertyPrimitive;
-
-import org.influxdb.InfluxDB;
-import org.influxdb.dto.Point;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.util.HashMap;
-import java.util.List;
-import java.util.Map;
-import java.util.concurrent.TimeUnit;
-import java.util.stream.Collectors;
-
-public class InfluxStore {
-
-  private static final Logger LOG = LoggerFactory.getLogger(InfluxStore.class);
-
-  private final DataLakeMeasure measure;
-  private final List<EventProperty> allEventProperties;
-  private final Map<String, String> sanitizedRuntimeNames = new HashMap<>();
-  private final InfluxDB influxDb;
-
-  private final PropertyHandler propertyHandler;
-
-
-  public InfluxStore(
-      DataLakeMeasure measure,
-      Environment environment,
-      InfluxClientProvider influxClientProvider
-  ) throws SpRuntimeException {
-    this.measure = measure;
-    storeSanitizedRuntimeNames(measure);
-    allEventProperties = getAllEventPropertiesExceptTimestamp(measure);
-    influxDb = influxClientProvider.getInitializedInfluxDBClient(environment);
-    propertyHandler = new PropertyHandler();
-  }
-
-  /**
-   * Takes an StreamPipes event, transforms it to an InfluxDB point and writes 
it to the InfluxDB
-   *
-   * @param event The event which should be saved
-   * @throws SpRuntimeException If the column name (key-value of the event 
map) is not allowed
-   */
-  public void onEvent(Event event) throws SpRuntimeException {
-
-    validateInputEventAndLogMissingFields(event);
-
-    sanitizeRuntimeNamesInEvent(event);
-
-    var point = initializePointWithTimestamp(event);
-
-    iterateOverallEventProperties(event, point);
-
-    influxDb.write(point.build());
-  }
-
-  private void validateInputEventAndLogMissingFields(Event event) {
-    checkEventIsNotNull(event);
-
-    logMissingFields(event);
-
-    logNullFields(event);
-  }
-
-  /**
-   * Logs all fields which are present in the schema, but not in the provided 
event
-   */
-  private void logMissingFields(Event event) {
-    var missingFields = getMissingProperties(allEventProperties, event);
-    if (!missingFields.isEmpty()) {
-      LOG.debug(
-          "Ignored {} fields which were present in the schema, but not in the 
provided event: {}",
-          missingFields.size(),
-          String.join(", ", missingFields)
-      );
-    }
-  }
-
-
-  /**
-   * Logs all fields that contain null values
-   */
-  private void logNullFields(Event event) {
-    List<String> nullFields = allEventProperties
-        .stream()
-        .filter(EventPropertyPrimitive.class::isInstance)
-        .filter(ep -> {
-          var runtimeName = ep.getRuntimeName();
-          var field = event.getOptionalFieldByRuntimeName(runtimeName);
-
-          return field.isPresent() && field.get()
-                                           .getAsPrimitive()
-                                           .getRawValue() == null;
-        })
-        .map(EventProperty::getRuntimeName)
-        .collect(Collectors.toList());
-
-    if (!nullFields.isEmpty()) {
-      LOG.warn("Ignored {} fields which had a value 'null': {}", 
nullFields.size(), String.join(", ", nullFields));
-    }
-  }
-
-  private void iterateOverallEventProperties(
-      Event event,
-      Point.Builder point
-  ) {
-
-    allEventProperties.forEach(ep -> {
-      var runtimeName = ep.getRuntimeName();
-      var sanitizedRuntimeName = sanitizedRuntimeNames.get(runtimeName);
-      var fieldOptional = event.getOptionalFieldByRuntimeName(runtimeName);
-
-      fieldOptional.ifPresent(field -> {
-        if (ep instanceof EventPropertyPrimitive) {
-          propertyHandler.handlePrimitiveProperty(
-              point,
-              (EventPropertyPrimitive) ep,
-              field.getAsPrimitive(),
-              sanitizedRuntimeName
-          );
-        } else {
-          propertyHandler.handleNonPrimitiveProperty(
-              point,
-              event,
-              sanitizedRuntimeName
-          );
-        }
-      });
-    });
-  }
-
-
-  /**
-   * Returns a list of the runtime names that are missing within the event
-   */
-  private List<String> getMissingProperties(
-      List<EventProperty> allEventProperties,
-      Event event
-  ) {
-    return allEventProperties.stream()
-                             .map(EventProperty::getRuntimeName)
-                             .filter(runtimeName -> 
event.getOptionalFieldByRuntimeName(runtimeName)
-                                                          .isEmpty())
-                             .collect(Collectors.toList());
-  }
-
-  private void checkEventIsNotNull(Event event) {
-    if (event == null) {
-      throw new SpRuntimeException("event is null");
-    }
-  }
-
-
-  /**
-   * Shuts down the connection to the InfluxDB server
-   */
-  public void close() throws SpRuntimeException {
-    influxDb.flush();
-    try {
-      Thread.sleep(1000);
-    } catch (InterruptedException e) {
-      throw new SpRuntimeException(e);
-    }
-    influxDb.close();
-  }
-
-
-  /**
-   * Creates a point object which is later written to the influxDB and adds 
the value of the timestamp field
-   */
-  private Point.Builder initializePointWithTimestamp(Event event) {
-    var timestampValue = event.getFieldBySelector(measure.getTimestampField())
-                              .getAsPrimitive()
-                              .getAsLong();
-    return Point.measurement(measure.getMeasureName())
-                .time((long) timestampValue, TimeUnit.MILLISECONDS);
-  }
-
-  /**
-   * store sanitized target property runtime names in local variable
-   */
-  private void storeSanitizedRuntimeNames(DataLakeMeasure measure) {
-    measure.getEventSchema()
-           .getEventProperties()
-           .forEach(ep -> sanitizedRuntimeNames.put(
-             ep.getRuntimeName(),
-             InfluxNameSanitizer.renameReservedKeywords(ep.getRuntimeName())
-           ));
-  }
-
-  /**
-   * Returns all measurements properties except the timestamp field
-   */
-  private List<EventProperty> 
getAllEventPropertiesExceptTimestamp(DataLakeMeasure dataLakeMeasure) {
-    return dataLakeMeasure.getEventSchema()
-                          .getEventProperties()
-                          .stream()
-                          .filter(ep -> !dataLakeMeasure.getTimestampField()
-                                                        
.endsWith(ep.getRuntimeName()))
-                          .collect(Collectors.toList());
-  }
-
-  /**
-   * Iterates over all properties of the event and renames the key if it is a 
reserved keywords in InfluxDB
-   */
-  private void sanitizeRuntimeNamesInEvent(Event event) {
-    // sanitize event
-    for (var key : event.getRaw()
-                        .keySet()) {
-      if (InfluxDbReservedKeywords.KEYWORD_LIST.stream()
-                                               .anyMatch(k -> 
k.equalsIgnoreCase(key))) {
-        event.renameFieldByRuntimeName(key, key + "_");
-      }
-    }
-  }
-
-}
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/TimeSeriesStorageInflux.java
 
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/TimeSeriesStorageInflux.java
new file mode 100644
index 0000000000..8c2f674a99
--- /dev/null
+++ 
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/TimeSeriesStorageInflux.java
@@ -0,0 +1,132 @@
+/*
+ * 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.dataexplorer.influx;
+
+import org.apache.streampipes.commons.environment.Environment;
+import org.apache.streampipes.commons.exceptions.SpRuntimeException;
+import org.apache.streampipes.dataexplorer.TimeSeriesStorage;
+import org.apache.streampipes.dataexplorer.influx.client.InfluxClientProvider;
+import org.apache.streampipes.dataexplorer.influx.sanitize.InfluxNameSanitizer;
+import org.apache.streampipes.model.datalake.DataLakeMeasure;
+import org.apache.streampipes.model.runtime.Event;
+import org.apache.streampipes.model.schema.EventPropertyPrimitive;
+
+import org.influxdb.InfluxDB;
+import org.influxdb.dto.Point;
+
+import java.util.concurrent.TimeUnit;
+
+public class TimeSeriesStorageInflux extends TimeSeriesStorage {
+
+  private final InfluxDB influxDb;
+
+  private final PropertyHandler propertyHandler;
+
+
+  public TimeSeriesStorageInflux(
+      DataLakeMeasure measure,
+      Environment environment,
+      InfluxClientProvider influxClientProvider
+  ) throws SpRuntimeException {
+    super(measure);
+    influxDb = influxClientProvider.getInitializedInfluxDBClient(environment);
+    propertyHandler = new PropertyHandler();
+  }
+
+  protected void writeToTimeSeriesStorage(Event event) throws 
SpRuntimeException {
+    var point = initializePointWithTimestamp(event);
+    iterateOverallEventProperties(event, point);
+    influxDb.write(point.build());
+  }
+
+  private void iterateOverallEventProperties(
+      Event event,
+      Point.Builder point
+  ) {
+
+    allEventProperties.forEach(ep -> {
+      var runtimeName = ep.getRuntimeName();
+      var sanitizedRuntimeName = sanitizedRuntimeNames.get(runtimeName);
+      var fieldOptional = event.getOptionalFieldByRuntimeName(runtimeName);
+
+      fieldOptional.ifPresent(field -> {
+        if (ep instanceof EventPropertyPrimitive) {
+          propertyHandler.handlePrimitiveProperty(
+              point,
+              (EventPropertyPrimitive) ep,
+              field.getAsPrimitive(),
+              sanitizedRuntimeName
+          );
+        } else {
+          propertyHandler.handleNonPrimitiveProperty(
+              point,
+              event,
+              sanitizedRuntimeName
+          );
+        }
+      });
+    });
+  }
+
+  /**
+   * Shuts down the connection to the InfluxDB server
+   */
+  public void close() throws SpRuntimeException {
+    influxDb.flush();
+    try {
+      Thread.sleep(1000);
+    } catch (InterruptedException e) {
+      throw new SpRuntimeException(e);
+    }
+    influxDb.close();
+  }
+
+  /**
+   * Creates a point object which is later written to the influxDB and adds 
the value of the timestamp field
+   */
+  private Point.Builder initializePointWithTimestamp(Event event) {
+    var timestampValue = event.getFieldBySelector(measure.getTimestampField())
+        .getAsPrimitive()
+        .getAsLong();
+    return Point.measurement(measure.getMeasureName())
+        .time((long) timestampValue, TimeUnit.MILLISECONDS);
+  }
+
+  /**
+   * store sanitized target property runtime names in local variable
+   */
+  protected void storeSanitizedRuntimeNames() {
+    measure.getEventSchema()
+           .getEventProperties()
+           .forEach(ep -> sanitizedRuntimeNames.put(
+             ep.getRuntimeName(),
+             InfluxNameSanitizer.renameReservedKeywords(ep.getRuntimeName())
+           ));
+  }
+
+  /**
+   * Iterates over all properties of the event and renames the key if it is a 
reserved keywords in InfluxDB
+   */
+  protected void sanitizeRuntimeNamesInEvent(Event event) {
+    // sanitize event
+    event.getRaw()
+         .keySet()
+         .forEach(key -> event.renameFieldByRuntimeName(key, 
InfluxNameSanitizer.renameReservedKeywords(key)));
+  }
+}
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/DataExplorerUtils.java
 
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/DataExplorerUtils.java
deleted file mode 100644
index 97d8b9089d..0000000000
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/DataExplorerUtils.java
+++ /dev/null
@@ -1,86 +0,0 @@
-/*
- * 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.dataexplorer.influx.migrate;
-
-import org.apache.streampipes.client.api.IStreamPipesClient;
-import org.apache.streampipes.commons.exceptions.SpRuntimeException;
-import org.apache.streampipes.dataexplorer.influx.sanitize.InfluxNameSanitizer;
-import 
org.apache.streampipes.dataexplorer.influx.sanitize.MeasureNameSanitizer;
-import org.apache.streampipes.model.datalake.DataLakeMeasure;
-import org.apache.streampipes.model.schema.EventProperty;
-
-import java.util.List;
-import java.util.stream.Collectors;
-
-public class DataExplorerUtils {
-  /**
-   * Sanitizes the event schema and stores the DataLakeMeasurement to the 
couchDB
-   *
-   * @param client  StreamPipes client to store measure
-   * @param measure DataLakeMeasurement
-   */
-  public static DataLakeMeasure 
sanitizeAndRegisterAtDataLake(IStreamPipesClient client,
-                                                              DataLakeMeasure 
measure) throws SpRuntimeException {
-    sanitizeDataLakeMeasure(measure);
-    registerAtDataLake(client, measure);
-
-    return measure;
-  }
-
-  public static DataLakeMeasure sanitizeAndUpdateAtDataLake(IStreamPipesClient 
client,
-                                                            DataLakeMeasure 
measure) throws SpRuntimeException {
-    sanitizeDataLakeMeasure(measure);
-    updateAtDataLake(client, measure);
-    return measure;
-  }
-
-  private static void registerAtDataLake(IStreamPipesClient client,
-                                         DataLakeMeasure measure) throws 
SpRuntimeException {
-    client.dataLakeMeasureApi().create(measure);
-  }
-
-  public static void updateAtDataLake(IStreamPipesClient client,
-                                      DataLakeMeasure measure) throws 
SpRuntimeException {
-    client.dataLakeMeasureApi().update(measure);
-  }
-
-
-  private static void sanitizeDataLakeMeasure(DataLakeMeasure measure) throws 
SpRuntimeException {
-
-    // Removes selected timestamp from event schema
-    removeTimestampsFromEventSchema(measure);
-
-    // Sanitize the data lake measure name
-    measure.setMeasureName(new 
MeasureNameSanitizer().sanitize(measure.getMeasureName()));
-
-    // Removes all spaces with _ and validates that no special terms are used 
as runtime names
-    measure.getEventSchema()
-        .getEventProperties()
-        .forEach(ep -> 
ep.setRuntimeName(InfluxNameSanitizer.renameReservedKeywords(ep.getRuntimeName())));
-
-  }
-
-  private static void removeTimestampsFromEventSchema(DataLakeMeasure measure) 
{
-    List<EventProperty> eventPropertiesWithoutTimestamp = 
measure.getEventSchema().getEventProperties()
-        .stream()
-        .filter(eventProperty -> 
!measure.getTimestampField().endsWith(eventProperty.getRuntimeName()))
-        .collect(Collectors.toList());
-    
measure.getEventSchema().setEventProperties(eventPropertiesWithoutTimestamp);
-  }
-
-}
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/sanitize/DataLakeMeasurementSanitizerInflux.java
 
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/sanitize/DataLakeMeasurementSanitizerInflux.java
new file mode 100644
index 0000000000..29fd6f4ccc
--- /dev/null
+++ 
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/sanitize/DataLakeMeasurementSanitizerInflux.java
@@ -0,0 +1,44 @@
+/*
+ * 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.dataexplorer.influx.sanitize;
+
+import org.apache.streampipes.client.api.IStreamPipesClient;
+import org.apache.streampipes.commons.exceptions.SpRuntimeException;
+import org.apache.streampipes.dataexplorer.DataLakeMeasurementSanitizer;
+import org.apache.streampipes.model.datalake.DataLakeMeasure;
+
+public class DataLakeMeasurementSanitizerInflux extends 
DataLakeMeasurementSanitizer {
+  public DataLakeMeasurementSanitizerInflux(
+    IStreamPipesClient client,
+    DataLakeMeasure measure
+  ) {
+    super(client, measure);
+  }
+
+  @Override
+  protected void cleanDataLakeMeasure() throws SpRuntimeException {
+    // Sanitize the data lake measure name
+    measure.setMeasureName(new 
MeasureNameSanitizer().sanitize(measure.getMeasureName()));
+
+    // Removes all spaces with _ and validates that no special terms are used 
as runtime names
+    measure.getEventSchema()
+           .getEventProperties()
+           .forEach(ep -> 
ep.setRuntimeName(InfluxNameSanitizer.renameReservedKeywords(ep.getRuntimeName())));
+  }
+}
diff --git 
a/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/InfluxStoreTest.java
 
b/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/TimeSeriesStorageInfluxTest.java
similarity index 97%
rename from 
streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/InfluxStoreTest.java
rename to 
streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/TimeSeriesStorageInfluxTest.java
index aec50f849d..3ddd6f7007 100644
--- 
a/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/InfluxStoreTest.java
+++ 
b/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/TimeSeriesStorageInfluxTest.java
@@ -53,7 +53,7 @@ import java.util.concurrent.TimeUnit;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 
-public class InfluxStoreTest {
+public class TimeSeriesStorageInfluxTest {
 
   private InfluxDB influxDBMock;
 
@@ -300,8 +300,8 @@ public class InfluxStoreTest {
   /**
    * Executes the onEvent method and returns the resuting data point using an 
argument captor
    */
-  private Point executeOnEvent(InfluxStore influxStore, Event event) {
-    influxStore.onEvent(event);
+  private Point executeOnEvent(TimeSeriesStorageInflux 
timeSeriesStorageInflux, Event event) {
+    timeSeriesStorageInflux.onEvent(event);
     var pointArgumentCaptor = ArgumentCaptor.forClass(Point.class);
 
     Mockito.verify(influxDBMock).write(pointArgumentCaptor.capture());
@@ -374,7 +374,7 @@ public class InfluxStoreTest {
   /**
    * Initializes an influx store with the given event schema
    */
-  private InfluxStore getInfluxStore(EventSchema eventSchema) {
+  private TimeSeriesStorageInflux getInfluxStore(EventSchema eventSchema) {
 
     DataLakeMeasure measure = new DataLakeMeasure(
         EXPECTED_MEASUREMENT,
@@ -386,7 +386,7 @@ public class InfluxStoreTest {
     
Mockito.when(influxClientProviderMock.getInitializedInfluxDBClient(ArgumentMatchers.any()))
            .thenReturn(influxDBMock);
 
-    return new InfluxStore(measure, null, influxClientProviderMock);
+    return new TimeSeriesStorageInflux(measure, null, 
influxClientProviderMock);
   }
 
 }
\ No newline at end of file
diff --git 
a/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/sanitize/DataLakeMeasurementSanitizerInfluxTest.java
 
b/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/sanitize/DataLakeMeasurementSanitizerInfluxTest.java
new file mode 100644
index 0000000000..5d6cb0542d
--- /dev/null
+++ 
b/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/sanitize/DataLakeMeasurementSanitizerInfluxTest.java
@@ -0,0 +1,93 @@
+/*
+ * 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.dataexplorer.influx.sanitize;
+
+import org.apache.streampipes.client.api.IDataLakeMeasureApi;
+import org.apache.streampipes.client.api.IStreamPipesClient;
+import org.apache.streampipes.commons.exceptions.SpRuntimeException;
+import org.apache.streampipes.model.datalake.DataLakeMeasure;
+import org.apache.streampipes.test.generator.EventPropertyPrimitiveTestBuilder;
+import org.apache.streampipes.test.generator.EventSchemaTestBuilder;
+
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.util.List;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+public class DataLakeMeasurementSanitizerInfluxTest {
+
+  private IStreamPipesClient clientMock;
+
+  @BeforeEach
+  public void setUp() {
+    clientMock = mock(IStreamPipesClient.class);
+
+    var apiMock = mock(IDataLakeMeasureApi.class);
+    when(clientMock.dataLakeMeasureApi()).thenReturn(apiMock);
+
+  }
+
+  @Test
+  public void cleanDataLakeMeasure() {
+    var measure = new DataLakeMeasure(
+      "test?Measurement",
+      "timestamp",
+      EventSchemaTestBuilder.create()
+                            .withEventProperties(List.of(
+                              EventPropertyPrimitiveTestBuilder.create()
+                                                               
.withRuntimeName("timestamp")
+                                                               .build(),
+                              EventPropertyPrimitiveTestBuilder.create()
+                                                               
.withRuntimeName("value")
+                                                               .build(),
+                              EventPropertyPrimitiveTestBuilder.create()
+                                                               
.withRuntimeName("all")
+                                                               .build()
+                            ))
+                            .build()
+    );
+
+    var result = new DataLakeMeasurementSanitizerInflux(clientMock, 
measure).sanitizeAndRegister();
+
+    assertEquals("test_Measurement", result.getMeasureName());
+    assertEquals(2, result.getEventSchema().getEventProperties().size());
+    assertEquals("value", 
result.getEventSchema().getEventProperties().get(0).getRuntimeName());
+    assertEquals("all_", 
result.getEventSchema().getEventProperties().get(1).getRuntimeName());
+  }
+
+  @Test
+  public void cleanDataLakeMeasureNoTimestampField() {
+    var measure = new DataLakeMeasure("test", EventSchemaTestBuilder.create()
+                                                                    
.withEventProperties(List.of(
+                                                                      
EventPropertyPrimitiveTestBuilder.create()
+                                                                               
                        .withRuntimeName(
+                                                                               
                          "value")
+                                                                               
                        .build()
+                                                                    ))
+                                                                    .build());
+
+    assertThrows(SpRuntimeException.class,
+                 () -> new DataLakeMeasurementSanitizerInflux(clientMock, 
measure).sanitizeAndRegister());
+  }
+}
diff --git a/streampipes-data-explorer-management/pom.xml 
b/streampipes-data-explorer-management/pom.xml
index 09c89dacff..56edc05e4e 100644
--- a/streampipes-data-explorer-management/pom.xml
+++ b/streampipes-data-explorer-management/pom.xml
@@ -44,6 +44,11 @@
             <artifactId>streampipes-data-explorer-api</artifactId>
             <version>0.95.0-SNAPSHOT</version>
         </dependency>
+        <dependency>
+            <groupId>org.apache.streampipes</groupId>
+            <artifactId>streampipes-data-explorer-influx</artifactId>
+            <version>0.95.0-SNAPSHOT</version>
+        </dependency>
     </dependencies>
 
     <properties>
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/ImageStoreUtils.java
 
b/streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/DataExplorerDispatcher.java
similarity index 50%
copy from 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/ImageStoreUtils.java
copy to 
streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/DataExplorerDispatcher.java
index 6504303a5a..96ba5b70ae 100644
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/ImageStoreUtils.java
+++ 
b/streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/DataExplorerDispatcher.java
@@ -16,22 +16,23 @@
  *
  */
 
-package org.apache.streampipes.dataexplorer.influx.migrate;
+package org.apache.streampipes.dataexplorer.management;
 
-import org.apache.streampipes.model.datalake.DataLakeMeasure;
-import org.apache.streampipes.model.schema.EventProperty;
-import org.apache.streampipes.vocabulary.SPSensor;
+import org.apache.streampipes.commons.environment.Environments;
+import org.apache.streampipes.dataexplorer.api.IDataExplorerManager;
+import org.apache.streampipes.dataexplorer.influx.DataExplorerManagerInflux;
 
-import java.util.List;
-import java.util.stream.Collectors;
+public class DataExplorerDispatcher {
 
-public class ImageStoreUtils {
+  public IDataExplorerManager getDataExplorerManager() {
 
-  public static List<EventProperty> getImageProperties(DataLakeMeasure 
measure) {
-    return measure.getEventSchema().getEventProperties().stream()
-        .filter(eventProperty -> eventProperty.getDomainProperties() != null
-            && eventProperty.getDomainProperties().size() > 0
-            && 
eventProperty.getDomainProperties().get(0).toString().equals(SPSensor.IMAGE))
-        .collect(Collectors.toList());
+    // currently this SWITCH CASE statement is not necessary
+    // but aims to give an idea how to deal with multiple data explorer 
storages
+    return switch (Environments.getEnvironment()
+                               .getTsStorage()
+                               .getValueOrDefault()) {
+      case SupportedDataExplorerStorages.INFLUX_DB -> 
DataExplorerManagerInflux.INSTANCE;
+      default -> DataExplorerManagerInflux.INSTANCE;
+    };
   }
 }
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/package-info.java
 
b/streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/SupportedDataExplorerStorages.java
similarity index 65%
rename from 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/package-info.java
rename to 
streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/SupportedDataExplorerStorages.java
index c96a83d257..70af654be8 100644
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/package-info.java
+++ 
b/streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/SupportedDataExplorerStorages.java
@@ -16,10 +16,16 @@
  *
  */
 
+package org.apache.streampipes.dataexplorer.management;
+
 /**
- * This package contains classes that have direct dependencies on 
InfluxDB-specific implementations.
- * To enhance code organization and promote agnostic design, these classes are 
slated for migration to other modules,
- * such as 'streampipes-data-explorer' or 'streampipes-data-explorer-api'.
- * Pending migration, they are temporarily housed within this package.
+ * Class containing constants for supported time series storage 
implementations.
+ * <p>
+ * Supported time series storage implementations:
+ * <ul>
+ *     <li>{@link #INFLUX_DB}: Represents the InfluxDB time series storage 
implementation.</li>
+ * </ul>
  */
-package org.apache.streampipes.dataexplorer.influx.migrate;
\ No newline at end of file
+public class SupportedDataExplorerStorages {
+  public static final String INFLUX_DB = "influxdb";
+}
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/AutoAggregationHandler.java
 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/AutoAggregationHandler.java
similarity index 90%
rename from 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/AutoAggregationHandler.java
rename to 
streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/AutoAggregationHandler.java
index 692e7412be..c0751c7fee 100644
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/AutoAggregationHandler.java
+++ 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/AutoAggregationHandler.java
@@ -15,15 +15,14 @@
  * limitations under the License.
  *
  */
-package org.apache.streampipes.dataexplorer.influx.migrate;
+package org.apache.streampipes.dataexplorer;
 
-import org.apache.streampipes.dataexplorer.DataExplorerSchemaManagement;
-import org.apache.streampipes.model.datalake.param.ProvidedRestQueryParams;
-import org.apache.streampipes.model.datalake.param.SupportedRestQueryParams;
+import org.apache.streampipes.dataexplorer.api.IDataExplorerQueryManagement;
 import org.apache.streampipes.dataexplorer.param.model.SelectColumn;
 import org.apache.streampipes.model.datalake.DataLakeQueryOrdering;
 import org.apache.streampipes.model.datalake.SpQueryResult;
-import org.apache.streampipes.storage.management.StorageDispatcher;
+import org.apache.streampipes.model.datalake.param.ProvidedRestQueryParams;
+import org.apache.streampipes.model.datalake.param.SupportedRestQueryParams;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -50,16 +49,10 @@ public class AutoAggregationHandler {
   private final IDataExplorerQueryManagement dataLakeQueryManagement;
   private final ProvidedRestQueryParams queryParams;
 
-  public AutoAggregationHandler(ProvidedRestQueryParams params) {
+  public AutoAggregationHandler(ProvidedRestQueryParams params,
+                                IDataExplorerQueryManagement 
dataExplorerQueryManagement) {
     this.queryParams = params;
-    this.dataLakeQueryManagement = getDataLakeQueryManagement();
-  }
-
-  private IDataExplorerQueryManagement getDataLakeQueryManagement() {
-    var dataLakeStorage = StorageDispatcher.INSTANCE
-        .getNoSqlStore()
-        .getDataLakeStorage();
-    return new DataExplorerQueryManagement(new 
DataExplorerSchemaManagement(dataLakeStorage));
+    this.dataLakeQueryManagement = dataExplorerQueryManagement;
   }
 
   public ProvidedRestQueryParams makeAutoAggregationQueryParams() throws 
IllegalArgumentException {
@@ -129,7 +122,7 @@ public class AutoAggregationHandler {
     singleEvent.update(SupportedRestQueryParams.QP_LIMIT, 1);
     singleEvent.update(SupportedRestQueryParams.QP_ORDER, order.name());
     singleEvent.update(SupportedRestQueryParams.QP_COLUMNS, 
transformColumns(singleEvent.getAsString(
-      SupportedRestQueryParams.QP_COLUMNS)));
+        SupportedRestQueryParams.QP_COLUMNS)));
 
     return fireQuery(singleEvent);
   }
diff --git 
a/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/DataLakeMeasurementSanitizer.java
 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/DataLakeMeasurementSanitizer.java
new file mode 100644
index 0000000000..9c8904ae8e
--- /dev/null
+++ 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/DataLakeMeasurementSanitizer.java
@@ -0,0 +1,118 @@
+/*
+ * 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.dataexplorer;
+
+import org.apache.streampipes.client.api.IStreamPipesClient;
+import org.apache.streampipes.commons.exceptions.SpRuntimeException;
+import org.apache.streampipes.dataexplorer.api.IDataLakeMeasurementSanitizer;
+import org.apache.streampipes.model.datalake.DataLakeMeasure;
+import org.apache.streampipes.model.schema.EventProperty;
+
+import java.util.List;
+
+/**
+ * Base class with shared implementation that is common for all time series 
storage backends.
+ * Leaves open the storage specific implementation
+ */
+public abstract class DataLakeMeasurementSanitizer implements 
IDataLakeMeasurementSanitizer {
+
+  protected final DataLakeMeasure measure;
+  protected final IStreamPipesClient client;
+
+  public DataLakeMeasurementSanitizer(IStreamPipesClient client, 
DataLakeMeasure measure){
+    this.client = client;
+    this.measure = measure;
+  }
+
+  /**
+   * Sanitizes the data lake measure and registers it with the data lake.
+   * <p>
+   * This method first sanitizes the data lake measure,
+   * then registers it at the data lake.
+   *
+   * @return The sanitized and registered data lake measure.
+   */
+  @Override
+  public DataLakeMeasure sanitizeAndRegister(){
+    sanitizeDataLakeMeasure();
+    registerAtDataLake();
+
+    return measure;
+  }
+
+  /**
+   * Sanitizes the data lake measure and updates it in the data lake.
+   * <p>
+   * This method first sanitizes the data lake measure,
+   * then updates it at the data lake.
+   *
+   * @return The sanitized and updated data lake measure.
+   */
+  @Override
+  public DataLakeMeasure sanitizeAndUpdate(){
+    sanitizeDataLakeMeasure();
+    updateAtDataLake();
+
+    return measure;
+  }
+
+
+
+  private void registerAtDataLake() throws SpRuntimeException {
+    client.dataLakeMeasureApi().create(measure);
+  }
+
+  private void updateAtDataLake() throws SpRuntimeException {
+    client.dataLakeMeasureApi().update(measure);
+  }
+
+  private void sanitizeDataLakeMeasure() throws SpRuntimeException {
+    removeTimestampsFromEventSchema();
+    cleanDataLakeMeasure();
+  }
+
+  /**
+   * Cleans the data lake measure to ensure compliance with the requirements 
of the respective time series storage.
+   * <p>
+   * This method performs the following steps:
+   * <ol>
+   *   <li>Sanitizes the name of the measure.</li>
+   *   <li>Sanitizes all runtime names associated with the measure.</li>
+   * </ol>
+   * @throws SpRuntimeException if an error occurs during the cleaning process.
+   */
+  protected abstract void cleanDataLakeMeasure() throws SpRuntimeException;
+
+  protected void removeTimestampsFromEventSchema() throws SpRuntimeException{
+    var timestampField = measure.getTimestampField();
+
+    if (timestampField == null){
+      throw new SpRuntimeException("Data lake measurement does not have a 
timestamp field - timestamp field is null.");
+    }
+
+    List<EventProperty> eventPropertiesWithoutTimestamp = 
measure.getEventSchema()
+                                                                 
.getEventProperties()
+                                                                 .stream()
+                                                                 
.filter(eventProperty -> !timestampField.endsWith(
+                                                                   
eventProperty.getRuntimeName()
+                                                                 ))
+                                                                 .toList();
+    
measure.getEventSchema().setEventProperties(eventPropertiesWithoutTimestamp);
+  }
+}
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/ImageStore.java
 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/ImageStore.java
similarity index 91%
rename from 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/ImageStore.java
rename to 
streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/ImageStore.java
index 6cfd2087be..0b961b6b50 100644
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/ImageStore.java
+++ 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/ImageStore.java
@@ -16,7 +16,7 @@
  *
  */
 
-package org.apache.streampipes.dataexplorer.influx.migrate;
+package org.apache.streampipes.dataexplorer;
 
 import org.apache.streampipes.commons.environment.Environment;
 import org.apache.streampipes.commons.exceptions.SpRuntimeException;
@@ -27,8 +27,6 @@ import org.apache.streampipes.model.schema.EventProperty;
 import org.apache.commons.codec.binary.Base64;
 import org.lightcouch.CouchDbClient;
 import org.lightcouch.CouchDbProperties;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
 
 import java.io.ByteArrayInputStream;
 import java.io.IOException;
@@ -36,12 +34,10 @@ import java.util.List;
 import java.util.UUID;
 
 public class ImageStore {
-
-  private static final Logger LOG = LoggerFactory.getLogger(ImageStore.class);
   private static final String DB_NAME = "images";
 
-  private List<EventProperty> imageProperties;
-  private CouchDbClient couchDbClient;
+  private final List<EventProperty> imageProperties;
+  private final CouchDbClient couchDbClient;
 
   public ImageStore(DataLakeMeasure measure,
                     Environment environment) {
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/ImageStoreUtils.java
 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/ImageStoreUtils.java
similarity index 91%
rename from 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/ImageStoreUtils.java
rename to 
streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/ImageStoreUtils.java
index 6504303a5a..59d4719a7d 100644
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/ImageStoreUtils.java
+++ 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/ImageStoreUtils.java
@@ -16,7 +16,7 @@
  *
  */
 
-package org.apache.streampipes.dataexplorer.influx.migrate;
+package org.apache.streampipes.dataexplorer;
 
 import org.apache.streampipes.model.datalake.DataLakeMeasure;
 import org.apache.streampipes.model.schema.EventProperty;
@@ -30,7 +30,8 @@ public class ImageStoreUtils {
   public static List<EventProperty> getImageProperties(DataLakeMeasure 
measure) {
     return measure.getEventSchema().getEventProperties().stream()
         .filter(eventProperty -> eventProperty.getDomainProperties() != null
-            && eventProperty.getDomainProperties().size() > 0
+            && !eventProperty.getDomainProperties()
+                             .isEmpty()
             && 
eventProperty.getDomainProperties().get(0).toString().equals(SPSensor.IMAGE))
         .collect(Collectors.toList());
   }
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/QueryResultProvider.java
 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/QueryResultProvider.java
similarity index 63%
rename from 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/QueryResultProvider.java
rename to 
streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/QueryResultProvider.java
index 85a23b188f..eeeadec24f 100644
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/QueryResultProvider.java
+++ 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/QueryResultProvider.java
@@ -16,44 +16,56 @@
  *
  */
 
-package org.apache.streampipes.dataexplorer.influx.migrate;
+package org.apache.streampipes.dataexplorer;
 
-import 
org.apache.streampipes.dataexplorer.influx.DataExplorerInfluxQueryExecutor;
+import org.apache.streampipes.dataexplorer.api.IDataExplorerQueryManagement;
 import 
org.apache.streampipes.dataexplorer.param.ProvidedRestQueryParamConverter;
 import org.apache.streampipes.dataexplorer.param.SelectQueryParams;
+import org.apache.streampipes.dataexplorer.query.DataExplorerQueryExecutor;
+import org.apache.streampipes.model.datalake.SpQueryResult;
 import org.apache.streampipes.model.datalake.param.ProvidedRestQueryParams;
 import org.apache.streampipes.model.datalake.param.SupportedRestQueryParams;
-import org.apache.streampipes.model.datalake.SpQueryResult;
+
+import java.util.Optional;
 
 
 public class QueryResultProvider {
 
   public static final String FOR_ID_KEY = "forId";
   protected final boolean ignoreMissingData;
+  protected final IDataExplorerQueryManagement dataExplorerQueryManagement;
+  protected final DataExplorerQueryExecutor<?, ?> queryExecutor;
   protected ProvidedRestQueryParams queryParams;
 
   public QueryResultProvider(ProvidedRestQueryParams queryParams,
+                             IDataExplorerQueryManagement 
dataExplorerQueryManagement,
+                             DataExplorerQueryExecutor<?, ?> queryExecutor,
                              boolean ignoreMissingData) {
     this.queryParams = queryParams;
     this.ignoreMissingData = ignoreMissingData;
+    this.dataExplorerQueryManagement = dataExplorerQueryManagement;
+    this.queryExecutor = queryExecutor;
   }
 
   public SpQueryResult getData() {
     if (queryParams.has(SupportedRestQueryParams.QP_AUTO_AGGREGATE)) {
-      queryParams = new 
AutoAggregationHandler(queryParams).makeAutoAggregationQueryParams();
+      queryParams = new AutoAggregationHandler(queryParams,
+                                               
dataExplorerQueryManagement).makeAutoAggregationQueryParams();
     }
     SelectQueryParams qp = 
ProvidedRestQueryParamConverter.getSelectQueryParams(queryParams);
 
     if 
(queryParams.getProvidedParams().containsKey(SupportedRestQueryParams.QP_MAXIMUM_AMOUNT_OF_EVENTS))
 {
-      int maximumAmountOfEvents = 
Integer.parseInt(queryParams.getProvidedParams().get(SupportedRestQueryParams.QP_MAXIMUM_AMOUNT_OF_EVENTS));
-      return new 
DataExplorerInfluxQueryExecutor(maximumAmountOfEvents).executeQuery(qp, 
ignoreMissingData);
+      int maximumAmountOfEvents = 
Integer.parseInt(queryParams.getProvidedParams()
+                                                              
.get(SupportedRestQueryParams.QP_MAXIMUM_AMOUNT_OF_EVENTS)
+      );
+      return queryExecutor.executeQuery(qp, maximumAmountOfEvents, 
Optional.empty(), ignoreMissingData);
     }
 
     if (queryParams.getProvidedParams().containsKey(FOR_ID_KEY)) {
       String forWidgetId = queryParams.getProvidedParams().get(FOR_ID_KEY);
-      return new DataExplorerInfluxQueryExecutor(forWidgetId).executeQuery(qp, 
ignoreMissingData);
+      return queryExecutor.executeQuery(qp, -1, Optional.of(forWidgetId), 
ignoreMissingData);
     } else {
-      return new DataExplorerInfluxQueryExecutor().executeQuery(qp, 
ignoreMissingData);
+      return queryExecutor.executeQuery(qp, -1, Optional.empty(), 
ignoreMissingData);
     }
   }
 }
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/StreamedQueryResultProvider.java
 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/StreamedQueryResultProvider.java
similarity index 90%
rename from 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/StreamedQueryResultProvider.java
rename to 
streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/StreamedQueryResultProvider.java
index 52ed5d7881..18e288ec98 100644
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/StreamedQueryResultProvider.java
+++ 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/StreamedQueryResultProvider.java
@@ -16,15 +16,17 @@
  *
  */
 
-package org.apache.streampipes.dataexplorer.influx.migrate;
+package org.apache.streampipes.dataexplorer;
 
-import org.apache.streampipes.model.datalake.param.ProvidedRestQueryParams;
-import org.apache.streampipes.model.datalake.param.SupportedRestQueryParams;
+import org.apache.streampipes.dataexplorer.api.IDataExplorerQueryManagement;
 import org.apache.streampipes.dataexplorer.export.ConfiguredOutputWriter;
 import org.apache.streampipes.dataexplorer.export.OutputFormat;
+import org.apache.streampipes.dataexplorer.query.DataExplorerQueryExecutor;
 import org.apache.streampipes.dataexplorer.utils.DataExplorerUtils;
 import org.apache.streampipes.model.datalake.DataLakeMeasure;
 import org.apache.streampipes.model.datalake.SpQueryResult;
+import org.apache.streampipes.model.datalake.param.ProvidedRestQueryParams;
+import org.apache.streampipes.model.datalake.param.SupportedRestQueryParams;
 
 import java.io.IOException;
 import java.io.OutputStream;
@@ -40,8 +42,10 @@ public class StreamedQueryResultProvider extends 
QueryResultProvider {
 
   public StreamedQueryResultProvider(ProvidedRestQueryParams params,
                                      OutputFormat format,
+                                     IDataExplorerQueryManagement 
dataExplorerQueryManagement,
+                                     DataExplorerQueryExecutor<?, ?> 
queryExecutor,
                                      boolean ignoreMissingValues) {
-    super(params, ignoreMissingValues);
+    super(params, dataExplorerQueryManagement, queryExecutor, 
ignoreMissingValues);
     this.format = format;
   }
 
diff --git 
a/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/TimeSeriesStorage.java
 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/TimeSeriesStorage.java
new file mode 100644
index 0000000000..ceaa6b9dc0
--- /dev/null
+++ 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/TimeSeriesStorage.java
@@ -0,0 +1,145 @@
+/*
+ * 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.dataexplorer;
+
+import org.apache.streampipes.commons.exceptions.SpRuntimeException;
+import org.apache.streampipes.dataexplorer.api.ITimeSeriesStorage;
+import org.apache.streampipes.model.datalake.DataLakeMeasure;
+import org.apache.streampipes.model.runtime.Event;
+import org.apache.streampipes.model.schema.EventProperty;
+import org.apache.streampipes.model.schema.EventPropertyPrimitive;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Collectors;
+
+public abstract class TimeSeriesStorage implements ITimeSeriesStorage {
+
+  private static final Logger LOG = 
LoggerFactory.getLogger(TimeSeriesStorage.class);
+
+  protected final DataLakeMeasure measure;
+  protected final List<EventProperty> allEventProperties;
+  protected final Map<String, String> sanitizedRuntimeNames = new HashMap<>();
+
+  public TimeSeriesStorage(DataLakeMeasure measure) {
+    this.measure = measure;
+    storeSanitizedRuntimeNames();
+    allEventProperties = getAllEventPropertiesExceptTimestamp();
+  }
+
+  @Override
+  public void onEvent(Event event) throws SpRuntimeException {
+    validateInputEventAndLogMissingFields(event);
+    sanitizeRuntimeNamesInEvent(event);
+    writeToTimeSeriesStorage(event);
+  }
+
+  private void validateInputEventAndLogMissingFields(Event event) {
+    checkEventIsNotNull(event);
+
+    logMissingFields(event);
+
+    logNullFields(event);
+  }
+
+  private void checkEventIsNotNull(Event event) {
+    if (event == null) {
+      throw new SpRuntimeException("Input event is null");
+    }
+  }
+
+  /**
+   * Logs all fields which are present in the schema, but not in the provided 
event
+   */
+  private void logMissingFields(Event event) {
+    var missingFields = getMissingProperties(allEventProperties, event);
+    if (!missingFields.isEmpty()) {
+      LOG.debug(
+          "Ignored {} fields which were present in the schema, but not in the 
provided event: {}",
+          missingFields.size(),
+          String.join(", ", missingFields)
+      );
+    }
+  }
+
+  /**
+   * Returns a list of the runtime names that are missing within the event
+   */
+  private List<String> getMissingProperties(
+      List<EventProperty> allEventProperties,
+      Event event
+  ) {
+    return allEventProperties.stream()
+                             .map(EventProperty::getRuntimeName)
+                             .filter(runtimeName -> 
event.getOptionalFieldByRuntimeName(runtimeName)
+                                                         .isEmpty())
+                             .toList();
+  }
+
+  /**
+   * Logs all fields that contain null values
+   */
+  private void logNullFields(Event event) {
+    List<String> nullFields = allEventProperties
+        .stream()
+        .filter(EventPropertyPrimitive.class::isInstance)
+        .filter(ep -> {
+          var runtimeName = ep.getRuntimeName();
+          var field = event.getOptionalFieldByRuntimeName(runtimeName);
+
+          return field.isPresent() && field.get()
+                                           .getAsPrimitive()
+                                           .getRawValue() == null;
+        })
+        .map(EventProperty::getRuntimeName)
+        .collect(Collectors.toList());
+
+    if (!nullFields.isEmpty()) {
+      LOG.warn("Ignored {} fields which had a value 'null': {}", 
nullFields.size(), String.join(", ", nullFields));
+    }
+  }
+
+  /**
+   * Returns all measurements properties except the timestamp field
+   */
+  private List<EventProperty> getAllEventPropertiesExceptTimestamp() {
+    return measure.getEventSchema()
+                  .getEventProperties()
+                  .stream()
+                  .filter(ep -> !measure.getTimestampField()
+                                        .endsWith(ep.getRuntimeName()))
+                  .toList();
+  }
+
+  /**
+   * store sanitized target property runtime names in variable 
`sanitizedRuntimeNames`
+   */
+  protected abstract void storeSanitizedRuntimeNames();
+
+  /**
+   * Iterates over all properties of the event and renames the key if it is a 
reserved keywords in InfluxDB
+   */
+  protected abstract void sanitizeRuntimeNamesInEvent(Event event);
+
+  protected abstract void writeToTimeSeriesStorage(Event event) throws 
SpRuntimeException;
+}
diff --git 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/TimeSeriesStore.java
 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/TimeSeriesStore.java
similarity index 72%
rename from 
streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/TimeSeriesStore.java
rename to 
streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/TimeSeriesStore.java
index ef8af9a24d..8a0a0fcc66 100644
--- 
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/migrate/TimeSeriesStore.java
+++ 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/TimeSeriesStore.java
@@ -16,13 +16,11 @@
  *
  */
 
-package org.apache.streampipes.dataexplorer.influx.migrate;
+package org.apache.streampipes.dataexplorer;
 
-import org.apache.streampipes.client.api.IStreamPipesClient;
 import org.apache.streampipes.commons.environment.Environment;
 import org.apache.streampipes.commons.exceptions.SpRuntimeException;
-import org.apache.streampipes.dataexplorer.influx.client.InfluxClientProvider;
-import org.apache.streampipes.dataexplorer.influx.InfluxStore;
+import org.apache.streampipes.dataexplorer.api.ITimeSeriesStorage;
 import org.apache.streampipes.model.datalake.DataLakeMeasure;
 import org.apache.streampipes.model.runtime.Event;
 
@@ -34,23 +32,21 @@ import java.io.IOException;
 public class TimeSeriesStore {
 
   private static final Logger LOG = 
LoggerFactory.getLogger(TimeSeriesStore.class);
-  private final InfluxStore influxStore;
+  private final ITimeSeriesStorage timeSeriesStorage;
   private ImageStore imageStore;
 
 
-  public TimeSeriesStore(Environment environment,
-                         IStreamPipesClient client,
-                         DataLakeMeasure measure,
-                         boolean enableImageStore) {
-
-    DataExplorerUtils.sanitizeAndRegisterAtDataLake(client, measure);
+  public TimeSeriesStore(
+      ITimeSeriesStorage timeSeriesStorage,
+      DataLakeMeasure measure,
+      Environment environment,
+      boolean enableImageStore
+  ) {
 
     if (enableImageStore) {
       this.imageStore = new ImageStore(measure, environment);
     }
-
-    this.influxStore = new InfluxStore(measure, environment, new 
InfluxClientProvider());
-
+    this.timeSeriesStorage = timeSeriesStorage;
   }
 
   public boolean onEvent(Event event) throws SpRuntimeException {
@@ -60,7 +56,7 @@ public class TimeSeriesStore {
     }
 
     // Store event in time series database
-    this.influxStore.onEvent(event);
+    this.timeSeriesStorage.onEvent(event);
 
     return true;
   }
@@ -75,6 +71,6 @@ public class TimeSeriesStore {
       }
     }
 
-    this.influxStore.close();
+    this.timeSeriesStorage.close();
   }
 }
diff --git 
a/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/query/DataExplorerQueryExecutor.java
 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/query/DataExplorerQueryExecutor.java
index d56856257f..347070a804 100644
--- 
a/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/query/DataExplorerQueryExecutor.java
+++ 
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/query/DataExplorerQueryExecutor.java
@@ -29,54 +29,39 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import java.util.Map;
+import java.util.Optional;
 
 public abstract class DataExplorerQueryExecutor<X, W> {
 
   private static final Logger LOG = 
LoggerFactory.getLogger(DataExplorerQueryExecutor.class);
-  protected int maximumAmountOfEvents;
-
-  protected boolean appendId = false;
-  protected String forId;
-
-  public DataExplorerQueryExecutor() {
-    this.maximumAmountOfEvents = -1;
-  }
-
-  public DataExplorerQueryExecutor(String forId) {
-    this();
-    this.appendId = true;
-    this.forId = forId;
-  }
-
-  public DataExplorerQueryExecutor(int maximumAmountOfEvents) {
-    this();
-    this.maximumAmountOfEvents = maximumAmountOfEvents;
-  }
 
   /**
    * Execute the data explorer query and return the result or a warning message
    * in case the maximum amount of events to return is defined
    */
   public SpQueryResult executeQuery(SelectQueryParams params,
+                                    int maximumAmountOfEvents,
+                                    Optional<String> forIdOpt,
                                     boolean ignoreMissingValues) throws 
RuntimeException {
     X query = makeSelectQuery(params);
-    var result = executeQuery(query, ignoreMissingValues);
-    if (this.maximumAmountOfEvents != -1) {
-      return validateAndReturnQueryResult(result, params.getLimit());
+    var result = executeQuery(query, forIdOpt, ignoreMissingValues);
+    if (maximumAmountOfEvents != -1) {
+      return validateAndReturnQueryResult(result, params.getLimit(), 
maximumAmountOfEvents);
     } else {
       return result;
     }
   }
 
   private SpQueryResult validateAndReturnQueryResult(SpQueryResult queryResult,
-                                                     int limit) {
+                                                     int limit,
+                                                     int 
maximumAmountOfEvents) {
     var amountOfResults = queryResult.getAllDataSeries()
         .stream()
         .mapToInt(DataSeries::getTotal)
         .sum();
 
     var amountOfQueryResults = limit == Integer.MIN_VALUE ? amountOfResults : 
Math.min(amountOfResults, limit);
-    if (amountOfQueryResults > this.maximumAmountOfEvents) {
+    if (amountOfQueryResults > maximumAmountOfEvents) {
       return makeTooMuchDataResult(amountOfQueryResults);
     } else {
       return queryResult;
@@ -91,10 +76,11 @@ public abstract class DataExplorerQueryExecutor<X, W> {
   }
 
   public SpQueryResult executeQuery(DeleteQueryParams params) {
-    return executeQuery(makeDeleteQuery(params), true);
+    return executeQuery(makeDeleteQuery(params), Optional.empty(), true);
   }
 
   public SpQueryResult executeQuery(X query,
+                                    Optional<String> forIdOpt,
                                     boolean ignoreMissingValues) {
     if (LOG.isDebugEnabled()) {
       LOG.debug("Data Lake Query {}", asQueryString(query));
@@ -105,10 +91,11 @@ public abstract class DataExplorerQueryExecutor<X, W> {
       LOG.debug("Data Lake Query Result: {}", result.toString());
     }
 
-    return postQuery(result, ignoreMissingValues);
+    return postQuery(result, forIdOpt, ignoreMissingValues);
   }
 
   protected abstract SpQueryResult postQuery(W queryResult,
+                                             Optional<String> forIdOpt,
                                              boolean ignoreMissingValues);
 
   public abstract W executeQuery(X query);
diff --git a/streampipes-extensions/streampipes-sinks-internal-jvm/pom.xml 
b/streampipes-extensions/streampipes-sinks-internal-jvm/pom.xml
index d9b9f7d216..239a09c3ac 100644
--- a/streampipes-extensions/streampipes-sinks-internal-jvm/pom.xml
+++ b/streampipes-extensions/streampipes-sinks-internal-jvm/pom.xml
@@ -41,19 +41,25 @@
         <!-- StreamPipes dependencies -->
         <dependency>
             <groupId>org.apache.streampipes</groupId>
-            <artifactId>streampipes-extensions-management</artifactId>
+            <artifactId>streampipes-data-explorer-influx</artifactId>
             <version>0.95.0-SNAPSHOT</version>
         </dependency>
         <dependency>
             <groupId>org.apache.streampipes</groupId>
-            <artifactId>streampipes-wrapper-standalone</artifactId>
+            <artifactId>streampipes-data-explorer-management</artifactId>
             <version>0.95.0-SNAPSHOT</version>
         </dependency>
         <dependency>
             <groupId>org.apache.streampipes</groupId>
-            <artifactId>streampipes-data-explorer-influx</artifactId>
+            <artifactId>streampipes-extensions-management</artifactId>
             <version>0.95.0-SNAPSHOT</version>
         </dependency>
+        <dependency>
+            <groupId>org.apache.streampipes</groupId>
+            <artifactId>streampipes-wrapper-standalone</artifactId>
+            <version>0.95.0-SNAPSHOT</version>
+        </dependency>
+
 
         <!-- External dependencies -->
         <!-- TODO remove after refactoring -->
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 2a8250cb82..ef88bbdcdc 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
@@ -20,11 +20,13 @@ package org.apache.streampipes.sinks.internal.jvm.datalake;
 
 import org.apache.streampipes.commons.environment.Environments;
 import org.apache.streampipes.commons.exceptions.SpRuntimeException;
-import org.apache.streampipes.dataexplorer.influx.migrate.TimeSeriesStore;
+import org.apache.streampipes.dataexplorer.TimeSeriesStore;
+import org.apache.streampipes.dataexplorer.management.DataExplorerDispatcher;
 import 
org.apache.streampipes.extensions.api.pe.context.EventSinkRuntimeContext;
 import org.apache.streampipes.model.DataSinkType;
 import org.apache.streampipes.model.datalake.DataLakeMeasure;
 import 
org.apache.streampipes.model.datalake.DataLakeMeasureSchemaUpdateStrategy;
+import org.apache.streampipes.model.extensions.ExtensionAssetType;
 import org.apache.streampipes.model.graph.DataSinkDescription;
 import org.apache.streampipes.model.runtime.Event;
 import org.apache.streampipes.model.schema.PropertyScope;
@@ -34,7 +36,6 @@ import org.apache.streampipes.sdk.helpers.EpRequirements;
 import org.apache.streampipes.sdk.helpers.Labels;
 import org.apache.streampipes.sdk.helpers.Locales;
 import org.apache.streampipes.sdk.helpers.Options;
-import org.apache.streampipes.sdk.utils.Assets;
 import org.apache.streampipes.wrapper.params.compat.SinkParams;
 import org.apache.streampipes.wrapper.standalone.StreamPipesDataSink;
 
@@ -54,25 +55,25 @@ public class DataLakeSink extends StreamPipesDataSink {
   @Override
   public DataSinkDescription declareModel() {
     return DataSinkBuilder
-        .create("org.apache.streampipes.sinks.internal.jvm.datalake", 1)
-        .withLocales(Locales.EN)
-        .withAssets(Assets.DOCUMENTATION, Assets.ICON)
-        .category(DataSinkType.INTERNAL)
-        .requiredStream(StreamRequirementsBuilder
-                            .create()
-                            .requiredPropertyWithUnaryMapping(
-                                EpRequirements.timestampReq(),
-                                Labels.withId(TIMESTAMP_MAPPING_KEY),
-                                PropertyScope.NONE
-                            )
-                            .build())
-        .requiredTextParameter(Labels.withId(DATABASE_MEASUREMENT_KEY))
-        .requiredSingleValueSelection(
-            Labels.withId(SCHEMA_UPDATE_KEY),
-            Options.from(SCHEMA_UPDATE_OPTION, EXTEND_EXISTING_SCHEMA_OPTION)
-        )
-
-        .build();
+      .create("org.apache.streampipes.sinks.internal.jvm.datalake", 1)
+      .withLocales(Locales.EN)
+      .withAssets(ExtensionAssetType.DOCUMENTATION, ExtensionAssetType.ICON)
+      .category(DataSinkType.INTERNAL)
+      .requiredStream(StreamRequirementsBuilder
+                        .create()
+                        .requiredPropertyWithUnaryMapping(
+                          EpRequirements.timestampReq(),
+                          Labels.withId(TIMESTAMP_MAPPING_KEY),
+                          PropertyScope.NONE
+                        )
+                        .build())
+      .requiredTextParameter(Labels.withId(DATABASE_MEASUREMENT_KEY))
+      .requiredSingleValueSelection(
+        Labels.withId(SCHEMA_UPDATE_KEY),
+        Options.from(SCHEMA_UPDATE_OPTION, EXTEND_EXISTING_SCHEMA_OPTION)
+      )
+
+      .build();
   }
 
   @Override
@@ -95,10 +96,14 @@ public class DataLakeSink extends StreamPipesDataSink {
       
measure.setSchemaUpdateStrategy(DataLakeMeasureSchemaUpdateStrategy.UPDATE_SCHEMA);
     }
 
+    measure = new DataExplorerDispatcher().getDataExplorerManager()
+        .getMeasurementSanitizer(runtimeContext.getStreamPipesClient(), 
measure)
+        .sanitizeAndRegister();
+
     this.timeSeriesStore = new TimeSeriesStore(
-        Environments.getEnvironment(),
-        runtimeContext.getStreamPipesClient(),
+        new 
DataExplorerDispatcher().getDataExplorerManager().getTimeseriesStorage(measure),
         measure,
+        Environments.getEnvironment(),
         true
     );
 
diff --git a/streampipes-platform-services/pom.xml 
b/streampipes-platform-services/pom.xml
index 13f92da128..ffe439094a 100644
--- a/streampipes-platform-services/pom.xml
+++ b/streampipes-platform-services/pom.xml
@@ -38,6 +38,11 @@
             <artifactId>streampipes-data-explorer-influx</artifactId>
             <version>0.95.0-SNAPSHOT</version>
         </dependency>
+        <dependency>
+            <groupId>org.apache.streampipes</groupId>
+            <artifactId>streampipes-data-explorer-management</artifactId>
+            <version>0.95.0-SNAPSHOT</version>
+        </dependency>
         <dependency>
             <groupId>org.apache.streampipes</groupId>
             <artifactId>streampipes-model</artifactId>
diff --git 
a/streampipes-platform-services/src/main/java/org/apache/streampipes/ps/DataLakeMeasureResourceV4.java
 
b/streampipes-platform-services/src/main/java/org/apache/streampipes/ps/DataLakeMeasureResourceV4.java
index 64c009e7a8..266bca7b27 100644
--- 
a/streampipes-platform-services/src/main/java/org/apache/streampipes/ps/DataLakeMeasureResourceV4.java
+++ 
b/streampipes-platform-services/src/main/java/org/apache/streampipes/ps/DataLakeMeasureResourceV4.java
@@ -18,12 +18,10 @@
 
 package org.apache.streampipes.ps;
 
-import org.apache.streampipes.dataexplorer.DataExplorerSchemaManagement;
 import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement;
-import 
org.apache.streampipes.dataexplorer.influx.migrate.DataLakeMeasurementCount;
+import org.apache.streampipes.dataexplorer.management.DataExplorerDispatcher;
 import org.apache.streampipes.model.datalake.DataLakeMeasure;
 import 
org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource;
-import org.apache.streampipes.storage.management.StorageDispatcher;
 
 import org.springframework.http.MediaType;
 import org.springframework.http.ResponseEntity;
@@ -48,10 +46,8 @@ public class DataLakeMeasureResourceV4 extends 
AbstractAuthGuardedRestResource {
   private final IDataExplorerSchemaManagement dataLakeMeasureManagement;
 
   public DataLakeMeasureResourceV4() {
-    var dataLakeStorage = StorageDispatcher.INSTANCE
-        .getNoSqlStore()
-        .getDataLakeStorage();
-    this.dataLakeMeasureManagement = new 
DataExplorerSchemaManagement(dataLakeStorage);
+    this.dataLakeMeasureManagement = new 
DataExplorerDispatcher().getDataExplorerManager()
+        .getSchemaManagement();
   }
 
   @PostMapping(
@@ -67,7 +63,12 @@ public class DataLakeMeasureResourceV4 extends 
AbstractAuthGuardedRestResource {
   public ResponseEntity<Map<String, Integer>> getDataLakeInfos(
       @RequestParam(value = "filter", required = false) List<String> 
measurementNames) {
     var allMeasurements = this.dataLakeMeasureManagement.getAllMeasurements();
-    return ok(new DataLakeMeasurementCount(allMeasurements, 
measurementNames).countMeasurementSizes());
+    return ok(new DataExplorerDispatcher().getDataExplorerManager()
+        .getMeasurementCounter(
+            allMeasurements,
+            measurementNames
+        )
+        .countMeasurementSizes());
   }
 
   @GetMapping(path = "{id}", produces = MediaType.APPLICATION_JSON_VALUE)
diff --git 
a/streampipes-platform-services/src/main/java/org/apache/streampipes/ps/DataLakeResourceV4.java
 
b/streampipes-platform-services/src/main/java/org/apache/streampipes/ps/DataLakeResourceV4.java
index 4716a9ec06..2bc5cf4a79 100644
--- 
a/streampipes-platform-services/src/main/java/org/apache/streampipes/ps/DataLakeResourceV4.java
+++ 
b/streampipes-platform-services/src/main/java/org/apache/streampipes/ps/DataLakeResourceV4.java
@@ -18,9 +18,10 @@
 
 package org.apache.streampipes.ps;
 
-import org.apache.streampipes.dataexplorer.DataExplorerSchemaManagement;
+import org.apache.streampipes.dataexplorer.api.IDataExplorerQueryManagement;
+import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement;
 import org.apache.streampipes.dataexplorer.export.OutputFormat;
-import 
org.apache.streampipes.dataexplorer.influx.migrate.DataExplorerQueryManagement;
+import org.apache.streampipes.dataexplorer.management.DataExplorerDispatcher;
 import org.apache.streampipes.model.datalake.DataLakeMeasure;
 import org.apache.streampipes.model.datalake.DataSeries;
 import org.apache.streampipes.model.datalake.SpQueryResult;
@@ -29,7 +30,6 @@ import org.apache.streampipes.model.message.Notifications;
 import org.apache.streampipes.model.monitoring.SpLogMessage;
 import org.apache.streampipes.rest.core.base.impl.AbstractRestResource;
 import org.apache.streampipes.rest.shared.exception.SpMessageException;
-import org.apache.streampipes.storage.management.StorageDispatcher;
 
 import io.swagger.v3.oas.annotations.Operation;
 import io.swagger.v3.oas.annotations.Parameter;
@@ -80,23 +80,23 @@ import static 
org.apache.streampipes.model.datalake.param.SupportedRestQueryPara
 @RequestMapping("/api/v4/datalake")
 public class DataLakeResourceV4 extends AbstractRestResource {
 
-  private final DataExplorerQueryManagement dataLakeManagement;
-  private final DataExplorerSchemaManagement dataExplorerSchemaManagement;
+  private final IDataExplorerQueryManagement dataExplorerQueryManagement;
+  private final IDataExplorerSchemaManagement dataExplorerSchemaManagement;
 
   public DataLakeResourceV4() {
-    var dataLakeStorage = StorageDispatcher.INSTANCE
-        .getNoSqlStore()
-        .getDataLakeStorage();
-    this.dataExplorerSchemaManagement = new 
DataExplorerSchemaManagement(dataLakeStorage);
-    this.dataLakeManagement = new 
DataExplorerQueryManagement(dataExplorerSchemaManagement);
+    this.dataExplorerSchemaManagement = new DataExplorerDispatcher()
+        .getDataExplorerManager()
+        .getSchemaManagement();
+    this.dataExplorerQueryManagement = new DataExplorerDispatcher()
+        .getDataExplorerManager()
+        .getQueryManagement(this.dataExplorerSchemaManagement);
   }
 
-  public DataLakeResourceV4(DataExplorerQueryManagement dataLakeManagement) {
-    var dataLakeStorage = StorageDispatcher.INSTANCE
-        .getNoSqlStore()
-        .getDataLakeStorage();
-    this.dataLakeManagement = dataLakeManagement;
-    this.dataExplorerSchemaManagement = new 
DataExplorerSchemaManagement(dataLakeStorage);
+  public DataLakeResourceV4(IDataExplorerQueryManagement 
dataExplorerQueryManagement) {
+    this.dataExplorerQueryManagement = dataExplorerQueryManagement;
+    this.dataExplorerSchemaManagement = new DataExplorerDispatcher()
+        .getDataExplorerManager()
+        .getSchemaManagement();
   }
 
   @DeleteMapping(path = "/measurements/{measurementID}")
@@ -112,7 +112,7 @@ public class DataLakeResourceV4 extends 
AbstractRestResource {
       , @Parameter(in = ParameterIn.QUERY, description = "end date for slicing 
operation")
       @RequestParam(value = "endDate", required = false) Long endDate) {
 
-    SpQueryResult result = this.dataLakeManagement.deleteData(measurementID, 
startDate, endDate);
+    SpQueryResult result = 
this.dataExplorerQueryManagement.deleteData(measurementID, startDate, endDate);
     return ok();
   }
 
@@ -132,7 +132,7 @@ public class DataLakeResourceV4 extends 
AbstractRestResource {
       @Parameter(in = ParameterIn.PATH, description = "the id of the 
measurement series", required = true)
       @PathVariable("measurementID") String measurementID) {
 
-    boolean isSuccessDataLake = 
this.dataLakeManagement.deleteData(measurementID);
+    boolean isSuccessDataLake = 
this.dataExplorerQueryManagement.deleteData(measurementID);
 
     if (isSuccessDataLake) {
       boolean isSuccessEventProperty = 
this.dataExplorerSchemaManagement.deleteMeasurementByName(measurementID);
@@ -165,7 +165,7 @@ public class DataLakeResourceV4 extends 
AbstractRestResource {
   @GetMapping(path = "/measurements/{measurementId}/tags", produces = 
MediaType.APPLICATION_JSON_VALUE)
   public ResponseEntity<Map<String, Object>> 
getTagValues(@PathVariable("measurementId") String measurementId,
                                                           
@RequestParam("fields") String fields) {
-    Map<String, Object> tagValues = 
dataLakeManagement.getTagValues(measurementId, fields);
+    Map<String, Object> tagValues = 
dataExplorerQueryManagement.getTagValues(measurementId, fields);
     return ok(tagValues);
   }
 
@@ -235,7 +235,7 @@ public class DataLakeResourceV4 extends 
AbstractRestResource {
       ProvidedRestQueryParams sanitizedParams = populate(measurementID, 
queryParams);
       try {
         SpQueryResult result =
-            this.dataLakeManagement.getData(sanitizedParams, 
isIgnoreMissingValues(missingValueBehaviour));
+            this.dataExplorerQueryManagement.getData(sanitizedParams, 
isIgnoreMissingValues(missingValueBehaviour));
         return ok(result);
       } catch (RuntimeException e) {
         return badRequest(SpLogMessage.from(e));
@@ -251,7 +251,7 @@ public class DataLakeResourceV4 extends 
AbstractRestResource {
     var results = queryParams
         .stream()
         .map(qp -> new ProvidedRestQueryParams(qp.get("measureName"), qp))
-        .map(params -> this.dataLakeManagement.getData(params, true))
+        .map(params -> this.dataExplorerQueryManagement.getData(params, true))
         .collect(Collectors.toList());
 
     return ok(results);
@@ -322,7 +322,7 @@ public class DataLakeResourceV4 extends 
AbstractRestResource {
       }
 
       OutputFormat outputFormat = format.equals("csv") ? OutputFormat.CSV : 
OutputFormat.JSON;
-      StreamingResponseBody streamingOutput = output -> 
dataLakeManagement.getDataAsStream(
+      StreamingResponseBody streamingOutput = output -> 
dataExplorerQueryManagement.getDataAsStream(
           sanitizedParams,
           outputFormat,
           isIgnoreMissingValues(missingValueBehaviour),
@@ -343,7 +343,7 @@ public class DataLakeResourceV4 extends 
AbstractRestResource {
       responses = {
           @ApiResponse(responseCode = "200", description = "All measurement 
series successfully removed")})
   public ResponseEntity<?> removeAll() {
-    boolean isSuccess = this.dataLakeManagement.deleteAllData();
+    boolean isSuccess = this.dataExplorerQueryManagement.deleteAllData();
     return ResponseEntity.ok(isSuccess);
   }
 
@@ -360,13 +360,6 @@ public class DataLakeResourceV4 extends 
AbstractRestResource {
 
   // Checks if the parameter for missing value behaviour is set
   private boolean isIgnoreMissingValues(String missingValueBehaviour) {
-    boolean ignoreMissingValues;
-    if ("ignore".equals(missingValueBehaviour)) {
-      ignoreMissingValues = true;
-    } else {
-      ignoreMissingValues = false;
-    }
-    return ignoreMissingValues;
+    return "ignore".equals(missingValueBehaviour);
   }
-
 }
diff --git a/streampipes-rest/pom.xml b/streampipes-rest/pom.xml
index 48b45acddc..5a007f5d3d 100644
--- a/streampipes-rest/pom.xml
+++ b/streampipes-rest/pom.xml
@@ -45,6 +45,11 @@
             <artifactId>streampipes-data-explorer-influx</artifactId>
             <version>0.95.0-SNAPSHOT</version>
         </dependency>
+        <dependency>
+            <groupId>org.apache.streampipes</groupId>
+            <artifactId>streampipes-data-explorer-management</artifactId>
+            <version>0.95.0-SNAPSHOT</version>
+        </dependency>
         <dependency>
             <groupId>org.apache.streampipes</groupId>
             <artifactId>streampipes-data-export</artifactId>
diff --git 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/ResetManagement.java
 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/ResetManagement.java
index 3c68acf0b9..3d57aa8ce7 100644
--- 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/ResetManagement.java
+++ 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/ResetManagement.java
@@ -22,8 +22,7 @@ import 
org.apache.streampipes.commons.exceptions.SpRuntimeException;
 import org.apache.streampipes.commons.exceptions.connect.AdapterException;
 import org.apache.streampipes.commons.prometheus.adapter.AdapterMetricsManager;
 import 
org.apache.streampipes.connect.management.management.AdapterMasterManagement;
-import org.apache.streampipes.dataexplorer.DataExplorerSchemaManagement;
-import 
org.apache.streampipes.dataexplorer.influx.migrate.DataExplorerQueryManagement;
+import org.apache.streampipes.dataexplorer.management.DataExplorerDispatcher;
 import org.apache.streampipes.manager.file.FileManager;
 import org.apache.streampipes.manager.pipeline.PipelineCacheManager;
 import 
org.apache.streampipes.manager.pipeline.PipelineCanvasMetadataCacheManager;
@@ -137,12 +136,12 @@ public class ResetManagement {
   }
 
   private static void removeAllDataInDataLake() {
-    var dataLakeStorage = StorageDispatcher.INSTANCE
-        .getNoSqlStore()
-        .getDataLakeStorage();
-    var dataLakeMeasureManagement = new 
DataExplorerSchemaManagement(dataLakeStorage);
-    var dataExplorerQueryManagement =
-        new DataExplorerQueryManagement(dataLakeMeasureManagement);
+    var dataLakeMeasureManagement = new DataExplorerDispatcher()
+        .getDataExplorerManager()
+        .getSchemaManagement();
+    var dataExplorerQueryManagement = new DataExplorerDispatcher()
+        .getDataExplorerManager()
+        .getQueryManagement(dataLakeMeasureManagement);
     List<DataLakeMeasure> allMeasurements = 
dataLakeMeasureManagement.getAllMeasurements();
     allMeasurements.forEach(measurement -> {
       boolean isSuccessDataLake = 
dataExplorerQueryManagement.deleteData(measurement.getMeasureName());

Reply via email to