From e561e268719dc353b15d03c0c0caf9fccf2abeb2 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Tue, 11 Jun 2024 10:41:12 +0300 Subject: [PATCH] added possibility to set ttl via config --- .../TbSaveToCustomCassandraTableNode.java | 52 ++++++++++++++++--- ...CustomCassandraTableNodeConfiguration.java | 2 + 2 files changed, 48 insertions(+), 6 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbSaveToCustomCassandraTableNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbSaveToCustomCassandraTableNode.java index 1aad815dfa..9a3950a1c2 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbSaveToCustomCassandraTableNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbSaveToCustomCassandraTableNode.java @@ -21,6 +21,8 @@ import com.datastax.oss.driver.api.core.cql.BoundStatement; import com.datastax.oss.driver.api.core.cql.BoundStatementBuilder; import com.datastax.oss.driver.api.core.cql.PreparedStatement; import com.datastax.oss.driver.api.core.cql.Statement; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.base.Function; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; @@ -28,6 +30,7 @@ import com.google.gson.JsonElement; import com.google.gson.JsonObject; import com.google.gson.JsonParser; import com.google.gson.JsonPrimitive; +import jakarta.annotation.Nullable; import lombok.extern.slf4j.Slf4j; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; @@ -35,20 +38,23 @@ import org.thingsboard.rule.engine.api.TbNode; import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; +import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.rule.RuleChainType; +import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; +import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.dao.cassandra.CassandraCluster; import org.thingsboard.server.dao.cassandra.guava.GuavaSession; import org.thingsboard.server.dao.nosql.CassandraStatementTask; import org.thingsboard.server.dao.nosql.TbResultSetFuture; -import jakarta.annotation.Nullable; import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import static org.thingsboard.common.util.DonAsynchron.withCallback; @@ -57,6 +63,7 @@ import static org.thingsboard.common.util.DonAsynchron.withCallback; @RuleNode(type = ComponentType.ACTION, name = "save to custom table", configClazz = TbSaveToCustomCassandraTableNodeConfiguration.class, + version = 1, nodeDescription = "Node stores data from incoming Message payload to the Cassandra database into the predefined custom table" + " that should have cs_tb_ prefix, to avoid the data insertion to the common TB tables.
" + "Note: rule node can be used only for Cassandra DB.", @@ -81,6 +88,7 @@ public class TbSaveToCustomCassandraTableNode implements TbNode { private PreparedStatement saveStmt; private ExecutorService readResultsProcessingExecutor; private Map fieldsMap; + private long ttl; @Override public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { @@ -88,15 +96,25 @@ public class TbSaveToCustomCassandraTableNode implements TbNode { cassandraCluster = ctx.getCassandraCluster(); if (cassandraCluster == null) { throw new RuntimeException("Unable to connect to Cassandra database"); - } else { - startExecutor(); - saveStmt = getSaveStmt(); + } + ctx.addTenantProfileListener(this::onTenantProfileUpdate); + onTenantProfileUpdate(ctx.getTenantProfile()); + startExecutor(); + saveStmt = getSaveStmt(); + } + + void onTenantProfileUpdate(TenantProfile tenantProfile) { + DefaultTenantProfileConfiguration configuration = (DefaultTenantProfileConfiguration) tenantProfile.getProfileData().getConfiguration(); + long tenantProfileDefaultStorageTtl = TimeUnit.DAYS.toSeconds(configuration.getDefaultStorageTtlDays()); + ttl = config.getDefaultTTL(); + if (ttl == 0L) { + ttl = tenantProfileDefaultStorageTtl; } } @Override public void onMsg(TbContext ctx, TbMsg msg) { - withCallback(save(msg, ctx), aVoid -> ctx.tellSuccess(msg), e -> ctx.tellFailure(msg, e), ctx.getDbCallbackExecutor()); + withCallback(save(msg, ctx, ttl), aVoid -> ctx.tellSuccess(msg), e -> ctx.tellFailure(msg, e), ctx.getDbCallbackExecutor()); } @Override @@ -163,10 +181,13 @@ public class TbSaveToCustomCassandraTableNode implements TbNode { query.append("?, "); } } + if (ttl > 0) { + query.append(" USING TTL ?"); + } return query.toString(); } - private ListenableFuture save(TbMsg msg, TbContext ctx) { + private ListenableFuture save(TbMsg msg, TbContext ctx, long ttl) { JsonElement data = JsonParser.parseString(msg.getData()); if (!data.isJsonObject()) { throw new IllegalStateException("Invalid message structure, it is not a JSON Object:" + data); @@ -204,6 +225,9 @@ public class TbSaveToCustomCassandraTableNode implements TbNode { } i.getAndIncrement(); }); + if (ttl > 0) { + stmtBuilder.setInt(i.get(), (int) ttl); + } return getFuture(executeAsyncWrite(ctx, stmtBuilder.build()), rs -> null); } } @@ -240,4 +264,20 @@ public class TbSaveToCustomCassandraTableNode implements TbNode { }, readResultsProcessingExecutor); } + @Override + public TbPair upgrade(int fromVersion, JsonNode oldConfiguration) throws TbNodeException { + boolean hasChanges = false; + switch (fromVersion) { + case 0: + if (!oldConfiguration.has("defaultTTL")) { + hasChanges = true; + ((ObjectNode) oldConfiguration).put("defaultTTL", 0); + } + break; + default: + break; + } + return new TbPair<>(hasChanges, oldConfiguration); + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbSaveToCustomCassandraTableNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbSaveToCustomCassandraTableNodeConfiguration.java index 0a5f153192..8d3d78e8b9 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbSaveToCustomCassandraTableNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbSaveToCustomCassandraTableNodeConfiguration.java @@ -27,11 +27,13 @@ public class TbSaveToCustomCassandraTableNodeConfiguration implements NodeConfig private String tableName; private Map fieldsMapping; + private long defaultTTL; @Override public TbSaveToCustomCassandraTableNodeConfiguration defaultConfiguration() { TbSaveToCustomCassandraTableNodeConfiguration configuration = new TbSaveToCustomCassandraTableNodeConfiguration(); + configuration.setDefaultTTL(0L); configuration.setTableName(""); Map map = new HashMap<>(); map.put("", "");