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

bossenti pushed a commit to branch add-session-pool-provider
in repository https://gitbox.apache.org/repos/asf/streampipes.git

commit 4aa62b66be06ce0fc86f4687c5ee2c67c7215033
Author: bossenti <[email protected]>
AuthorDate: Mon May 6 10:11:45 2024 +0200

    feat: introduce IoTDB session handling
---
 .../apache/streampipes/commons/constants/Envs.java |  4 ++
 .../commons/environment/DefaultEnvironment.java    | 20 +++++++++
 .../commons/environment/Environment.java           |  8 ++++
 .../ts/store/iotdb/IotDbSessionProvider.java       | 47 ++++++++++++++++++++++
 4 files changed, 79 insertions(+)

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 edcbb32999..672014852e 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
@@ -64,6 +64,10 @@ public enum Envs {
   SP_TS_STORAGE_ORG("SP_TS_STORAGE_ORG", "sp"),
 
   SP_TS_STORAGE_BUCKET("SP_TS_STORAGE_BUCKET", "sp"),
+  
SP_TS_STORAGE_IOT_DB_SESSION_POOL_SIZE("SP_TS_STORAGE_IOT_DB_SESSION_POOL_SIZE",
 "10"),
+  
SP_TS_STORAGE_IOT_DB_SESSION_POOL_ENABLE_COMPRESSION("SP_TS_STORAGE_IOT_DB_SESSION_POOL_ENABLE_COMPRESSION",
 "false"),
+  SP_TS_STORAGE_IOT_DB_USER("SP_TS_STORAGE_IOT_DB_USER", "root"),
+  SP_TS_STORAGE_IOT_DB_PASSWORD("SP_TS_STORAGE_IOT_DB_PASSWORD", "root"),
 
   SP_FLINK_JAR_FILE_LOC(
       "SP_FLINK_JAR_FILE_LOC",
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 69b4033551..463a3a1d3e 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
@@ -90,6 +90,26 @@ public class DefaultEnvironment implements Environment {
     return new StringEnvironmentVariable(Envs.SP_TS_STORAGE_BUCKET);
   }
 
+  @Override
+  public IntEnvironmentVariable getIotDbSessionPoolSize(){
+    return new 
IntEnvironmentVariable(Envs.SP_TS_STORAGE_IOT_DB_SESSION_POOL_SIZE);
+  }
+
+  @Override
+  public BooleanEnvironmentVariable getIotDbSessionEnableCompression(){
+    return new 
BooleanEnvironmentVariable(Envs.SP_TS_STORAGE_IOT_DB_SESSION_POOL_ENABLE_COMPRESSION);
+  }
+
+  @Override
+  public StringEnvironmentVariable getIotDbUser(){
+    return new StringEnvironmentVariable(Envs.SP_TS_STORAGE_IOT_DB_USER);
+  }
+
+  @Override
+  public StringEnvironmentVariable getIotDbPassword(){
+    return new StringEnvironmentVariable(Envs.SP_TS_STORAGE_IOT_DB_PASSWORD);
+  }
+
   @Override
   public StringEnvironmentVariable getCouchDbProtocol() {
     return new StringEnvironmentVariable(Envs.SP_COUCHDB_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 e190402378..aabae364b4 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
@@ -52,6 +52,14 @@ public interface Environment {
 
   StringEnvironmentVariable getTsStorageBucket();
 
+  IntEnvironmentVariable getIotDbSessionPoolSize();
+
+  BooleanEnvironmentVariable getIotDbSessionEnableCompression();
+
+  StringEnvironmentVariable getIotDbUser();
+
+  StringEnvironmentVariable getIotDbPassword();
+
   // CouchDB env variables
 
   StringEnvironmentVariable getCouchDbProtocol();
diff --git 
a/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/IotDbSessionProvider.java
 
b/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/IotDbSessionProvider.java
new file mode 100644
index 0000000000..4a46ede043
--- /dev/null
+++ 
b/streampipes-ts-store-iotdb/src/main/java/org/apache/streampipes/ts/store/iotdb/IotDbSessionProvider.java
@@ -0,0 +1,47 @@
+/*
+ * 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.session.pool.SessionPool;
+import org.apache.streampipes.commons.environment.Environment;
+
+/**
+ * This class provides a method to retrieve a session pool for IoT DB 
operations based on the given environment configuration.
+ */
+public class IotDbSessionProvider {
+
+  /**
+   * Retrieves a session pool for IoT DB operations.
+   * <p>
+   * The session pool is configured by the StreamPipes environment and 
respective environment variables.
+   *
+   * @param environment the environment configuration containing IoT DB 
connection details.
+   * @return a SessionPool configured based on the provided environment.
+   */
+  public SessionPool getSessionPool(Environment environment) {
+    return new SessionPool.Builder()
+        .maxSize(environment.getIotDbSessionPoolSize().getValueOrDefault())
+        
.enableCompression(environment.getIotDbSessionEnableCompression().getValueOrDefault())
+        .host(environment.getTsStorageHost().getValueOrDefault())
+        .port(environment.getTsStoragePort().getValueOrDefault())
+        .user(environment.getIotDbUser().getValueOrDefault())
+        .password(environment.getIotDbPassword().getValueOrDefault())
+        .build();
+  }
+}

Reply via email to