Browse Source

Merge branch 'fix/cf-restore' of github.com:thingsboard/thingsboard into fix/cf-restore-master

pull/14573/head
Viacheslav Klimov 10 months ago
parent
commit
5bb8f896fd
  1. 1
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
  2. 7
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java
  3. 3
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldStateRestoreMsg.java
  4. 24
      application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java
  5. 8
      application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java
  6. 36
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java
  7. 18
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java

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

@ -132,6 +132,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
} else { } else {
removeState(cfId); removeState(cfId);
} }
msg.getCallback().onSuccess();
} }
public void process(CalculatedFieldStatePartitionRestoreMsg msg) { public void process(CalculatedFieldStatePartitionRestoreMsg msg) {

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

@ -167,10 +167,13 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
if (ctx != null) { if (ctx != null) {
msg.setCtx(ctx); msg.setCtx(ctx);
log.debug("Pushing CF state restore msg to specific actor [{}]", msg.getId().entityId()); log.debug("[{}] Pushing CF state restore msg to specific actor [{}]", tenantId, msg.getId().entityId());
getOrCreateActor(msg.getId().entityId()).tellWithHighPriority(msg); getOrCreateActor(msg.getId().entityId()).tellWithHighPriority(msg);
} else { } else if (msg.getState() != null) {
log.debug("[{}] Received CF state restore msg for non-existing CF [{}]. Removing state", tenantId, cfId);
cfStateService.deleteState(msg.getId(), msg.getCallback()); cfStateService.deleteState(msg.getId(), msg.getCallback());
} else {
msg.getCallback().onSuccess();
} }
} }

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

@ -19,6 +19,7 @@ import lombok.Data;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.MsgType; import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.ToCalculatedFieldSystemMsg; import org.thingsboard.server.common.msg.ToCalculatedFieldSystemMsg;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx;
@ -30,6 +31,7 @@ public class CalculatedFieldStateRestoreMsg implements ToCalculatedFieldSystemMs
private final CalculatedFieldEntityCtxId id; private final CalculatedFieldEntityCtxId id;
private final CalculatedFieldState state; private final CalculatedFieldState state;
private final TopicPartitionInfo partition; private final TopicPartitionInfo partition;
private final TbCallback callback;
private CalculatedFieldCtx ctx; private CalculatedFieldCtx ctx;
@Override @Override
@ -41,4 +43,5 @@ public class CalculatedFieldStateRestoreMsg implements ToCalculatedFieldSystemMs
public TenantId getTenantId() { public TenantId getTenantId() {
return id.tenantId(); return id.tenantId();
} }
} }

24
application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java

