Browse Source

EDQS readiness check; refactor API enabling

pull/13362/head
ViacheslavKlimov 1 year ago
parent
commit
b97c1888f1
  1. 7
      application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java
  2. 2
      application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java
  3. 26
      application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsApiService.java
  4. 137
      application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsService.java
  5. 2
      application/src/main/resources/thingsboard.yml
  6. 13
      application/src/test/java/org/thingsboard/server/controller/EdqsEntityQueryControllerTest.java
  7. 7
      application/src/test/java/org/thingsboard/server/controller/EntityQueryControllerTest.java
  8. 9
      application/src/test/java/org/thingsboard/server/service/entitiy/EdqsEntityServiceTest.java
  9. 71
      common/data/src/main/java/org/thingsboard/server/common/data/edqs/EdqsState.java
  10. 4
      common/data/src/main/java/org/thingsboard/server/common/data/edqs/ToCoreEdqsMsg.java
  11. 6
      common/data/src/main/java/org/thingsboard/server/common/data/edqs/ToCoreEdqsRequest.java
  12. 21
      common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java
  13. 6
      common/message/src/main/java/org/thingsboard/server/common/msg/edqs/EdqsApiService.java
  14. 5
      common/message/src/main/java/org/thingsboard/server/common/msg/edqs/EdqsService.java
  15. 1
      common/proto/src/main/proto/queue.proto
  16. 8
      common/queue/src/main/java/org/thingsboard/server/queue/common/state/KafkaQueueStateService.java
  17. 21
      common/queue/src/main/java/org/thingsboard/server/queue/common/state/QueueStateService.java
  18. 11
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java
  19. 2
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/DiscoveryService.java
  20. 9
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java
  21. 2
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java
  22. 12
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java
  23. 10
      dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java
  24. 15
      dao/src/main/java/org/thingsboard/server/dao/sql/query/DummyEdqsApiService.java
  25. 11
      dao/src/main/java/org/thingsboard/server/dao/sql/query/DummyEdqsService.java

7
application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java

@ -29,6 +29,7 @@ import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.ResponseBody;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.context.request.async.DeferredResult;
import org.thingsboard.server.common.data.edqs.EdqsState;
import org.thingsboard.server.common.data.edqs.ToCoreEdqsRequest;
import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.TenantId;
@ -149,9 +150,9 @@ public class EntityQueryController extends BaseController {
}
@PreAuthorize("hasAnyAuthority('SYS_ADMIN')")
@GetMapping("/edqs/enabled")
public boolean isEdqsApiEnabled() {
return edqsApiService.isEnabled();
@GetMapping("/edqs/state")
public EdqsState getEdqsState() {
return edqsService.getState();
}
}

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

@ -72,7 +72,7 @@ public abstract class AbstractCalculatedFieldStateService implements CalculatedF
@Override
public void restore(QueueKey queueKey, Set<TopicPartitionInfo> partitions) {
stateService.update(queueKey, partitions);
stateService.update(queueKey, partitions, null);
}
@Override

26
application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsApiService.java

