From 4e5abf98da5685314e1a34725ee08fb4a36a7a15 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Wed, 30 Oct 2024 17:26:31 +0200 Subject: [PATCH] Process edge-event consumer with correct polling strategy --- .../service/edge/EdgeContextComponent.java | 3 +- .../edge/EdgeEventSourcingListener.java | 31 ++++-- .../edge/rpc/AbstractEdgeGrpcSession.java | 3 +- .../service/edge/rpc/EdgeGrpcService.java | 66 ++++++------ .../edge/rpc/KafkaEdgeEventService.java | 3 - .../edge/rpc/KafkaEdgeGrpcSession.java | 101 +++++++++--------- .../edge/rpc/PostgresEdgeGrpcSession.java | 1 - .../src/main/resources/thingsboard.yml | 2 +- 8 files changed, 107 insertions(+), 103 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java index 79259c36b6..ee50bbb4a3 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java @@ -17,6 +17,7 @@ package org.thingsboard.server.service.edge; import lombok.Data; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Component; import org.thingsboard.server.cache.limits.RateLimitService; @@ -141,7 +142,7 @@ public class EdgeContextComponent { @Autowired private EdgeRequestsService edgeRequestsService; - @Autowired + @Autowired(required = false) private EdgeRpcService edgeRpcService; @Autowired diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java index 7706244680..ff5f62b655 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java @@ -53,8 +53,10 @@ import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; import org.thingsboard.server.dao.eventsourcing.RelationActionEvent; import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; import org.thingsboard.server.dao.tenant.TenantService; -import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.discovery.TopicService; +import org.thingsboard.server.queue.kafka.TbKafkaAdmin; +import org.thingsboard.server.queue.kafka.TbKafkaSettings; +import org.thingsboard.server.queue.kafka.TbKafkaTopicConfigs; /** * This event listener does not support async event processing because relay on ThreadLocal @@ -77,11 +79,14 @@ import org.thingsboard.server.queue.discovery.TopicService; public class EdgeEventSourcingListener { private final TopicService topicService; - private final TbQueueAdmin tbQueueAdmin; - private final TenantService tenantService; private final TbClusterService tbClusterService; + + private final TenantService tenantService; private final EdgeSynchronizationManager edgeSynchronizationManager; + private final TbKafkaSettings kafkaSettings; + private final TbKafkaTopicConfigs kafkaTopicConfigs; + @Value("#{'${queue.type:null}' == 'kafka'}") private boolean isKafkaSupported; @@ -117,17 +122,13 @@ public class EdgeEventSourcingListener { return; } try { - if (EntityType.EDGE.equals(entityType)) { - if (isKafkaSupported) { - String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, (EdgeId) event.getEntityId()).getTopic(); - tbQueueAdmin.deleteTopic(topic); - } else { - return; - } - } if (EntityType.TENANT.equals(entityType)) { return; } + if (EntityType.EDGE.equals(entityType)) { + handleEdgeEntityDeletion((EdgeId) event.getEntityId(), tenantId); + return; + } log.trace("[{}] DeleteEntityEvent called: {}", tenantId, event); EdgeEventType type = getEdgeEventTypeForEntityEvent(event.getEntity()); EdgeEventActionType actionType = getEdgeEventActionTypeForEntityEvent(event.getEntity()); @@ -139,6 +140,14 @@ public class EdgeEventSourcingListener { } } + private void handleEdgeEntityDeletion(EdgeId edgeId, TenantId tenantId) { + if (isKafkaSupported) { + String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); + TbKafkaAdmin kafkaAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdgeEventConfigs()); + kafkaAdmin.deleteTopic(topic); + } + } + private EdgeEventActionType getEdgeEventActionTypeForEntityEvent(Object entity) { if (entity instanceof AlarmComment) { return EdgeEventActionType.DELETED_COMMENT; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/AbstractEdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/AbstractEdgeGrpcSession.java index 397e1f6137..475826fee7 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/AbstractEdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/AbstractEdgeGrpcSession.java @@ -76,6 +76,7 @@ import org.thingsboard.server.gen.edge.v1.UserCredentialsRequestMsg; import org.thingsboard.server.gen.edge.v1.WidgetBundleTypesRequestMsg; import org.thingsboard.server.service.edge.EdgeContextComponent; import org.thingsboard.server.service.edge.rpc.fetch.EdgeEventFetcher; +import org.thingsboard.server.service.edge.rpc.fetch.GeneralEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.processor.alarm.AlarmProcessor; import org.thingsboard.server.service.edge.rpc.processor.asset.AssetProcessor; import org.thingsboard.server.service.edge.rpc.processor.asset.profile.AssetProfileProcessor; @@ -147,8 +148,6 @@ public abstract class AbstractEdgeGrpcSession processEdgeEvents() throws Exception; - public void initInputStream() { this.inputStream = new StreamObserver<>() { @Override diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java index bc63eb251a..bed4272ec1 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java @@ -89,7 +89,7 @@ import java.util.function.Consumer; @TbCoreComponent public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase implements EdgeRpcService { - private final ConcurrentMap sessions = new ConcurrentHashMap<>(); + private final ConcurrentMap> sessions = new ConcurrentHashMap<>(); private final ConcurrentMap sessionNewEventsLocks = new ConcurrentHashMap<>(); private final Map sessionNewEvents = new HashMap<>(); private final ConcurrentMap> sessionEdgeEventChecks = new ConcurrentHashMap<>(); @@ -212,22 +212,16 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Override public StreamObserver handleMsgs(StreamObserver outputStream) { - if (isKafkaSupported) { - return new KafkaEdgeGrpcSession(ctx, - outputStream, - this::onEdgeConnect, - this::onEdgeDisconnect, - sendDownlinkExecutorService, - this.maxInboundMessageSize, - this.maxHighPriorityQueueSizePerSession).getInputStream(); - } - return new PostgresEdgeGrpcSession(ctx, - outputStream, - this::onEdgeConnect, - this::onEdgeDisconnect, - sendDownlinkExecutorService, - this.maxInboundMessageSize, - this.maxHighPriorityQueueSizePerSession).getInputStream(); + AbstractEdgeGrpcSession session = createEdgeGrpcSession(outputStream); + return session.getInputStream(); + } + + private AbstractEdgeGrpcSession createEdgeGrpcSession(StreamObserver outputStream) { + return isKafkaSupported + ? new KafkaEdgeGrpcSession(ctx, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, + sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession) + : new PostgresEdgeGrpcSession(ctx, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, + sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession); } @Override @@ -258,7 +252,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Override public void updateEdge(TenantId tenantId, Edge edge) { - AbstractEdgeGrpcSession session = sessions.get(edge.getId()); + AbstractEdgeGrpcSession session = sessions.get(edge.getId()); if (session != null && session.isConnected()) { log.debug("[{}] Updating configuration for edge [{}] [{}]", tenantId, edge.getName(), edge.getId()); session.onConfigurationUpdate(edge); @@ -269,7 +263,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Override public void deleteEdge(TenantId tenantId, EdgeId edgeId) { - AbstractEdgeGrpcSession session = sessions.get(edgeId); + AbstractEdgeGrpcSession session = sessions.get(edgeId); if (session != null && session.isConnected()) { log.info("[{}] Closing and removing session for edge [{}]", tenantId, edgeId); session.close(); @@ -286,7 +280,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } private void onEdgeEventUpdate(TenantId tenantId, EdgeId edgeId) { - AbstractEdgeGrpcSession session = sessions.get(edgeId); + AbstractEdgeGrpcSession session = sessions.get(edgeId); if (session != null && session.isConnected()) { log.trace("[{}] onEdgeEventUpdate [{}]", tenantId, edgeId.getId()); updateSessionEventsFlag(tenantId, edgeId); @@ -297,7 +291,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i TenantId tenantId = msg.getTenantId(); EdgeEvent edgeEvent = msg.getEdgeEvent(); EdgeId edgeId = edgeEvent.getEdgeId(); - AbstractEdgeGrpcSession session = sessions.get(edgeId); + AbstractEdgeGrpcSession session = sessions.get(edgeId); if (session != null && session.isConnected()) { log.trace("[{}] onEdgeEvent [{}]", tenantId, edgeId); session.addEventToHighPriorityQueue(edgeEvent); @@ -318,7 +312,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } } - private void onEdgeConnect(EdgeId edgeId, AbstractEdgeGrpcSession edgeGrpcSession) { + private void onEdgeConnect(EdgeId edgeId, AbstractEdgeGrpcSession edgeGrpcSession) { Edge edge = edgeGrpcSession.getEdge(); TenantId tenantId = edge.getTenantId(); log.info("[{}][{}] edge [{}] connected successfully.", tenantId, edgeGrpcSession.getSessionId(), edgeId); @@ -334,17 +328,24 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i long lastConnectTs = System.currentTimeMillis(); save(tenantId, edgeId, DefaultDeviceStateService.LAST_CONNECT_TIME, lastConnectTs); edgeIdServiceIdCache.put(edgeId, serviceInfoProvider.getServiceId()); - if (isKafkaSupported) { - TbQueueConsumer> consumer = tbCoreQueueFactory.createEdgeEventMsgConsumer(tenantId, edgeId); - ((KafkaEdgeGrpcSession) edgeGrpcSession).initConsumer(consumer); + if (edgeGrpcSession instanceof KafkaEdgeGrpcSession) { + initializeKafkaConsumer((KafkaEdgeGrpcSession) edgeGrpcSession, tenantId, edgeId); } pushRuleEngineMessage(tenantId, edge, lastConnectTs, TbMsgType.CONNECT_EVENT); cancelScheduleEdgeEventsCheck(edgeId); - scheduleEdgeEventsCheck(edgeGrpcSession); + if (edgeGrpcSession instanceof PostgresEdgeGrpcSession) { + scheduleEdgeEventsCheck((PostgresEdgeGrpcSession) edgeGrpcSession); + } + } + + private void initializeKafkaConsumer(KafkaEdgeGrpcSession kafkaEdgeGrpcSession, TenantId tenantId, EdgeId edgeId) { + TbQueueConsumer> consumer = tbCoreQueueFactory.createEdgeEventMsgConsumer(tenantId, edgeId); + kafkaEdgeGrpcSession.initConsumer(() -> consumer, schedulerPoolSize); + kafkaEdgeGrpcSession.startConsumers(); } private void startSyncProcess(TenantId tenantId, EdgeId edgeId, UUID requestId, String requestServiceId) { - AbstractEdgeGrpcSession session = sessions.get(edgeId); + AbstractEdgeGrpcSession session = sessions.get(edgeId); if (session != null) { if (!session.isSyncCompleted()) { clusterService.pushEdgeSyncResponseToCore(new FromEdgeSyncResponse(requestId, tenantId, edgeId, false, "Sync process is active at the moment"), requestServiceId); @@ -364,7 +365,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i ToEdgeSyncRequest request = new ToEdgeSyncRequest(UUID.randomUUID(), tenantId, edgeId, serviceInfoProvider.getServiceId()); UUID requestId = request.getId(); - AbstractEdgeGrpcSession session = sessions.get(request.getEdgeId()); + AbstractEdgeGrpcSession session = sessions.get(request.getEdgeId()); if (session != null && !session.isSyncCompleted()) { responseConsumer.accept(new FromEdgeSyncResponse(requestId, request.getTenantId(), request.getEdgeId(), false, "Sync process is active at the moment")); } else { @@ -398,7 +399,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } } - private void scheduleEdgeEventsCheck(AbstractEdgeGrpcSession session) { + private void scheduleEdgeEventsCheck(PostgresEdgeGrpcSession session) { EdgeId edgeId = session.getEdge().getId(); UUID tenantId = session.getEdge().getTenantId().getId(); if (sessions.containsKey(edgeId)) { @@ -410,7 +411,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i if (Boolean.TRUE.equals(sessionNewEvents.get(edgeId))) { log.trace("[{}][{}] Set session new events flag to false", tenantId, edgeId.getId()); sessionNewEvents.put(edgeId, false); - Futures.addCallback(session.processEdgeEvents(), new FutureCallback() { + Futures.addCallback(session.processEdgeEvents(), new FutureCallback<>() { @Override public void onSuccess(Boolean newEventsAdded) { if (Boolean.TRUE.equals(newEventsAdded)) { @@ -418,6 +419,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } scheduleEdgeEventsCheck(session); } + @Override public void onFailure(Throwable t) { log.warn("[{}] Failed to process edge events for edge [{}]!", tenantId, session.getEdge().getId().getId(), t); @@ -456,7 +458,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i private void onEdgeDisconnect(Edge edge, UUID sessionId) { EdgeId edgeId = edge.getId(); log.info("[{}][{}] edge disconnected!", edgeId, sessionId); - AbstractEdgeGrpcSession toRemove = sessions.get(edgeId); + AbstractEdgeGrpcSession toRemove = sessions.get(edgeId); if (toRemove.getSessionId().equals(sessionId)) { toRemove = sessions.remove(edgeId); final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); @@ -506,6 +508,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } private static class AttributeSaveCallback implements FutureCallback { + private final TenantId tenantId; private final EdgeId edgeId; private final String key; @@ -527,6 +530,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i public void onFailure(Throwable t) { log.warn("[{}][{}] Failed to update attribute [{}] with value [{}]", tenantId, edgeId, key, value, t); } + } private void pushRuleEngineMessage(TenantId tenantId, Edge edge, long ts, TbMsgType msgType) { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java index 3caea5476a..13255dc4dc 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java @@ -35,7 +35,6 @@ import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.common.msg.tools.TbRateLimitsException; import org.thingsboard.server.common.util.ProtoUtils; import org.thingsboard.server.dao.edge.EdgeEventService; -import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeEventNotificationMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg; @@ -51,7 +50,6 @@ import java.util.UUID; public class KafkaEdgeEventService implements EdgeEventService { private final RateLimitService rateLimitService; - private final ApplicationEventPublisher eventPublisher; private final DataValidator edgeEventValidator; @Lazy private final TbQueueProducerProvider producerProvider; @@ -71,7 +69,6 @@ public class KafkaEdgeEventService implements EdgeEventService { ToEdgeEventNotificationMsg msg = ToEdgeEventNotificationMsg.newBuilder().setEdgeEventMsg(ProtoUtils.toProto(edgeEvent)).build(); producerProvider.getTbEdgeEventsMsgProducer().send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), msg), null); - eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(edgeEvent.getTenantId()).entity(edgeEvent).entityId(edgeEvent.getEdgeId()).build()); return Futures.immediateFuture(null); } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java index b7857d4e7d..e60523397f 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java @@ -17,33 +17,38 @@ package org.thingsboard.server.service.edge.rpc; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.ListenableFuture; import io.grpc.stub.StreamObserver; import jakarta.annotation.PreDestroy; import lombok.extern.slf4j.Slf4j; import org.checkerframework.checker.nullness.qual.Nullable; +import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.util.ProtoUtils; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.ResponseMsg; -import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeEventNotificationMsg; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; +import org.thingsboard.server.queue.common.consumer.QueueConsumerManager; import org.thingsboard.server.service.edge.EdgeContextComponent; import java.util.ArrayList; import java.util.List; import java.util.UUID; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.function.BiConsumer; +import java.util.function.Supplier; @Slf4j public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession { - private TbQueueConsumer> edgeEventsConsumer; + private ExecutorService consumerExecutor; + + private QueueConsumerManager> consumer; public KafkaEdgeGrpcSession(EdgeContextComponent ctx, StreamObserver outputStream, BiConsumer sessionOpenListener, @@ -53,75 +58,65 @@ public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession processEdgeEvents() { + protected void initConsumer(Supplier>> edgeEventsConsumer, long schedulerPoolSize) { + this.consumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("edge-event-consumer")); + this.consumer = QueueConsumerManager.>builder() + .name("TB Edge events") + .msgPackProcessor(this::processMsgs) + .pollInterval(schedulerPoolSize) + .consumerCreator(edgeEventsConsumer) + .consumerExecutor(consumerExecutor) + .threadPrefix("edge-events") + .build(); + } + + private void processMsgs(List> msgs, TbQueueConsumer> consumer) { log.trace("[{}][{}] starting processing edge events", this.tenantId, this.sessionId); if (isConnected() && isSyncCompleted()) { if (!highPriorityQueue.isEmpty()) { processHighPriorityEvents(); } else { - return processKafkaEvents(); + List edgeEvents = new ArrayList<>(); + for (TbProtoQueueMsg msg : msgs) { + EdgeEvent edgeEvent = ProtoUtils.fromProto(msg.getValue().getEdgeEventMsg()); + edgeEvents.add(edgeEvent); + } + List downlinkMsgsPack = convertToDownlinkMsgsPack(edgeEvents); + Futures.addCallback(sendDownlinkMsgsPack(downlinkMsgsPack), new FutureCallback<>() { + @Override + public void onSuccess(@Nullable Boolean isInterrupted) { + if (Boolean.TRUE.equals(isInterrupted)) { + log.debug("[{}][{}][{}] Send downlink messages task was interrupted", tenantId, edge.getId(), sessionId); + } else { + consumer.commit(); + } + } + @Override + public void onFailure(Throwable t) { + log.error("[{}] Failed to send downlink msgs pack", sessionId, t); + } + }, ctx.getGrpcCallbackExecutorService()); } } else { log.trace("[{}][{}] edge is not connected or sync is not completed. Skipping iteration", tenantId, sessionId); } - return null; } - private ListenableFuture processKafkaEvents() { - List edgeEvents = new ArrayList<>(); - try { - edgeEventsConsumer.subscribe(); - List> messages = edgeEventsConsumer.poll(100); - - if (messages.isEmpty()) { - return Futures.immediateFuture(Boolean.FALSE); - } - - for (TbProtoQueueMsg msg : messages) { - EdgeEvent edgeEvent = ProtoUtils.fromProto(msg.getValue().getEdgeEventMsg()); - edgeEvents.add(edgeEvent); - } - - List downlinkMsgsPack = convertToDownlinkMsgsPack(edgeEvents); - Futures.addCallback(sendDownlinkMsgsPack(downlinkMsgsPack), new FutureCallback<>() { - @Override - public void onSuccess(@Nullable Boolean isInterrupted) { - if (Boolean.TRUE.equals(isInterrupted)) { - log.debug("[{}][{}][{}] Send downlink messages task was interrupted", tenantId, edge.getId(), sessionId); - } else { - edgeEventsConsumer.commit(); - processEdgeEvents(); - } - } - @Override - public void onFailure(Throwable t) { - log.error("[{}] Failed to send downlink msgs pack", sessionId, t); - } - }, ctx.getGrpcCallbackExecutorService()); - return Futures.immediateFuture(Boolean.TRUE); - } catch (Exception e) { - log.error("[{}][{}] Error occurred while polling edge events from Kafka: {}", tenantId, edge.getId(), e.getMessage()); - return Futures.immediateFailedFuture(e); - } - } - - protected void initConsumer(TbQueueConsumer> edgeEventsConsumer) { - this.edgeEventsConsumer = edgeEventsConsumer; + public void startConsumers() { + consumer.subscribe(); + consumer.launch(); } @PreDestroy private void destroy() { - if (edgeEventsConsumer != null) { - edgeEventsConsumer.unsubscribe(); - } + consumer.stop(); + consumerExecutor.shutdown(); } public void stopConsumer() { - if (edgeEventsConsumer != null) { - edgeEventsConsumer.stop(); - } + consumer.stop(); + consumerExecutor.shutdown(); } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java index 8293942523..4d0fda21f6 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java @@ -63,7 +63,6 @@ public class PostgresEdgeGrpcSession extends AbstractEdgeGrpcSession processEdgeEvents() throws Exception { SettableFuture result = SettableFuture.create(); log.trace("[{}][{}] starting processing edge events", this.tenantId, this.sessionId); diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 180a978744..037957d1c4 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1439,7 +1439,7 @@ swagger: # Queue configuration parameters queue: - type: "${TB_QUEUE_TYPE:kafka}" # in-memory or kafka (Apache Kafka) or aws-sqs (AWS SQS) or pubsub (PubSub) or service-bus (Azure Service Bus) or rabbitmq (RabbitMQ) + type: "${TB_QUEUE_TYPE:in-memory}" # in-memory or kafka (Apache Kafka) or aws-sqs (AWS SQS) or pubsub (PubSub) or service-bus (Azure Service Bus) or rabbitmq (RabbitMQ) prefix: "${TB_QUEUE_PREFIX:}" # Global queue prefix. If specified, prefix is added before default topic name: 'prefix.default_topic_name'. Prefix is applied to all topics (and consumer groups for kafka). in_memory: stats: