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 15a87ec3f1..e9a33520ad 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 @@ -44,7 +44,7 @@ import java.util.function.BiConsumer; @Slf4j public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession { - private final TbQueueConsumer> edgeEventsConsumer; + private final TbCoreQueueFactory tbCoreQueueFactory; private volatile boolean isHighPriorityProcessing; @@ -57,7 +57,7 @@ public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession { BiConsumer sessionCloseListener, ScheduledExecutorService sendDownlinkExecutorService, int maxInboundMessageSize, int maxHighPriorityQueueSizePerSession) { super(ctx, outputStream, sessionOpenListener, sessionCloseListener, sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession); - edgeEventsConsumer = tbCoreQueueFactory.createEdgeEventMsgConsumer(tenantId, edge.getId()); + this.tbCoreQueueFactory = tbCoreQueueFactory; } private void processMsgs(List> msgs, TbQueueConsumer> consumer) { @@ -102,7 +102,7 @@ public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession { .name("TB Edge events") .msgPackProcessor(this::processMsgs) .pollInterval(ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval()) - .consumerCreator(() -> edgeEventsConsumer) + .consumerCreator(() -> tbCoreQueueFactory.createEdgeEventMsgConsumer(tenantId, edge.getId())) .consumerExecutor(consumerExecutor) .threadPrefix("edge-events") .build();