Browse Source

Process edge-event consumer with correct polling strategy

pull/11924/head
Andrii Landiak 2 years ago
parent
commit
4e5abf98da
  1. 3
      application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java
  2. 31
      application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java
  3. 3
      application/src/main/java/org/thingsboard/server/service/edge/rpc/AbstractEdgeGrpcSession.java
  4. 66
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java
  5. 3
      application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java
  6. 101
      application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java
  7. 1
      application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java
  8. 2
      application/src/main/resources/thingsboard.yml

3
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

31
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;

3
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<T extends AbstractEdgeGrpcSession<
initInputStream();
}
protected abstract ListenableFuture<Boolean> processEdgeEvents() throws Exception;
public void initInputStream() {
this.inputStream = new StreamObserver<>() {
@Override

66
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<EdgeId, AbstractEdgeGrpcSession> sessions = new ConcurrentHashMap<>();
private final ConcurrentMap<EdgeId, AbstractEdgeGrpcSession<?>> sessions = new ConcurrentHashMap<>();
private final ConcurrentMap<EdgeId, Lock> sessionNewEventsLocks = new ConcurrentHashMap<>();
private final Map<EdgeId, Boolean> sessionNewEvents = new HashMap<>();
private final ConcurrentMap<EdgeId, ScheduledFuture<?>> sessionEdgeEventChecks = new ConcurrentHashMap<>();
@ -212,22 +212,16 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
@Override
public StreamObserver<RequestMsg> handleMsgs(StreamObserver<ResponseMsg> 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<ResponseMsg> 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<TbProtoQueueMsg<ToEdgeEventNotificationMsg>> 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<TbProtoQueueMsg<ToEdgeEventNotificationMsg>> 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<Boolean>() {
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<Void> {
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) {

3
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<EdgeEvent> 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);
}

101
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<KafkaEdgeGrpcSession> {
private TbQueueConsumer<TbProtoQueueMsg<ToEdgeEventNotificationMsg>> edgeEventsConsumer;
private ExecutorService consumerExecutor;
private QueueConsumerManager<TbProtoQueueMsg<ToEdgeEventNotificationMsg>> consumer;
public KafkaEdgeGrpcSession(EdgeContextComponent ctx, StreamObserver<ResponseMsg> outputStream,
BiConsumer<EdgeId, KafkaEdgeGrpcSession> sessionOpenListener,
@ -53,75 +58,65 @@ public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession<KafkaEdgeGrpcS
super(ctx, outputStream, sessionOpenListener, sessionCloseListener, sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession);
}
@Override
protected ListenableFuture<Boolean> processEdgeEvents() {
protected void initConsumer(Supplier<TbQueueConsumer<TbProtoQueueMsg<ToEdgeEventNotificationMsg>>> edgeEventsConsumer, long schedulerPoolSize) {
this.consumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("edge-event-consumer"));
this.consumer = QueueConsumerManager.<TbProtoQueueMsg<ToEdgeEventNotificationMsg>>builder()
.name("TB Edge events")
.msgPackProcessor(this::processMsgs)
.pollInterval(schedulerPoolSize)
.consumerCreator(edgeEventsConsumer)
.consumerExecutor(consumerExecutor)
.threadPrefix("edge-events")
.build();
}
private void processMsgs(List<TbProtoQueueMsg<ToEdgeEventNotificationMsg>> msgs, TbQueueConsumer<TbProtoQueueMsg<ToEdgeEventNotificationMsg>> consumer) {
log.trace("[{}][{}] starting processing edge events", this.tenantId, this.sessionId);
if (isConnected() && isSyncCompleted()) {
if (!highPriorityQueue.isEmpty()) {
processHighPriorityEvents();
} else {
return processKafkaEvents();
List<EdgeEvent> edgeEvents = new ArrayList<>();
for (TbProtoQueueMsg<ToEdgeEventNotificationMsg> msg : msgs) {
EdgeEvent edgeEvent = ProtoUtils.fromProto(msg.getValue().getEdgeEventMsg());
edgeEvents.add(edgeEvent);
}
List<DownlinkMsg> 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<Boolean> processKafkaEvents() {
List<EdgeEvent> edgeEvents = new ArrayList<>();
try {
edgeEventsConsumer.subscribe();
List<TbProtoQueueMsg<ToEdgeEventNotificationMsg>> messages = edgeEventsConsumer.poll(100);
if (messages.isEmpty()) {
return Futures.immediateFuture(Boolean.FALSE);
}
for (TbProtoQueueMsg<ToEdgeEventNotificationMsg> msg : messages) {
EdgeEvent edgeEvent = ProtoUtils.fromProto(msg.getValue().getEdgeEventMsg());
edgeEvents.add(edgeEvent);
}
List<DownlinkMsg> 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<TbProtoQueueMsg<TransportProtos.ToEdgeEventNotificationMsg>> 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();
}
}

1
application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java

@ -63,7 +63,6 @@ public class PostgresEdgeGrpcSession extends AbstractEdgeGrpcSession<PostgresEdg
initInputStream();
}
@Override
protected ListenableFuture<Boolean> processEdgeEvents() throws Exception {
SettableFuture<Boolean> result = SettableFuture.create();
log.trace("[{}][{}] starting processing edge events", this.tenantId, this.sessionId);

2
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:

Loading…
Cancel
Save