|
|
|
@ -43,7 +43,6 @@ import org.thingsboard.server.queue.provider.TbRuleEngineQueueFactory; |
|
|
|
import org.thingsboard.server.service.cf.AbstractCalculatedFieldStateService; |
|
|
|
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; |
|
|
|
|
|
|
|
import java.util.Set; |
|
|
|
import java.util.concurrent.atomic.AtomicInteger; |
|
|
|
|
|
|
|
import static org.thingsboard.server.queue.common.AbstractTbQueueTemplate.bytesToString; |
|
|
|
@ -105,11 +104,6 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta |
|
|
|
this.stateProducer = (TbKafkaProducerTemplate<TbProtoQueueMsg<CalculatedFieldStateProto>>) queueFactory.createCalculatedFieldStateProducer(); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void restore(QueueKey queueKey, Set<TopicPartitionInfo> partitions) { |
|
|
|
stateService.update(queueKey, partitions, null); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
protected void doPersist(CalculatedFieldEntityCtxId stateId, CalculatedFieldStateProto stateMsgProto, TbCallback callback) { |
|
|
|
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, DataConstants.CF_STATES_QUEUE_NAME, stateId.tenantId(), stateId.entityId()); |
|
|
|
|