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(); + } +}