@ -22,7 +22,6 @@ import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil;
@ -51,11 +50,6 @@ public class DefaultEdqsApiService implements EdqsApiService {
private final EdqsClientQueueFactory queueFactory;
private TbQueueRequestTemplate<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<FromEdqsMsg>> requestTemplate;
@Value("${queue.edqs.api.auto_enable:true}")
private boolean autoEnable;
private Boolean apiEnabled = null;
@PostConstruct
private void init() {
requestTemplate = queueFactory.createEdqsRequestTemplate();
@ -85,31 +79,11 @@ public class DefaultEdqsApiService implements EdqsApiService {
}, MoreExecutors.directExecutor());
}
@Override
public boolean isEnabled() {
return Boolean.TRUE.equals(apiEnabled);
}
@Override
public void setEnabled(boolean enabled) {
if (enabled) {
log.info("Enabling EDQS API");
} else {
log.info("Disabling EDQS API");
}
apiEnabled = enabled;
}
@Override
public boolean isSupported() {
return true;
}
@Override
public boolean isAutoEnable() {
return autoEnable;
}
@PreDestroy
private void stop() {
requestTemplate.stop();

137
application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsService.java

@ -20,11 +20,13 @@ import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.Getter;
import lombok.NoArgsConstructor;
import lombok.RequiredArgsConstructor;
import lombok.SneakyThrows;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service;
@ -36,6 +38,8 @@ import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.ObjectType;
import org.thingsboard.server.common.data.edqs.EdqsEventType;
import org.thingsboard.server.common.data.edqs.EdqsObject;
import org.thingsboard.server.common.data.edqs.EdqsState;
import org.thingsboard.server.common.data.edqs.EdqsState.EdqsSyncStatus;
import org.thingsboard.server.common.data.edqs.EdqsSyncRequest;
import org.thingsboard.server.common.data.edqs.Entity;
import org.thingsboard.server.common.data.edqs.ToCoreEdqsMsg;
@ -53,18 +57,22 @@ import org.thingsboard.server.edqs.processor.EdqsProducer;
import org.thingsboard.server.edqs.state.EdqsPartitionService;
import org.thingsboard.server.edqs.util.EdqsConverter;
import org.thingsboard.server.gen.transport.TransportProtos.EdqsEventMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ServiceInfo;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToEdqsCoreServiceMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToEdqsMsg;
import org.thingsboard.server.queue.discovery.DiscoveryService;
import org.thingsboard.server.queue.discovery.HashPartitionService;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
import org.thingsboard.server.queue.discovery.TopicService;
import org.thingsboard.server.queue.environment.DistributedLock;
import org.thingsboard.server.queue.environment.DistributedLockService;
import org.thingsboard.server.queue.provider.EdqsClientQueueFactory;
import org.thingsboard.server.queue.util.AfterStartUp;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
@Service
@ -80,25 +88,36 @@ public class DefaultEdqsService implements EdqsService {
private final DistributedLockService distributedLockService;
private final AttributesService attributesService;
private final EdqsPartitionService edqsPartitionService;
private final TopicService topicService;
private final TbServiceInfoProvider serviceInfoProvider;
private final DiscoveryService discoveryService;
@Autowired @Lazy
private TbClusterService clusterService;
@Autowired @Lazy
private HashPartitionService hashPartitionService;
@Value("${queue.edqs.api.auto_enable:true}")
private boolean autoEnableApi;
@Value("${queue.edqs.readiness_check_interval:60000}")
private int edqsReadinessCheckInterval;
private EdqsProducer eventsProducer;
private ExecutorService executor;
private ScheduledExecutorService scheduler;
private DistributedLock syncLock;
@Getter
private EdqsState state;
@PostConstruct
private void init() {
executor = ThingsBoardExecutors.newWorkStealingPool(12, getClass());
scheduler = ThingsBoardExecutors.newSingleThreadScheduledExecutor("edqs-check");
eventsProducer = EdqsProducer.builder()
.producer(queueFactory.createEdqsEventsProducer())
.partitionService(edqsPartitionService)
.build();
syncLock = distributedLockService.getLock("edqs_sync");
state = new EdqsState();
}
@AfterStartUp(order = AfterStartUp.REGULAR_SERVICE)
@ -106,6 +125,26 @@ public class DefaultEdqsService implements EdqsService {
if (!serviceInfoProvider.isService(ServiceType.TB_CORE)) {
return;
}
if (edqsApiService.isSupported()) {
scheduler.scheduleWithFixedDelay(() -> {
if (!hashPartitionService.isSystemPartitionMine(ServiceType.TB_CORE)) {
return;
}
List<ServiceInfo> servers = new ArrayList<>(discoveryService.getOtherServers());
servers.add(serviceInfoProvider.getServiceInfo());
List<ServiceInfo> readyEdqsServers = servers.stream()
.filter(serviceInfo -> serviceInfo.getServiceTypesList().contains(ServiceType.EDQS.name()))
.filter(ServiceInfo::getReady)
.toList();
boolean changed = state.setEdqsReady(!readyEdqsServers.isEmpty());
if (changed) {
broadcastEdqsReady(state.getEdqsReady());
}
}, 0, edqsReadinessCheckInterval, TimeUnit.MILLISECONDS);
}
executor.submit(() -> {
try {
EdqsSyncState syncState = getSyncState();
@ -115,9 +154,9 @@ public class DefaultEdqsService implements EdqsService {
.syncRequest(new EdqsSyncRequest())
.build());
}
} else if (edqsApiService.isSupported() && edqsApiService.isAutoEnable()) {
} else {
// only if topic/RocksDB is not empty and sync is finished
edqsApiService.setEnabled(true);
onSyncStatusUpdate(EdqsSyncStatus.FINISHED);
}
} catch (Throwable e) {
log.error("Failed to start EDQS service", e);
@ -131,7 +170,10 @@ public class DefaultEdqsService implements EdqsService {
if (request.getSyncRequest() != null) {
saveSyncState(EdqsSyncStatus.REQUESTED);
}
broadcast(request.toInternalMsg());
broadcast(ToCoreEdqsMsg.builder()
.syncRequest(request.getSyncRequest())
.apiEnabled(request.getApiEnabled())
.build());
}
@Override
@ -140,7 +182,13 @@ public class DefaultEdqsService implements EdqsService {
log.info("Processing system msg {}", msg);
try {
if (msg.getApiEnabled() != null) {
edqsApiService.setEnabled(msg.getApiEnabled());
state.setApiEnabled(msg.getApiEnabled());
}
if (msg.getEdqsReady() != null) {
onEdqsReady(msg.getEdqsReady());
}
if (msg.getSyncStatus() != null) {
onSyncStatusUpdate(msg.getSyncStatus());
}
if (msg.getSyncRequest() != null) {
@ -154,23 +202,16 @@ public class DefaultEdqsService implements EdqsService {
return;
}
}
saveSyncState(EdqsSyncStatus.STARTED);
edqsSyncService.sync();
saveSyncState(EdqsSyncStatus.FINISHED);
if (edqsApiService.isSupported())
if (edqsApiService.isAutoEnable()) {
log.info("EDQS sync is finished, auto-enabling API");
broadcast(ToCoreEdqsMsg.builder()
.apiEnabled(Boolean.TRUE)
.build());
} else {
log.info("EDQS sync is finished, but leaving API disabled");
}
saveSyncState(EdqsSyncStatus.FINISHED);
broadcastSyncStatusUpdate(EdqsSyncStatus.FINISHED);
} catch (Exception e) {
log.error("Failed to complete sync", e);
saveSyncState(EdqsSyncStatus.FAILED);
broadcastSyncStatusUpdate(EdqsSyncStatus.FAILED);
} finally {
syncLock.unlock();
}
@ -181,6 +222,60 @@ public class DefaultEdqsService implements EdqsService {
});
}
private void broadcastEdqsReady(boolean ready) {
broadcast(ToCoreEdqsMsg.builder()
.edqsReady(ready)
.build());
}
private void onEdqsReady(boolean ready) {
state.setEdqsReady(ready);
checkState();
}
private void broadcastSyncStatusUpdate(EdqsSyncStatus status) {
broadcast(ToCoreEdqsMsg.builder()
.syncStatus(status)
.build());
}
private void onSyncStatusUpdate(EdqsSyncStatus status) {
state.setSyncStatus(status);
checkState();
}
private void checkState() {
if (!edqsApiService.isSupported()) {
log.info("New state: {}. EDQS API not supported", state);
return;
}
if (state.isApiReady()) {
if (autoEnableApi) {
if (state.getApiEnabled() == null) {
state.setApiEnabled(true);
log.info("New state: {}. Auto-enabled EDQS API", state);
} else {
log.info("New state: {}. API mode left as is", state);
}
} else {
log.info("New state: {}. API auto-enabling is disabled", state);
}
} else {
if (state.isApiEnabled()) {
state.setApiEnabled(false);
log.info("New state: {}. Disabled EDQS API", state);
} else {
log.info("New state: {}. API left disabled", state);
}
}
}
@Override
public boolean isApiEnabled() {
return state.isApiEnabled();
}
@Override
public void onUpdate(TenantId tenantId, EntityId entityId, Object entity) {
EntityType entityType = entityId.getEntityType();
@ -278,6 +373,7 @@ public class DefaultEdqsService implements EdqsService {
@PreDestroy
private void stop() {
executor.shutdown();
scheduler.shutdownNow();
eventsProducer.stop();
}
@ -288,11 +384,4 @@ public class DefaultEdqsService implements EdqsService {
private EdqsSyncStatus status;
}
private enum EdqsSyncStatus {
REQUESTED,
STARTED,
FINISHED,
FAILED
}
}

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

@ -1755,6 +1755,8 @@ queue:
supported: "${TB_EDQS_API_SUPPORTED:false}"
# Whether to auto-enable EDQS API (if queue.edqs.api.supported is true) when sync of data to Kafka is finished
auto_enable: "${TB_EDQS_API_AUTO_ENABLE:true}"
# Interval in milliseconds to check for ready EDQS servers
readiness_check_interval: "${TB_EDQS_READINESS_CHECK_INTERVAL_MS:60000}"
# Mode of EDQS: local (for monolith) or remote (with separate EDQS microservices)
mode: "${TB_EDQS_MODE:local}"
local:

13
application/src/test/java/org/thingsboard/server/controller/EdqsEntityQueryControllerTest.java

@ -23,9 +23,8 @@ import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.query.EntityCountQuery;
import org.thingsboard.server.common.data.query.EntityData;
import org.thingsboard.server.common.data.query.EntityDataQuery;
import org.thingsboard.server.common.msg.edqs.EdqsApiService;
import org.thingsboard.server.common.msg.edqs.EdqsService;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.edqs.state.EdqsStateService;
import org.thingsboard.server.edqs.util.EdqsRocksDb;
import java.util.concurrent.TimeUnit;
@ -39,22 +38,20 @@ import static org.awaitility.Awaitility.await;
"queue.edqs.sync.enabled=true",
"queue.edqs.api.supported=true",
"queue.edqs.api.auto_enable=true",
"queue.edqs.mode=local"
"queue.edqs.mode=local",
"queue.edqs.readiness_check_interval=1000"
})
public class EdqsEntityQueryControllerTest extends EntityQueryControllerTest {
@Autowired
private EdqsApiService edqsApiService;
@Autowired
private EdqsStateService edqsStateService;
private EdqsService edqsService;
@MockBean // so that we don't do backup for tests
private EdqsRocksDb edqsRocksDb;
@Before
public void before() {
await().atMost(TIMEOUT, TimeUnit.SECONDS).until(() -> edqsApiService.isEnabled() && edqsStateService.isReady());
await().atMost(TIMEOUT, TimeUnit.SECONDS).until(() -> edqsService.getState().isApiEnabled());
}
@Override

7
application/src/test/java/org/thingsboard/server/controller/EntityQueryControllerTest.java

@ -22,6 +22,7 @@ import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.test.context.TestPropertySource;
import org.springframework.test.web.servlet.ResultActions;
import org.testcontainers.shaded.org.apache.commons.lang3.RandomStringUtils;
import org.thingsboard.common.util.JacksonUtil;
@ -78,6 +79,10 @@ import static org.awaitility.Awaitility.await;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
@DaoSqlTest
@TestPropertySource(properties = {
"queue.edqs.sync.enabled=true", // only enabling sync
"queue.edqs.api.supported=false",
})
public class EntityQueryControllerTest extends AbstractControllerTest {
private static final String CUSTOMER_USER_EMAIL = "entityQueryCustomer@thingsboard.org";
@ -803,7 +808,7 @@ public class EntityQueryControllerTest extends AbstractControllerTest {
//assign dashboard
doPost("/api/customer/" + savedCustomer.getId().getId().toString()
+ "/dashboard/" + savedDashboard.getId().getId().toString(), Dashboard.class);
+ "/dashboard/" + savedDashboard.getId().getId().toString(), Dashboard.class);
// check entity data query by customer
User customerUser = new User();

9
application/src/test/java/org/thingsboard/server/service/entitiy/EdqsEntityServiceTest.java

@ -33,7 +33,7 @@ import org.thingsboard.server.common.data.query.EntityKeyType;
import org.thingsboard.server.common.data.query.RelationsQueryFilter;
import org.thingsboard.server.common.data.relation.EntitySearchDirection;
import org.thingsboard.server.common.data.relation.RelationEntityTypeFilter;
import org.thingsboard.server.common.msg.edqs.EdqsApiService;
import org.thingsboard.server.common.msg.edqs.EdqsService;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.edqs.util.EdqsRocksDb;
@ -53,19 +53,20 @@ import static org.awaitility.Awaitility.await;
"queue.edqs.sync.enabled=true",
"queue.edqs.api.supported=true",
"queue.edqs.api.auto_enable=true",
"queue.edqs.mode=local"
"queue.edqs.mode=local",
"queue.edqs.readiness_check_interval=1000"
})
public class EdqsEntityServiceTest extends EntityServiceTest {
@Autowired
private EdqsApiService edqsApiService;
private EdqsService edqsService;
@MockBean
private EdqsRocksDb edqsRocksDb;
@Before
public void beforeEach() {
await().atMost(TIMEOUT, TimeUnit.SECONDS).until(() -> edqsApiService.isEnabled());
await().atMost(TIMEOUT, TimeUnit.SECONDS).until(() -> edqsService.isApiEnabled());
}
// sql implementation has a bug with data duplication, edqs implementation returns correct value

71
common/data/src/main/java/org/thingsboard/server/common/data/edqs/EdqsState.java

@ -0,0 +1,71 @@
/**
* Copyright © 2016-2025 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.data.edqs;
import lombok.Getter;
import lombok.NoArgsConstructor;
import org.apache.commons.lang3.BooleanUtils;
@Getter
@NoArgsConstructor
public class EdqsState {
private Boolean edqsReady;
private EdqsSyncStatus syncStatus;
private Boolean apiEnabled; // null until auto-enabled or set manually
public boolean setEdqsReady(boolean ready) {
boolean changed = BooleanUtils.toBooleanDefaultIfNull(this.edqsReady, false) != ready;
this.edqsReady = ready;
return changed;
}
public void setSyncStatus(EdqsSyncStatus syncStatus) {
this.syncStatus = syncStatus;
}
public boolean setApiEnabled(boolean apiEnabled) {
boolean changed = BooleanUtils.toBooleanDefaultIfNull(this.apiEnabled, false) != apiEnabled;
this.apiEnabled = apiEnabled;
return changed;
}
public boolean isApiReady() {
return edqsReady && syncStatus == EdqsSyncStatus.FINISHED;
}
public boolean isApiEnabled() {
return apiEnabled != null && apiEnabled;
}
@Override
public String toString() {
return '[' +
"EDQS ready: " + edqsReady +
", sync status: " + syncStatus +
", API enabled: " + apiEnabled +
']';
}
public enum EdqsSyncStatus {
REQUESTED,
STARTED,
FINISHED,
FAILED
}
}

4
common/data/src/main/java/org/thingsboard/server/common/data/edqs/ToCoreEdqsMsg.java

@ -19,6 +19,7 @@ import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.thingsboard.server.common.data.edqs.EdqsState.EdqsSyncStatus;
@Data
@AllArgsConstructor
@ -29,4 +30,7 @@ public class ToCoreEdqsMsg {
private EdqsSyncRequest syncRequest;
private Boolean apiEnabled;
private EdqsSyncStatus syncStatus;
private Boolean edqsReady;
}

6
common/data/src/main/java/org/thingsboard/server/common/data/edqs/ToCoreEdqsRequest.java

@ -15,7 +15,6 @@
*/
package org.thingsboard.server.common.data.edqs;
import com.fasterxml.jackson.annotation.JsonIgnore;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
@ -30,9 +29,4 @@ public class ToCoreEdqsRequest {
private EdqsSyncRequest syncRequest;
private Boolean apiEnabled;
@JsonIgnore
public ToCoreEdqsMsg toInternalMsg() {
return new ToCoreEdqsMsg(syncRequest, apiEnabled);
}
}

21
common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java

@ -34,7 +34,7 @@ import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager;
import org.thingsboard.server.queue.common.consumer.QueueConsumerManager;
import org.thingsboard.server.queue.common.state.KafkaQueueStateService;
import org.thingsboard.server.queue.common.state.QueueStateService;
import org.thingsboard.server.queue.discovery.DiscoveryService;
import org.thingsboard.server.queue.discovery.QueueKey;
import org.thingsboard.server.queue.edqs.EdqsConfig;
import org.thingsboard.server.queue.edqs.EdqsExecutors;
@ -61,20 +61,22 @@ public class KafkaEdqsStateService implements EdqsStateService {
private final EdqsConfig config;
private final EdqsPartitionService partitionService;
private final KafkaEdqsQueueFactory queueFactory;
private final DiscoveryService discoveryService;
private final EdqsExecutors edqsExecutors;
@Autowired
@Lazy
private EdqsProcessor edqsProcessor;
private PartitionedQueueConsumerManager<TbProtoQueueMsg<ToEdqsMsg>> stateConsumer;
private QueueStateService<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<ToEdqsMsg>> queueStateService;
private KafkaQueueStateService<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<ToEdqsMsg>> queueStateService;
private QueueConsumerManager<TbProtoQueueMsg<ToEdqsMsg>> eventsToBackupConsumer;
private EdqsProducer stateProducer;
private final VersionsStore versionsStore = new VersionsStore();
private final AtomicInteger stateReadCount = new AtomicInteger();
private final AtomicInteger eventsReadCount = new AtomicInteger();
private Boolean ready;
private boolean ready = false;
@Override
public void init(PartitionedQueueConsumerManager<TbProtoQueueMsg<ToEdqsMsg>> eventConsumer, List<PartitionedQueueConsumerManager<?>> otherConsumers) {
@ -187,7 +189,10 @@ public class KafkaEdqsStateService implements EdqsStateService {
eventsToBackupConsumer.subscribe(allPartitions);
eventsToBackupConsumer.launch();
}
queueStateService.update(new QueueKey(ServiceType.EDQS), partitions);
queueStateService.update(new QueueKey(ServiceType.EDQS), partitions, () -> {
ready = true;
discoveryService.setReady(true);
});
}
@Override
@ -197,13 +202,7 @@ public class KafkaEdqsStateService implements EdqsStateService {
@Override
public boolean isReady() {
if (ready == null) {
Set<TopicPartitionInfo> partitionsInProgress = queueStateService.getPartitionsInProgress();
if (partitionsInProgress != null && partitionsInProgress.isEmpty()) {
ready = true; // once true - always true, not to change readiness status on each repartitioning
}
}
return ready != null && ready;
return ready;
}
private TenantId getTenantId(ToEdqsMsg edqsMsg) {

6
common/message/src/main/java/org/thingsboard/server/common/msg/edqs/EdqsApiService.java

@ -25,12 +25,6 @@ public interface EdqsApiService {
ListenableFuture<EdqsResponse> processRequest(TenantId tenantId, CustomerId customerId, EdqsRequest request);
boolean isEnabled();
void setEnabled(boolean enabled);
boolean isSupported();
boolean isAutoEnable();
}

5
common/message/src/main/java/org/thingsboard/server/common/msg/edqs/EdqsService.java

@ -17,6 +17,7 @@ package org.thingsboard.server.common.msg.edqs;
import org.thingsboard.server.common.data.ObjectType;
import org.thingsboard.server.common.data.edqs.EdqsObject;
import org.thingsboard.server.common.data.edqs.EdqsState;
import org.thingsboard.server.common.data.edqs.ToCoreEdqsMsg;
import org.thingsboard.server.common.data.edqs.ToCoreEdqsRequest;
import org.thingsboard.server.common.data.id.EntityId;
@ -36,4 +37,8 @@ public interface EdqsService {
void processSystemMsg(ToCoreEdqsMsg request);
boolean isApiEnabled();
EdqsState getState();
}

1
common/proto/src/main/proto/queue.proto

@ -87,6 +87,7 @@ message ServiceInfo {
SystemInfoProto systemInfo = 10;
repeated string assignedTenantProfiles = 11;
string label = 12;
bool ready = 13;
}
message SystemInfoProto {

8
common/queue/src/main/java/org/thingsboard/server/queue/common/state/KafkaQueueStateService.java

@ -26,6 +26,7 @@ import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.Supplier;
import static org.thingsboard.server.common.msg.queue.TopicPartitionInfo.withTopic;
@ -36,6 +37,8 @@ public class KafkaQueueStateService<E extends TbQueueMsg, S extends TbQueueMsg>
private final PartitionedQueueConsumerManager<S> stateConsumer;
private final Supplier<Map<String, Long>> eventsStartOffsetsProvider;
private final Set<TopicPartitionInfo> partitionsInProgress = ConcurrentHashMap.newKeySet();
@Builder
public KafkaQueueStateService(PartitionedQueueConsumerManager<E> eventConsumer,
PartitionedQueueConsumerManager<S> stateConsumer,
@ -47,7 +50,7 @@ public class KafkaQueueStateService<E extends TbQueueMsg, S extends TbQueueMsg>
}
@Override
protected void addPartitions(QueueKey queueKey, Set<TopicPartitionInfo> partitions) {
protected void addPartitions(QueueKey queueKey, Set<TopicPartitionInfo> partitions, Runnable whenAllProcessed) {
Map<String, Long> eventsStartOffsets = eventsStartOffsetsProvider != null ? eventsStartOffsetsProvider.get() : null; // remembering the offsets before subscribing to states
Set<TopicPartitionInfo> statePartitions = withTopic(partitions, stateConsumer.getTopic());
@ -60,6 +63,9 @@ public class KafkaQueueStateService<E extends TbQueueMsg, S extends TbQueueMsg>
log.info("Finished partition {} (still in progress: {})", statePartition, partitionsInProgress);
if (partitionsInProgress.isEmpty()) {
log.info("All partitions processed");
if (whenAllProcessed != null) {
whenAllProcessed.run();
}
}
TopicPartitionInfo eventPartition = statePartition.withTopic(eventConsumer.getTopic());

21
common/queue/src/main/java/org/thingsboard/server/queue/common/state/QueueStateService.java

@ -28,7 +28,6 @@ import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
@ -42,8 +41,6 @@ public abstract class QueueStateService<E extends TbQueueMsg, S extends TbQueueM
@Getter
protected final Map<QueueKey, Set<TopicPartitionInfo>> partitions = new HashMap<>();
protected final Set<TopicPartitionInfo> partitionsInProgress = ConcurrentHashMap.newKeySet();
protected boolean initialized;
protected final ReadWriteLock partitionsLock = new ReentrantReadWriteLock();
@ -52,7 +49,7 @@ public abstract class QueueStateService<E extends TbQueueMsg, S extends TbQueueM
this.otherConsumers = otherConsumers;
}
public void update(QueueKey queueKey, Set<TopicPartitionInfo> newPartitions) {
public void update(QueueKey queueKey, Set<TopicPartitionInfo> newPartitions, Runnable whenAllProcessed) {
newPartitions = withTopic(newPartitions, eventConsumer.getTopic());
var writeLock = partitionsLock.writeLock();
writeLock.lock();
@ -74,12 +71,18 @@ public abstract class QueueStateService<E extends TbQueueMsg, S extends TbQueueM
}
if (!addedPartitions.isEmpty()) {
addPartitions(queueKey, addedPartitions);
addPartitions(queueKey, addedPartitions, whenAllProcessed);
} else {
if (whenAllProcessed != null) {
whenAllProcessed.run();
}
}
initialized = true;
}
protected void addPartitions(QueueKey queueKey, Set<TopicPartitionInfo> partitions) {
protected void addPartitions(QueueKey queueKey, Set<TopicPartitionInfo> partitions, Runnable whenAllProcessed) {
if (whenAllProcessed != null) {
whenAllProcessed.run();
}
eventConsumer.addPartitions(partitions);
for (PartitionedQueueConsumerManager<?> consumer : otherConsumers) {
consumer.addPartitions(withTopic(partitions, consumer.getTopic()));
@ -114,10 +117,6 @@ public abstract class QueueStateService<E extends TbQueueMsg, S extends TbQueueM
}
}
public Set<TopicPartitionInfo> getPartitionsInProgress() {
return initialized ? partitionsInProgress : null;
}
public void stop() {
eventConsumer.stop();
eventConsumer.awaitStop();

11
common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java

@ -68,6 +68,8 @@ public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider {
private List<ServiceType> serviceTypes;
private ServiceInfo serviceInfo;
private boolean ready = true;
@PostConstruct
public void init() {
if (StringUtils.isEmpty(serviceId)) {
@ -87,6 +89,7 @@ public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider {
assignedTenantProfiles = Collections.emptySet();
}
if (serviceTypes.contains(ServiceType.EDQS)) {
ready = false;
if (StringUtils.isBlank(edqsConfig.getLabel())) {
edqsConfig.setLabel(serviceId);
}
@ -128,9 +131,17 @@ public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider {
builder.addAllAssignedTenantProfiles(assignedTenantProfiles.stream().map(UUID::toString).collect(Collectors.toList()));
}
builder.setLabel(edqsConfig.getLabel());
builder.setReady(ready);
return serviceInfo = builder.build();
}
@Override
public boolean setReady(boolean ready) {
boolean changed = this.ready != ready;
this.ready = ready;
return changed;
}
private TransportProtos.SystemInfoProto getCurrentSystemInfoProto() {
TransportProtos.SystemInfoProto.Builder builder = TransportProtos.SystemInfoProto.newBuilder();

2
common/queue/src/main/java/org/thingsboard/server/queue/discovery/DiscoveryService.java

@ -25,4 +25,6 @@ public interface DiscoveryService {
boolean isMonolith();
void setReady(boolean ready);
}

9
common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java

@ -55,4 +55,13 @@ public class DummyDiscoveryService implements DiscoveryService {
public boolean isMonolith() {
return true;
}
@Override
public void setReady(boolean ready) {
boolean changed = serviceInfoProvider.setReady(ready);
if (changed) {
serviceInfoProvider.generateNewServiceInfoWithCurrentSystemInfo();
}
}
}

2
common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java

@ -35,4 +35,6 @@ public interface TbServiceInfoProvider {
Set<UUID> getAssignedTenantProfiles();
boolean setReady(boolean ready);
}

12
common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java

@ -177,6 +177,18 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
}
}
@Override
public void setReady(boolean ready) {
boolean changed = serviceInfoProvider.setReady(ready);
if (changed) {
try {
publishCurrentServer();
} catch (Exception e) {
log.error("Failed to update server readiness status", e);
}
}
}
private boolean currentServerExists() {
if (nodePath == null) {
return false;

10
dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java

@ -46,6 +46,7 @@ import org.thingsboard.server.common.data.query.EntityTypeFilter;
import org.thingsboard.server.common.data.query.KeyFilter;
import org.thingsboard.server.common.data.query.RelationsQueryFilter;
import org.thingsboard.server.common.msg.edqs.EdqsApiService;
import org.thingsboard.server.common.msg.edqs.EdqsService;
import org.thingsboard.server.common.stats.EdqsStatsService;
import org.thingsboard.server.dao.exception.IncorrectParameterException;
@ -86,6 +87,9 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe
@Lazy
EntityServiceRegistry entityServiceRegistry;
@Autowired
private EdqsService edqsService;
@Autowired
@Lazy
private EdqsApiService edqsApiService;
@ -102,7 +106,7 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe
long startNs = System.nanoTime();
Long result;
if (edqsApiService.isEnabled() && validForEdqs(query) && !tenantId.isSysTenantId()) {
if (edqsService.isApiEnabled() && validForEdqs(query) && !tenantId.isSysTenantId()) {
EdqsRequest request = EdqsRequest.builder()
.entityCountQuery(query)
.build();
@ -124,7 +128,7 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe
long startNs = System.nanoTime();
PageData<EntityData> result;
if (edqsApiService.isEnabled() && validForEdqs(query)) {
if (edqsService.isApiEnabled() && validForEdqs(query)) {
EdqsRequest request = EdqsRequest.builder()
.entityDataQuery(query)
.build();
@ -285,7 +289,7 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe
}
if ((query.getEntityFields() == null || query.getEntityFields().isEmpty()) &&
(query.getLatestValues() == null || query.getLatestValues().isEmpty())) {
(query.getLatestValues() == null || query.getLatestValues().isEmpty())) {
return false;
}

15
dao/src/main/java/org/thingsboard/server/dao/sql/query/DummyEdqsApiService.java

@ -35,24 +35,9 @@ public class DummyEdqsApiService implements EdqsApiService {
throw new UnsupportedOperationException();
}
@Override
public boolean isEnabled() {
return false;
}
@Override
public void setEnabled(boolean enabled) {
log.warn("Got request to enable EDQS API, but it isn't supported", new RuntimeException("stacktrace"));
}
@Override
public boolean isSupported() {
return false;
}
@Override
public boolean isAutoEnable() {
return false;
}
}

11
dao/src/main/java/org/thingsboard/server/dao/sql/query/DummyEdqsService.java

@ -19,6 +19,7 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.ObjectType;
import org.thingsboard.server.common.data.edqs.EdqsObject;
import org.thingsboard.server.common.data.edqs.EdqsState;
import org.thingsboard.server.common.data.edqs.ToCoreEdqsMsg;
import org.thingsboard.server.common.data.edqs.ToCoreEdqsRequest;
import org.thingsboard.server.common.data.id.EntityId;
@ -47,4 +48,14 @@ public class DummyEdqsService implements EdqsService {
@Override
public void processSystemMsg(ToCoreEdqsMsg request) {}
@Override
public boolean isApiEnabled() {
return getState().isApiEnabled();
}
@Override
public EdqsState getState() {
return new EdqsState();
}
}

Loading…
Cancel
Save