diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java index 6218f38edf..6a62c0b1e3 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java @@ -15,12 +15,14 @@ */ package org.thingsboard.server.service.cf.ctx.state; +import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonTypeInfo; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import java.util.List; +import java.util.TreeMap; import java.util.stream.Collectors; @JsonTypeInfo( @@ -34,16 +36,18 @@ import java.util.stream.Collectors; }) public interface ArgumentEntry { + @JsonIgnore ArgumentType getType(); Object getValue(); static ArgumentEntry createSingleValueArgument(KvEntry kvEntry) { - return new SingleValueArgumentEntry(kvEntry.getValue()); + return new SingleValueArgumentEntry(kvEntry); } static ArgumentEntry createLastRecordsArgument(List kvEntries) { - return new LastRecordsArgumentEntry(kvEntries.stream() .collect(Collectors.toMap(TsKvEntry::getTs, TsKvEntry::getValue))); + return new LastRecordsArgumentEntry(kvEntries.stream(). + collect(Collectors.toMap(TsKvEntry::getTs, TsKvEntry::getValue, (oldValue, newValue) -> newValue, TreeMap::new))); } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/LastRecordsArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/LastRecordsArgumentEntry.java index 39da1838bc..93fabd3cb8 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/LastRecordsArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/LastRecordsArgumentEntry.java @@ -15,24 +15,26 @@ */ package org.thingsboard.server.service.cf.ctx.state; +import com.fasterxml.jackson.annotation.JsonIgnore; import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; -import java.util.Map; +import java.util.TreeMap; @Data @NoArgsConstructor @AllArgsConstructor public class LastRecordsArgumentEntry implements ArgumentEntry { - private Map tsRecords; + private TreeMap tsRecords; @Override public ArgumentType getType() { return ArgumentType.LAST_RECORDS; } + @JsonIgnore @Override public Object getValue() { return tsRecords.values(); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/LastRecordsCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/LastRecordsCalculatedFieldState.java index 437707bfdb..dd69790236 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/LastRecordsCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/LastRecordsCalculatedFieldState.java @@ -19,16 +19,13 @@ import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.Data; import org.thingsboard.server.common.data.cf.CalculatedFieldType; +import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.common.data.cf.configuration.Output; -import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.service.cf.CalculatedFieldResult; -import java.util.ArrayList; -import java.util.Comparator; import java.util.HashMap; -import java.util.List; import java.util.Map; -import java.util.stream.Collectors; +import java.util.TreeMap; @Data public class LastRecordsCalculatedFieldState extends BaseCalculatedFieldState { @@ -44,16 +41,15 @@ public class LastRecordsCalculatedFieldState extends BaseCalculatedFieldState { @Override public void initState(Map argumentValues) { if (arguments == null) { - arguments = new HashMap<>(); + arguments = new TreeMap<>(); } argumentValues.forEach((key, argumentEntry) -> { LastRecordsArgumentEntry existingArgumentEntry = (LastRecordsArgumentEntry) - arguments.computeIfAbsent(key, k -> new LastRecordsArgumentEntry(new HashMap<>())); + arguments.computeIfAbsent(key, k -> new LastRecordsArgumentEntry(new TreeMap<>())); if (argumentEntry instanceof LastRecordsArgumentEntry lastRecordsArgumentEntry) { existingArgumentEntry.getTsRecords().putAll(lastRecordsArgumentEntry.getTsRecords()); - } else if (argumentEntry instanceof SingleValueArgumentEntry singleValueArgumentEntry - && singleValueArgumentEntry.getValue() instanceof TsKvEntry tsKvEntry) { - existingArgumentEntry.getTsRecords().put(tsKvEntry.getTs(), tsKvEntry.getValue()); + } else if (argumentEntry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { + existingArgumentEntry.getTsRecords().put(singleValueArgumentEntry.getTs(), singleValueArgumentEntry.getValue()); } }); } @@ -61,26 +57,22 @@ public class LastRecordsCalculatedFieldState extends BaseCalculatedFieldState { @Override public ListenableFuture performCalculation(CalculatedFieldCtx ctx) { Map resultMap = new HashMap<>(); - arguments.replaceAll((key, argumentEntry) -> { - int limit = ctx.getArguments().get(key).getLimit(); - - // TODO: implement removing if size > limit - - -// List limitedEntries = entries.stream() -// .sorted(Comparator.comparingLong(TsKvEntry::getTs).reversed()) -// .limit(limit) -// .collect(Collectors.toList()); -// -// Map valueWithTs = limitedEntries.stream() -// .collect(Collectors.toMap(TsKvEntry::getTs, TsKvEntry::getValue)); -// resultMap.put(key, valueWithTs); - -// return new LastRecordsArgumentEntry(limitedEntries); - return null; + arguments.forEach((key, argumentEntry) -> { + Argument argument = ctx.getArguments().get(key); + TreeMap tsRecords = ((LastRecordsArgumentEntry) argumentEntry).getTsRecords(); + if (tsRecords.size() > argument.getLimit()) { + tsRecords.pollFirstEntry(); + } + long necessaryIntervalTs = calculateIntervalStart(System.currentTimeMillis(), argument.getTimeWindow()); + tsRecords.entrySet().removeIf(tsRecord -> calculateIntervalStart(tsRecord.getKey(), argument.getTimeWindow()) < necessaryIntervalTs); + resultMap.put(key, tsRecords); }); Output output = ctx.getOutput(); return Futures.immediateFuture(new CalculatedFieldResult(output.getType(), output.getScope(), resultMap)); } + private long calculateIntervalStart(long ts, long interval) { + return (ts / interval) * interval; + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java index 64eb409d3a..fc97141806 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java @@ -21,9 +21,7 @@ import lombok.Data; import net.objecthunter.exp4j.Expression; import net.objecthunter.exp4j.ExpressionBuilder; import org.thingsboard.server.common.data.cf.CalculatedFieldType; -import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.common.data.cf.configuration.Output; -import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.service.cf.CalculatedFieldResult; import java.util.HashMap; @@ -51,7 +49,7 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState { customExpression.set(expr); } Map variables = new HashMap<>(); - this.arguments.forEach((k, v) -> variables.put(k, Double.parseDouble(((KvEntry) v.getValue()).getValueAsString()))); + this.arguments.forEach((k, v) -> variables.put(k, Double.parseDouble(v.getValue().toString()))); expr.setVariables(variables); double expressionResult = expr.evaluate(); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java index 504d213748..e0db8c50fb 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java @@ -15,17 +15,29 @@ */ package org.thingsboard.server.service.cf.ctx.state; -import lombok.AllArgsConstructor; import lombok.Data; -import lombok.NoArgsConstructor; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; +import org.thingsboard.server.common.data.kv.KvEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; @Data -@NoArgsConstructor -@AllArgsConstructor public class SingleValueArgumentEntry implements ArgumentEntry { + private long ts; private Object value; + public SingleValueArgumentEntry() { + } + + public SingleValueArgumentEntry(KvEntry entry) { + if (entry instanceof TsKvEntry) { + this.ts = ((TsKvEntry) entry).getTs(); + } else if (entry instanceof AttributeKvEntry) { + this.ts = ((AttributeKvEntry) entry).getLastUpdateTs(); + } + this.value = entry.getValue(); + } + @Override public ArgumentType getType() { return ArgumentType.SINGLE_VALUE;