Browse Source

review fixes

pull/14253/head
IrynaMatveieva 11 months ago
parent
commit
7b9b15d0f2
  1. 19
      application/src/main/java/org/thingsboard/server/actors/app/AppActor.java
  2. 13
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
  3. 5
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java
  4. 3
      application/src/main/java/org/thingsboard/server/actors/calculatedField/EntityInitCalculatedFieldMsg.java
  5. 16
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java
  6. 6
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java
  7. 3
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java
  8. 2
      application/src/main/resources/thingsboard.yml

19
application/src/main/java/org/thingsboard/server/actors/app/AppActor.java

@ -166,16 +166,17 @@ public class AppActor extends ContextAwareActor {
private void onComponentLifecycleMsg(ComponentLifecycleMsg msg) { private void onComponentLifecycleMsg(ComponentLifecycleMsg msg) {
TbActorRef target = null; TbActorRef target = null;
if (TenantId.SYS_TENANT_ID.equals(msg.getTenantId())) { if (TenantId.SYS_TENANT_ID.equals(msg.getTenantId())) {
if (msg.getEntityId() instanceof TenantProfileId tenantProfileId) { if (systemContext.isTenantComponentsInitEnabled()) {
tenantService.findTenantIdsByTenantProfileId(tenantProfileId).forEach(tenantId -> { if (msg.getEntityId() instanceof TenantProfileId tenantProfileId) {
TbActorRef tenantActor = getOrCreateTenantActor(tenantId).orElseGet(() -> { tenantService.findTenantIdsByTenantProfileId(tenantProfileId).forEach(tenantId -> {
log.debug("Ignoring component lifecycle msg for tenant {} because it is not managed by this service", tenantId); getOrCreateTenantActor(tenantId).ifPresentOrElse(tenantActor -> {
return null; log.debug("[{}] Sending component lifecycle msg for tenant.", tenantId);
tenantActor.tellWithHighPriority(msg);
}, () -> {
log.debug("Ignoring component lifecycle msg for tenant {} because it is not managed by this service", tenantId);
});
}); });
if (tenantActor != null) { }
tenantActor.tellWithHighPriority(msg);
}
});
} }
if (!msg.getEntityId().getEntityType().isOneOf(EntityType.TENANT_PROFILE, EntityType.TB_RESOURCE)) { if (!msg.getEntityId().getEntityType().isOneOf(EntityType.TENANT_PROFILE, EntityType.TB_RESOURCE)) {
log.warn("Message has system tenant id: {}", msg); log.warn("Message has system tenant id: {}", msg);

13
application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java

@ -123,7 +123,6 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
if (state != null) { if (state != null) {
state.setCtx(msg.getCtx(), actorCtx); state.setCtx(msg.getCtx(), actorCtx);
state.setPartition(msg.getPartition()); state.setPartition(msg.getPartition());
state.init(true);
states.put(cfId, state); states.put(cfId, state);
} else { } else {
removeState(cfId); removeState(cfId);
@ -134,7 +133,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
log.debug("Processing CF state partition restore msg: {}", msg); log.debug("Processing CF state partition restore msg: {}", msg);
for (CalculatedFieldState state : states.values()) { for (CalculatedFieldState state : states.values()) {
if (msg.getPartition().equals(state.getPartition())) { if (msg.getPartition().equals(state.getPartition())) {
state.init(false); state.init(true);
} }
} }
} }
@ -159,12 +158,10 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
} else { } else {
state.setCtx(ctx, actorCtx); state.setCtx(ctx, actorCtx);
} }
if (msg.getStateAction() != StateAction.REFRESH_CTX) { if (state.isSizeOk()) {
if (state.isSizeOk()) { processStateIfReady(state, Collections.emptyMap(), ctx, Collections.singletonList(ctx.getCfId()), null, null, msg.getCallback());
processStateIfReady(state, Collections.emptyMap(), ctx, Collections.singletonList(ctx.getCfId()), null, null, msg.getCallback()); } else {
} else { throw new RuntimeException(ctx.getSizeExceedsLimitMessage());
throw new RuntimeException(ctx.getSizeExceedsLimitMessage());
}
} }
} catch (Exception e) { } catch (Exception e) {
log.debug("[{}][{}] Failed to initialize CF state", entityId, ctx.getCfId(), e); log.debug("[{}][{}] Failed to initialize CF state", entityId, ctx.getCfId(), e);

5
application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java

@ -260,10 +260,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
calculatedFields.values().stream(), calculatedFields.values().stream(),
entityIdCalculatedFields.values().stream().flatMap(Collection::stream) entityIdCalculatedFields.values().stream().flatMap(Collection::stream)
).forEach(CalculatedFieldCtx::updateTenantProfileProperties); ).forEach(CalculatedFieldCtx::updateTenantProfileProperties);
callback.onSuccess();
calculatedFields.values().forEach(ctx -> {
applyToTargetCfEntityActors(ctx, callback, (id, cb) -> initCfForEntity(id, ctx, StateAction.REFRESH_CTX, cb));
});
} }
private void onEntityCreated(ComponentLifecycleMsg msg, TbCallback callback) { private void onEntityCreated(ComponentLifecycleMsg msg, TbCallback callback) {

3
application/src/main/java/org/thingsboard/server/actors/calculatedField/EntityInitCalculatedFieldMsg.java

@ -39,7 +39,6 @@ public class EntityInitCalculatedFieldMsg implements ToCalculatedFieldSystemMsg
INIT, INIT,
REINIT, REINIT,
RECREATE, RECREATE,
REPROCESS, REPROCESS
REFRESH_CTX
} }
} }

