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

bossenti pushed a commit to branch basic-iotdb-support
in repository https://gitbox.apache.org/repos/asf/streampipes.git

commit 4a605c3150bb6c84e606f4f118714931e733b7f3
Author: bossenti <[email protected]>
AuthorDate: Tue May 7 10:15:10 2024 +0200

    feat: provide basic support for IotDB as time series storage
---
 streampipes-data-explorer-management/pom.xml       |   5 +
 .../management/DataExplorerDispatcher.java         |   3 +-
 .../management/SupportedDataExplorerStorages.java  |   1 +
 streampipes-ts-store-iotdb/pom.xml                 |  16 +++
 .../ts/store/iotdb/DataExplorerManagerIotDb.java   |  60 ++++++++
 .../ts/store/iotdb/TimeSeriesStorageIotDb.java     | 158 +++++++++++++++++++++
 .../DataLakeMeasurementSanitizerIotDb.java         |  45 ++++++
 .../iotdb/{ => sanitize}/IotDbNameSanitizer.java   |   2 +-
 .../{ => sanitize}/IotDbReservedKeywords.java      |   2 +-
 .../{ => sanitize}/MeasureNameSanitizerIotDb.java  |   2 +-
 .../DataLakeMeasurementSanitizerIotDbTest.java     |  75 ++++++++++
 .../{ => sanitize}/IotDbNameSanitizerTest.java     |   2 +-
 .../sanitize/MeasureNameSanitizerIotDbTest.java    |   2 +-
 13 files changed, 367 insertions(+), 6 deletions(-)

diff --git a/streampipes-data-explorer-management/pom.xml 
b/streampipes-data-explorer-management/pom.xml
index 9e1d6a9600..3c4922669c 100644
--- a/streampipes-data-explorer-management/pom.xml
+++ b/streampipes-data-explorer-management/pom.xml
@@ -47,6 +47,11 @@
             <artifactId>streampipes-data-explorer-influx</artifactId>
             <version>0.97.0-SNAPSHOT</version>
         </dependency>
+        <dependency>
+            <groupId>org.apache.streampipes</groupId>
+            <artifactId>streampipes-ts-store-iotdb</artifactId>
+            <version>0.97.0-SNAPSHOT</version>
+        </dependency>
     </dependencies>
 
     <properties>
diff --git 
a/streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/DataExplorerDispatcher.java
 
b/streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/DataExplorerDispatcher.java
index 96ba5b70ae..ef9b6105e8 100644
--- 
a/streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/DataExplorerDispatcher.java
+++ 
b/streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/DataExplorerDispatcher.java
@@ -21,6 +21,7 @@ package org.apache.streampipes.dataexplorer.management;
 import org.apache.streampipes.commons.environment.Environments;
 import org.apache.streampipes.dataexplorer.api.IDataExplorerManager;
 import org.apache.streampipes.dataexplorer.influx.DataExplorerManagerInflux;