@ -43,7 +43,6 @@ import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChain;
import org.thingsboard.server.common.data.rule.RuleChainType; import org.thingsboard.server.common.data.rule.RuleChainType;
import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.TbActorMsg;
import org.thingsboard.server.common.msg.TbActorStopReason; import org.thingsboard.server.common.msg.TbActorStopReason;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
@ -139,13 +138,22 @@ public class TenantActor extends RuleChainManagerActor {
@Override @Override
protected boolean doProcess(TbActorMsg msg) { protected boolean doProcess(TbActorMsg msg) {
if (cantFindTenant) { if (cantFindTenant) {
log.info("[{}] Processing missing Tenant msg: {}", tenantId, msg); log.debug("[{}] Processing message for non-existing tenant: {}", tenantId, msg);
if (msg.getMsgType().equals(MsgType.QUEUE_TO_RULE_ENGINE_MSG)) { switch (msg.getMsgType()) {
QueueToRuleEngineMsg queueMsg = (QueueToRuleEngineMsg) msg; case QUEUE_TO_RULE_ENGINE_MSG -> {
queueMsg.getMsg().getCallback().onSuccess(); ((QueueToRuleEngineMsg) msg).getMsg().getCallback().onSuccess();
} else if (msg.getMsgType().equals(MsgType.TRANSPORT_TO_DEVICE_ACTOR_MSG)) { }
TransportToDeviceActorMsgWrapper transportMsg = (TransportToDeviceActorMsgWrapper) msg; case TRANSPORT_TO_DEVICE_ACTOR_MSG -> {
transportMsg.getCallback().onSuccess(); ((TransportToDeviceActorMsgWrapper) msg).getCallback().onSuccess();
}
case CF_STATE_RESTORE_MSG -> {
((CalculatedFieldStateRestoreMsg) msg).getCallback().onSuccess();
}
default -> {
if (!log.isDebugEnabled()) {
log.info("[{}] Processing message for non-existing tenant: {}", tenantId, msg);
}
}
} }
return true; return true;
} }

8
application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java

@ -68,7 +68,7 @@ public abstract class AbstractCalculatedFieldStateService implements CalculatedF
protected abstract void doRemove(CalculatedFieldEntityCtxId stateId, TbCallback callback); protected abstract void doRemove(CalculatedFieldEntityCtxId stateId, TbCallback callback);
protected void processRestoredState(CalculatedFieldStateProto stateMsg, TopicPartitionInfo partition) { protected void processRestoredState(CalculatedFieldStateProto stateMsg, TopicPartitionInfo partition, TbCallback callback) {
var id = fromProto(stateMsg.getId()); var id = fromProto(stateMsg.getId());
if (partition == null) { if (partition == null) {
try { try {
@ -79,12 +79,12 @@ public abstract class AbstractCalculatedFieldStateService implements CalculatedF
} }
} }
var state = fromProto(id, stateMsg); var state = fromProto(id, stateMsg);
processRestoredState(id, state, partition); processRestoredState(id, state, partition, callback);
} }
protected void processRestoredState(CalculatedFieldEntityCtxId id, CalculatedFieldState state, TopicPartitionInfo partition) { protected void processRestoredState(CalculatedFieldEntityCtxId id, CalculatedFieldState state, TopicPartitionInfo partition, TbCallback callback) {
partition = partition.withTopic(DataConstants.CF_STATES_QUEUE_NAME); partition = partition.withTopic(DataConstants.CF_STATES_QUEUE_NAME);
actorSystemContext.tellWithHighPriority(new CalculatedFieldStateRestoreMsg(id, state, partition)); actorSystemContext.tellWithHighPriority(new CalculatedFieldStateRestoreMsg(id, state, partition, callback));
} }
@Override @Override

36
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java

@ -43,6 +43,8 @@ import org.thingsboard.server.queue.provider.TbRuleEngineQueueFactory;
import org.thingsboard.server.service.cf.AbstractCalculatedFieldStateService; import org.thingsboard.server.service.cf.AbstractCalculatedFieldStateService;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicInteger;
import static org.thingsboard.server.queue.common.AbstractTbQueueTemplate.bytesToString; import static org.thingsboard.server.queue.common.AbstractTbQueueTemplate.bytesToString;
@ -61,6 +63,8 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta
@Value("${queue.calculated_fields.poll_interval:25}") @Value("${queue.calculated_fields.poll_interval:25}")
private long pollInterval; private long pollInterval;
@Value("${queue.calculated_fields.pack_processing_timeout:60000}")
private long packProcessingTimeout;
private TbKafkaProducerTemplate<TbProtoQueueMsg<CalculatedFieldStateProto>> stateProducer; private TbKafkaProducerTemplate<TbProtoQueueMsg<CalculatedFieldStateProto>> stateProducer;
@ -74,21 +78,39 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta
.topic(partitionService.getTopic(queueKey)) .topic(partitionService.getTopic(queueKey))
.pollInterval(pollInterval) .pollInterval(pollInterval)
.msgPackProcessor((msgs, consumer, consumerKey, config) -> { .msgPackProcessor((msgs, consumer, consumerKey, config) -> {
CountDownLatch completionLatch = new CountDownLatch(msgs.size());
for (TbProtoQueueMsg<CalculatedFieldStateProto> msg : msgs) { for (TbProtoQueueMsg<CalculatedFieldStateProto> msg : msgs) {
TbCallback callback = new TbCallback() {
@Override
public void onSuccess() {
int processedMsgCount = counter.incrementAndGet();
if (processedMsgCount % 10000 == 0) {
log.info("Processed {} CF state messages", processedMsgCount);
}
completionLatch.countDown();
}
@Override
public void onFailure(Throwable t) {
log.error("Failed to process CF state message: {}", msg, t);
completionLatch.countDown();
}
};
try { try {
if (msg.getValue() != null) { if (msg.getValue() != null) {
processRestoredState(msg.getValue(), consumerKey.partition()); processRestoredState(msg.getValue(), consumerKey.partition(), callback);
} else { } else {
processRestoredState(getStateId(msg.getHeaders()), null, consumerKey.partition()); processRestoredState(getStateId(msg.getHeaders()), null, consumerKey.partition(), callback);
} }
} catch (Throwable t) { } catch (Throwable t) {
log.error("Failed to process state message: {}", msg, t); callback.onFailure(t);
} }
}
int processedMsgCount = counter.incrementAndGet(); boolean success = completionLatch.await(packProcessingTimeout, TimeUnit.MILLISECONDS);
if (processedMsgCount % 10000 == 0) { if (!success) {
log.info("Processed {} calculated field state msgs", processedMsgCount); log.error("Timeout to process CF state messages pack of size {}", msgs.size());
}
} }
}) })
.consumerCreator((queueConfig, tpi) -> queueFactory.createCalculatedFieldStateConsumer()) .consumerCreator((queueConfig, tpi) -> queueFactory.createCalculatedFieldStateConsumer())

18
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.service.cf.ctx.state; package org.thingsboard.server.service.cf.ctx.state;
import com.google.protobuf.InvalidProtocolBufferException;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
@ -62,11 +63,22 @@ public class RocksDBCalculatedFieldStateService extends AbstractCalculatedFieldS
public void restore(QueueKey queueKey, Set<TopicPartitionInfo> partitions) { public void restore(QueueKey queueKey, Set<TopicPartitionInfo> partitions) {
if (stateService.getPartitions().isEmpty()) { if (stateService.getPartitions().isEmpty()) {
cfRocksDb.forEach((key, value) -> { cfRocksDb.forEach((key, value) -> {
CalculatedFieldStateProto stateMsg;
try { try {
processRestoredState(CalculatedFieldStateProto.parseFrom(value), null); stateMsg = CalculatedFieldStateProto.parseFrom(value);
} catch (Exception e) { } catch (InvalidProtocolBufferException e) {
log.error("[{}] Failed to process restored state", key, e); log.error("Failed to parse CalculatedFieldStateProto for key {}", key, e);
return;
} }
processRestoredState(stateMsg, null, new TbCallback() {
@Override
public void onSuccess() {}
@Override
public void onFailure(Throwable t) {
log.error("Failed to process CF state message: {}", stateMsg, t);
}
});
}); });
} }
super.restore(queueKey, partitions); super.restore(queueKey, partitions);

Loading…
Cancel
Save