16
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java

@ -221,16 +221,18 @@ public class CalculatedFieldCtx implements Closeable {
return true; return true;
} }
} }
long reevaluationIntervalMillis = TimeUnit.SECONDS.toMillis(systemContext.getAlarmRulesReevaluationInterval());
boolean requiresScheduledReevaluation = calculatedField.getConfiguration().requiresScheduledReevaluation(); boolean requiresScheduledReevaluation = calculatedField.getConfiguration().requiresScheduledReevaluation();
if (requiresScheduledReevaluation) { if (calculatedField.getConfiguration() instanceof AlarmCalculatedFieldConfiguration) {
if (now + cfCheckIntervalMillis >= lastReevaluationTs + reevaluationIntervalMillis) { long reevaluationIntervalMillis = TimeUnit.SECONDS.toMillis(systemContext.getAlarmRulesReevaluationInterval());
lastReevaluationTs = now; if (requiresScheduledReevaluation) {
return true; if (now + cfCheckIntervalMillis >= lastReevaluationTs + reevaluationIntervalMillis) {
lastReevaluationTs = now;
return true;
}
return false;
} }
return false;
} }
return false; return requiresScheduledReevaluation;
} }
public void init() { public void init() {

6
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java

@ -127,7 +127,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt
}); });
} }
public void fillMissingIntervals() { private void fillMissingIntervals() {
ZoneId zoneId = interval.getZoneId(); ZoneId zoneId = interval.getZoneId();
long currentIntervalEndTs = interval.getCurrentIntervalEndTs(); long currentIntervalEndTs = interval.getCurrentIntervalEndTs();
@ -253,12 +253,12 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt
metricsNode.put(metricName, JacksonUtil.toString(resultValue)); metricsNode.put(metricName, JacksonUtil.toString(resultValue));
} }
} }
ObjectNode resultNode = JacksonUtil.newObjectNode();
if (!metricsNode.isEmpty()) { if (!metricsNode.isEmpty()) {
ObjectNode resultNode = JacksonUtil.newObjectNode();
resultNode.put("ts", interval.getEndTs() - 1); resultNode.put("ts", interval.getEndTs() - 1);
resultNode.set("values", metricsNode); resultNode.set("values", metricsNode);
result.add(resultNode);
} }
result.add(resultNode);
}); });
return result; return result;
} }

3
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java

@ -124,9 +124,6 @@ public class AlarmCalculatedFieldState extends BaseCalculatedFieldState {
@Override @Override
public void init(boolean restored) { public void init(boolean restored) {
super.init(restored); super.init(restored);
if (restored) {
return;
}
AtomicBoolean reevalNeeded = new AtomicBoolean(false); AtomicBoolean reevalNeeded = new AtomicBoolean(false);
Map<AlarmSeverity, AlarmRule> createRules = configuration.getCreateRules(); Map<AlarmSeverity, AlarmRule> createRules = configuration.getCreateRules();
for (AlarmSeverity severity : AlarmSeverity.values()) { for (AlarmSeverity severity : AlarmSeverity.values()) {

2
application/src/main/resources/thingsboard.yml

@ -529,7 +529,7 @@ actors:
configuration: "${ACTORS_CALCULATED_FIELD_DEBUG_MODE_RATE_LIMITS_PER_TENANT_CONFIGURATION:50000:3600}" configuration: "${ACTORS_CALCULATED_FIELD_DEBUG_MODE_RATE_LIMITS_PER_TENANT_CONFIGURATION:50000:3600}"
# Time in seconds to receive calculation result. # Time in seconds to receive calculation result.
calculation_timeout: "${ACTORS_CALCULATION_TIMEOUT_SEC:5}" calculation_timeout: "${ACTORS_CALCULATION_TIMEOUT_SEC:5}"
# Interval in seconds to re-evaluate calculated fields that have a time schedule. 1 minute by default. # Interval in seconds to check calculated fields for re-evaluation interval. 1 minute by default.
check_interval: "${ACTORS_CALCULATED_FIELDS_CHECK_INTERVAL_SEC:60}" check_interval: "${ACTORS_CALCULATED_FIELDS_CHECK_INTERVAL_SEC:60}"
alarms: alarms:
# Interval in seconds to re-evaluate Alarm rules that have a time schedule. 2 minutes by default. # Interval in seconds to re-evaluate Alarm rules that have a time schedule. 2 minutes by default.

Loading…
Cancel
Save