+import org.apache.streampipes.ts.store.iotdb.DataExplorerManagerIotDb;
 
 public class DataExplorerDispatcher {
 
@@ -31,7 +32,7 @@ public class DataExplorerDispatcher {
     return switch (Environments.getEnvironment()
                                .getTsStorage()
                                .getValueOrDefault()) {
-      case SupportedDataExplorerStorages.INFLUX_DB -> 
DataExplorerManagerInflux.INSTANCE;
+      case SupportedDataExplorerStorages.IOT_DB -> 
DataExplorerManagerIotDb.INSTANCE;
       default -> DataExplorerManagerInflux.INSTANCE;
     };
   }
diff --git 
a/streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/SupportedDataExplorerStorages.java
 
b/streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/SupportedDataExplorerStorages.java
index 70af654be8..1ee48e2acc 100644
--- 
a/streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/SupportedDataExplorerStorages.java
+++ 
b/streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/SupportedDataExplorerStorages.java
@@ -28,4 +28,5 @@ package org.apache.streampipes.dataexplorer.management;
  */
 public class SupportedDataExplorerStorages {
   public static final String INFLUX_DB = "influxdb";
+  public static final String IOT_DB = "iotdb";
 }
diff --git a/streampipes-ts-store-iotdb/pom.xml 
b/streampipes-ts-store-iotdb/pom.xml
index ff294b445d..c1e09c66c4 100644
--- a/streampipes-ts-store-iotdb/pom.xml
+++ b/streampipes-ts-store-iotdb/pom.xml
@@ -44,6 +44,11 @@
             <artifactId>streampipes-commons</artifactId>
             <version>0.97.0-SNAPSHOT</version>
         </dependency>
+        <dependency>
+            <groupId>org.apache.streampipes</groupId>
+            <artifactId>streampipes-data-explorer-api</artifactId>
+            <version>0.97.0-SNAPSHOT</version>
+        </dependency>
         <dependency>
             <groupId>org.apache.streampipes</groupId>
             <artifactId>streampipes-data-explorer</artifactId>
@@ -72,11 +77,22 @@
         </dependency>
 
         <!--Test dependencies -->
+        <dependency>
+            <groupId>org.apache.streampipes</groupId>
+            <artifactId>streampipes-test-utils</artifactId>
+            <version>0.97.0-SNAPSHOT</version>
+            <scope>test</scope>
+        </dependency>
         <dependency>
             <groupId>org.junit.jupiter</groupId>
             <artifactId>junit-jupiter-api</artifactId>
             <scope>test</scope>
         </dependency>
+        <dependency>
+            <groupId>org.mockito</groupId>
+            <artifactId>mockito-core</artifactId>
+            <scope>test</scope>
+        </dependency>
     </dependencies>
 
 </project>
\ No newline at end of file
diff --git 
a/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/DataExplorerManagerIotDb.java
 
b/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/DataExplorerManagerIotDb.java
new file mode 100644
index 0000000000..2af5a59537
--- /dev/null
+++ 
b/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/DataExplorerManagerIotDb.java
@@ -0,0 +1,60 @@
+/*
+ * 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.ts.store.iotdb;
+
+import org.apache.streampipes.client.api.IStreamPipesClient;
+import org.apache.streampipes.dataexplorer.DataExplorerSchemaManagement;
+import org.apache.streampipes.dataexplorer.api.*;
+import org.apache.streampipes.model.datalake.DataLakeMeasure;
+import org.apache.streampipes.storage.management.StorageDispatcher;
+import 
org.apache.streampipes.ts.store.iotdb.sanitize.DataLakeMeasurementSanitizerIotDb;
+
+import java.util.List;
+
+public enum DataExplorerManagerIotDb implements IDataExplorerManager {
+
+  INSTANCE;
+
+  @Override
+  public IDataLakeMeasurementCounter 
getMeasurementCounter(List<DataLakeMeasure> allMeasurements, List<String> 
measurementsToCount) {
+    return null;
+  }
+
+  @Override
+  public IDataExplorerQueryManagement 
getQueryManagement(IDataExplorerSchemaManagement dataExplorerSchemaManagement) {
+    return null;
+  }
+
+  @Override
+  public IDataExplorerSchemaManagement getSchemaManagement() {
+    return new DataExplorerSchemaManagement(StorageDispatcher.INSTANCE
+                                                             .getNoSqlStore()
+                                                             
.getDataLakeStorage());
+  }
+
+  @Override
+  public ITimeSeriesStorage getTimeseriesStorage(DataLakeMeasure measure) {
+    return new TimeSeriesStorageIotDb(measure, new IotDbPropertyConverter(), 
new IotDbSessionProvider());
+  }
+
+  @Override
+  public IDataLakeMeasurementSanitizer 
getMeasurementSanitizer(IStreamPipesClient client, DataLakeMeasure measure) {
+    return new DataLakeMeasurementSanitizerIotDb(client, measure);
+  }
+}
diff --git 
a/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/TimeSeriesStorageIotDb.java
 
b/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/TimeSeriesStorageIotDb.java
new file mode 100644
index 0000000000..04252d99a7
--- /dev/null
+++ 
b/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/TimeSeriesStorageIotDb.java
@@ -0,0 +1,158 @@
+/*
+ * 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.ts.store.iotdb;
+
+import org.apache.iotdb.rpc.IoTDBConnectionException;
+import org.apache.iotdb.rpc.StatementExecutionException;
+import org.apache.iotdb.session.pool.SessionPool;
+import org.apache.streampipes.commons.environment.Environments;
+import org.apache.streampipes.commons.exceptions.SpRuntimeException;
+import org.apache.streampipes.dataexplorer.TimeSeriesStorage;
+import org.apache.streampipes.model.datalake.DataLakeMeasure;
+import org.apache.streampipes.model.runtime.Event;
+import org.apache.streampipes.model.schema.EventPropertyPrimitive;
+import org.apache.streampipes.ts.store.iotdb.sanitize.IotDbNameSanitizer;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.ArrayList;
+import java.util.List;
+
+public class TimeSeriesStorageIotDb extends TimeSeriesStorage {
+
+  private static final Logger LOG = 
LoggerFactory.getLogger(TimeSeriesStorageIotDb.class);
+
+  private final IotDbPropertyConverter propertyConverter;
+  private final SessionPool sessionPool;
+
+  public TimeSeriesStorageIotDb(DataLakeMeasure measure,
+                                IotDbPropertyConverter propertyConverter,
+                                IotDbSessionProvider iotDbSessionProvider) {
+    super(measure);
+    this.propertyConverter = propertyConverter;
+    this.sessionPool = 
iotDbSessionProvider.getSessionPool(Environments.getEnvironment());
+  }
+
+  @Override
+  protected void storeSanitizedRuntimeNames() {
+    measure.getEventSchema()
+        .getEventProperties()
+        .forEach(ep -> sanitizedRuntimeNames.put(
+            ep.getRuntimeName(),
+            new 
IotDbNameSanitizer().renameReservedKeywords(ep.getRuntimeName())
+        ));
+  }
+
+  @Override
+  protected void sanitizeRuntimeNamesInEvent(Event event) {
+    event.getRaw()
+        .keySet()
+        .forEach(runtimeName -> {
+              // timestamp field does not need to be renamed since it is not 
written as measurement
+              // and taken care of separately
+              if (!runtimeName.equals(measure.getTimestampFieldName())) {
+                event.renameFieldByRuntimeName(
+                    runtimeName,
+                    new 
IotDbNameSanitizer().renameReservedKeywords(runtimeName)
+                );
+              }
+            }
+        );
+  }
+
+  @Override
+  protected void writeToTimeSeriesStorage(Event event) throws 
SpRuntimeException {
+    if (event == null) {
+      LOG.warn("Input event is null - skipping event");
+      return;
+    }
+
+    var timestampValue = event.getFieldBySelector(measure.getTimestampField())
+                              .getAsPrimitive()
+                              .getAsLong();
+
+    if (timestampValue == null) {
+      LOG.warn("Timestamp of input event is null - skipping event");
+      return;
+    }
+
+    if (event.getRaw().size() <= 1) {
+      LOG.warn("Event only consists of timestamp and does not include any 
measurement - skipping event");
+      return;
+    }
+
+    var iotDbRecords = extractMeasurementRecords(event);
+
+    if (iotDbRecords.isEmpty()) {
+      return;
+    }
+
+    insertIntoIotDb(timestampValue, iotDbRecords);
+  }
+
+  @Override
+  public void close() throws SpRuntimeException {
+    this.sessionPool.close();
+  }
+
+  /**
+   * Extracts all relevant information for the IotDb from the incoming event.
+   *
+   * @param event The event from which measurement records are extracted.
+   * @return An ArrayList of IotDbMeasurementRecord objects containing the 
extracted measurement records.
+   */
+  private ArrayList<IotDbMeasurementRecord> extractMeasurementRecords(Event 
event){
+    var iotDbRecords = new ArrayList<IotDbMeasurementRecord>();
+    allEventProperties.forEach(ep -> {
+      try {
+        // timestamp is already known and is not part of the measurement thus 
we ignore it here
+        if (!ep.getRuntimeName().equals(measure.getTimestampFieldName())) {
+          if (ep instanceof EventPropertyPrimitive) {
+            
iotDbRecords.add(propertyConverter.convertPrimitiveProperty((EventPropertyPrimitive)
 ep, event.getFieldByRuntimeName(ep.getRuntimeName()).getAsPrimitive(), 
sanitizedRuntimeNames.get(ep.getRuntimeName())));
+          } else {
+            iotDbRecords.add(propertyConverter.convertNonPrimitiveProperty(ep, 
sanitizedRuntimeNames.get(ep.getRuntimeName())));
+          }
+        }
+      } catch (SpRuntimeException e) {
+        LOG.error("Event could not be converted to the IotDB representation: 
{}", e.getMessage());
+      }
+    });
+    return iotDbRecords;
+  }
+
+  /**
+   * Inserts measurement records into an IoTDB database as an aligned time 
series of one device.
+   *
+   * @param timestampValue The timestamp value associated with the measurement 
records.
+   * @param iotDbRecords   An ArrayList of IotDbMeasurementRecord objects 
containing the measurement records to be inserted.
+   */
+  private void insertIntoIotDb(long timestampValue, 
ArrayList<IotDbMeasurementRecord> iotDbRecords) {
+    try {
+      sessionPool.insertAlignedRecordsOfOneDevice(
+          "root.streampipes.%s".formatted(measure.getMeasureName()),
+          List.of(timestampValue),
+          
List.of(iotDbRecords.stream().map(IotDbMeasurementRecord::measurementName).toList()),
+          
List.of(iotDbRecords.stream().map(IotDbMeasurementRecord::dataType).toList()),
+          
List.of(iotDbRecords.stream().map(IotDbMeasurementRecord::value).toList())
+      );
+    } catch (IoTDBConnectionException | StatementExecutionException e) {
+      LOG.error("Failed to write event to IoTDB  - {}", e.getMessage());
+    }
+  }
+}
diff --git 
a/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/sanitize/DataLakeMeasurementSanitizerIotDb.java
 
b/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/sanitize/DataLakeMeasurementSanitizerIotDb.java
new file mode 100644
index 0000000000..8bf25321df
--- /dev/null
+++ 
b/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/sanitize/DataLakeMeasurementSanitizerIotDb.java
@@ -0,0 +1,45 @@
+/*
+ * 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.ts.store.iotdb.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 DataLakeMeasurementSanitizerIotDb extends 
DataLakeMeasurementSanitizer {
+  public DataLakeMeasurementSanitizerIotDb(IStreamPipesClient client, 
DataLakeMeasure measure) {
+    super(client, measure);
+  }
+
+  @Override
+  protected void cleanDataLakeMeasure() throws SpRuntimeException {
+      measure.setMeasureName(new 
MeasureNameSanitizerIotDb().sanitize(measure.getMeasureName()));
+
+      measure.getEventSchema()
+             .getEventProperties()
+             .forEach(eventProperty -> eventProperty.setRuntimeName(
+                 new 
IotDbNameSanitizer().renameReservedKeywords(eventProperty.getRuntimeName()))
+             );
+  }
+
+  protected DataLakeMeasure getMeasure() {
+    return measure;
+  }
+}
diff --git 
a/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/IotDbNameSanitizer.java
 
b/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/sanitize/IotDbNameSanitizer.java
similarity index 96%
rename from 
streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/IotDbNameSanitizer.java
rename to 
streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/sanitize/IotDbNameSanitizer.java
index 8d66978dc9..467875624b 100644
--- 
a/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/IotDbNameSanitizer.java
+++ 
b/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/sanitize/IotDbNameSanitizer.java
@@ -16,7 +16,7 @@
  *
  */
 
-package org.apache.streampipes.ts.store.iotdb;
+package org.apache.streampipes.ts.store.iotdb.sanitize;
 
 public class IotDbNameSanitizer {
 
diff --git 
a/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/IotDbReservedKeywords.java
 
b/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/sanitize/IotDbReservedKeywords.java
similarity index 98%
rename from 
streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/IotDbReservedKeywords.java
rename to 
streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/sanitize/IotDbReservedKeywords.java
index 5813cffce1..ae78632adb 100644
--- 
a/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/IotDbReservedKeywords.java
+++ 
b/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/sanitize/IotDbReservedKeywords.java
@@ -16,7 +16,7 @@
  *
  */
 
-package org.apache.streampipes.ts.store.iotdb;
+package org.apache.streampipes.ts.store.iotdb.sanitize;
 
 import java.util.Arrays;
 import java.util.List;
diff --git 
a/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/MeasureNameSanitizerIotDb.java
 
b/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/sanitize/MeasureNameSanitizerIotDb.java
similarity index 97%
rename from 
streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/MeasureNameSanitizerIotDb.java
rename to 
streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/sanitize/MeasureNameSanitizerIotDb.java
index bc30a4ed2d..c1aa656aee 100644
--- 
a/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/MeasureNameSanitizerIotDb.java
+++ 
b/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/sanitize/MeasureNameSanitizerIotDb.java
@@ -16,7 +16,7 @@
  *
  */
 
-package org.apache.streampipes.ts.store.iotdb;
+package org.apache.streampipes.ts.store.iotdb.sanitize;
 
 /**
  * Ensures that measurement names comply with IoTDB path specifications.
diff --git 
a/streampipes-ts-store-iotdb/src/test/java/org/apache/streampipes/ts/store/iotdb/sanitize/DataLakeMeasurementSanitizerIotDbTest.java
 
b/streampipes-ts-store-iotdb/src/test/java/org/apache/streampipes/ts/store/iotdb/sanitize/DataLakeMeasurementSanitizerIotDbTest.java
new file mode 100644
index 0000000000..13462de390
--- /dev/null
+++ 
b/streampipes-ts-store-iotdb/src/test/java/org/apache/streampipes/ts/store/iotdb/sanitize/DataLakeMeasurementSanitizerIotDbTest.java
@@ -0,0 +1,75 @@
+/*
+ * 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.ts.store.iotdb.sanitize;
+
+import org.apache.streampipes.client.api.IStreamPipesClient;
+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 static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.Mockito.mock;
+
+public class DataLakeMeasurementSanitizerIotDbTest {
+
+  private IStreamPipesClient clientMock;
+
+  @BeforeEach
+  public void setUp() {
+    this.clientMock = mock(IStreamPipesClient.class);
+  }
+
+  @Test
+  public void cleanDataLakeMeasure() {
+    var eventSchema = EventSchemaTestBuilder.create()
+        .withEventProperty(
+            EventPropertyPrimitiveTestBuilder
+                .create()
+                .withRuntimeName("timestamp")
+                .build())
+        .withEventProperty(
+            EventPropertyPrimitiveTestBuilder
+                .create()
+                .withRuntimeName("all")
+                .build())
+        .withEventProperty(
+            EventPropertyPrimitiveTestBuilder
+                .create()
+                .withRuntimeName("pressure")
+                .build())
+        .build();
+    var measure = new DataLakeMeasure(
+        "invalid.Measure",
+        "s0::%s".formatted("timestamp"),
+        eventSchema
+    );
+    var sanitizer = new DataLakeMeasurementSanitizerIotDb(clientMock, measure);
+
+    sanitizer.cleanDataLakeMeasure();
+    var result = sanitizer.getMeasure();
+
+    assertEquals("invalid_Measure", result.getMeasureName());
+    assertEquals(3, result.getEventSchema().getEventProperties().size());
+    assertEquals("timestamp_", 
result.getEventSchema().getEventProperties().get(0).getRuntimeName());
+    assertEquals("all_", 
result.getEventSchema().getEventProperties().get(1).getRuntimeName());
+    assertEquals("pressure", 
result.getEventSchema().getEventProperties().get(2).getRuntimeName());
+  }
+}
diff --git 
a/streampipes-ts-store-iotdb/src/test/java/org/apache/streampipes/ts/store/iotdb/IotDbNameSanitizerTest.java
 
b/streampipes-ts-store-iotdb/src/test/java/org/apache/streampipes/ts/store/iotdb/sanitize/IotDbNameSanitizerTest.java
similarity index 96%
rename from 
streampipes-ts-store-iotdb/src/test/java/org/apache/streampipes/ts/store/iotdb/IotDbNameSanitizerTest.java
rename to 
streampipes-ts-store-iotdb/src/test/java/org/apache/streampipes/ts/store/iotdb/sanitize/IotDbNameSanitizerTest.java
index 16cbf865b2..c900d5c5ac 100644
--- 
a/streampipes-ts-store-iotdb/src/test/java/org/apache/streampipes/ts/store/iotdb/IotDbNameSanitizerTest.java
+++ 
b/streampipes-ts-store-iotdb/src/test/java/org/apache/streampipes/ts/store/iotdb/sanitize/IotDbNameSanitizerTest.java
@@ -16,7 +16,7 @@
  *
  */
 
-package org.apache.streampipes.ts.store.iotdb;
+package org.apache.streampipes.ts.store.iotdb.sanitize;
 
 import org.junit.jupiter.api.Test;
 
diff --git 
a/streampipes-ts-store-iotdb/src/test/java/org/apache/streampipes/ts/store/iotdb/sanitize/MeasureNameSanitizerIotDbTest.java
 
b/streampipes-ts-store-iotdb/src/test/java/org/apache/streampipes/ts/store/iotdb/sanitize/MeasureNameSanitizerIotDbTest.java
index 6a11f4dee6..b71e8900d5 100644
--- 
a/streampipes-ts-store-iotdb/src/test/java/org/apache/streampipes/ts/store/iotdb/sanitize/MeasureNameSanitizerIotDbTest.java
+++ 
b/streampipes-ts-store-iotdb/src/test/java/org/apache/streampipes/ts/store/iotdb/sanitize/MeasureNameSanitizerIotDbTest.java
@@ -18,7 +18,6 @@
 
 package org.apache.streampipes.ts.store.iotdb.sanitize;
 
-import org.apache.streampipes.ts.store.iotdb.MeasureNameSanitizerIotDb;
 import org.junit.jupiter.api.Test;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
@@ -30,6 +29,7 @@ public class MeasureNameSanitizerIotDbTest {
     var sanitizer = new MeasureNameSanitizerIotDb();
 
     assertEquals("myMeasure", sanitizer.sanitize("myMeasure"));
+    assertEquals("my_Measure", sanitizer.sanitize("my.Measure"));
     assertEquals("我的措施", sanitizer.sanitize("我的措施"));
     assertEquals("my_Measure", sanitizer.sanitize("my-Measure"));
     assertEquals("my_Measure_", sanitizer.sanitize("my-Measure?"));

Reply via email to