@ -39,7 +39,7 @@ import org.thingsboard.server.queue.TbQueueRequestTemplate;
import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate ;
import org.thingsboard.server.queue.common.TbProtoJsQueueMsg ;
import org.thingsboard.server.queue.common.TbProtoQueueMsg ;
import org.thingsboard.server.queue.discovery.Notifications TopicService ;
import org.thingsboard.server.queue.discovery.TopicService ;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider ;
import org.thingsboard.server.queue.kafka.TbKafkaAdmin ;
import org.thingsboard.server.queue.kafka.TbKafkaConsumerStatsService ;
@ -62,7 +62,7 @@ import java.util.concurrent.atomic.AtomicLong;
@ConditionalOnExpression ( "'${queue.type:null}'=='kafka' && '${service.type:null}'=='monolith'" )
public class KafkaMonolithQueueFactory implements TbCoreQueueFactory , TbRuleEngineQueueFactory , TbVersionControlQueueFactory {
private final Notifications TopicService no tificationsT opicService;
private final TopicService topicService ;
private final TbKafkaSettings kafkaSettings ;
private final TbServiceInfoProvider serviceInfoProvider ;
private final TbQueueCoreSettings coreSettings ;
@ -85,7 +85,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
private final AtomicLong consumerCount = new AtomicLong ( ) ;
public KafkaMonolithQueueFactory ( Notifications TopicService no tificationsT opicService, TbKafkaSettings kafkaSettings ,
public KafkaMonolithQueueFactory ( TopicService topicService , TbKafkaSettings kafkaSettings ,
TbServiceInfoProvider serviceInfoProvider ,
TbQueueCoreSettings coreSettings ,
TbQueueRuleEngineSettings ruleEngineSettings ,
@ -95,7 +95,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
TbQueueVersionControlSettings vcSettings ,
TbKafkaConsumerStatsService consumerStatsService ,
TbKafkaTopicConfigs kafkaTopicConfigs ) {
this . no tificationsT opicService = no tificationsT opicService;
this . topicService = topicService ;
this . kafkaSettings = kafkaSettings ;
this . serviceInfoProvider = serviceInfoProvider ;
this . coreSettings = coreSettings ;
@ -122,7 +122,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
TbKafkaProducerTemplate . TbKafkaProducerTemplateBuilder < TbProtoQueueMsg < ToTransportMsg > > requestBuilder = TbKafkaProducerTemplate . builder ( ) ;
requestBuilder . settings ( kafkaSettings ) ;
requestBuilder . clientId ( "monolith-transport-notifications-" + serviceInfoProvider . getServiceId ( ) ) ;
requestBuilder . defaultTopic ( transportNotificationSettings . getNotificationsTopic ( ) ) ;
requestBuilder . defaultTopic ( topicService . buildTopicName ( t ransportNotificationSettings . getNotificationsTopic ( ) ) ) ;
requestBuilder . admin ( notificationAdmin ) ;
return requestBuilder . build ( ) ;
}
@ -132,7 +132,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
TbKafkaProducerTemplate . TbKafkaProducerTemplateBuilder < TbProtoQueueMsg < ToRuleEngineMsg > > requestBuilder = TbKafkaProducerTemplate . builder ( ) ;
requestBuilder . settings ( kafkaSettings ) ;
requestBuilder . clientId ( "monolith-rule-engine-" + serviceInfoProvider . getServiceId ( ) ) ;
requestBuilder . defaultTopic ( ruleEngineSettings . getTopic ( ) ) ;
requestBuilder . defaultTopic ( topicService . buildTopicName ( ruleEngineSettings . getTopic ( ) ) ) ;
requestBuilder . admin ( ruleEngineAdmin ) ;
return requestBuilder . build ( ) ;
}
@ -142,7 +142,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
TbKafkaProducerTemplate . TbKafkaProducerTemplateBuilder < TbProtoQueueMsg < ToRuleEngineNotificationMsg > > requestBuilder = TbKafkaProducerTemplate . builder ( ) ;
requestBuilder . settings ( kafkaSettings ) ;
requestBuilder . clientId ( "monolith-rule-engine-notifications-" + serviceInfoProvider . getServiceId ( ) ) ;
requestBuilder . defaultTopic ( ruleEngineSettings . getTopic ( ) ) ;
requestBuilder . defaultTopic ( topicService . buildTopicName ( ruleEngineSettings . getTopic ( ) ) ) ;
requestBuilder . admin ( notificationAdmin ) ;
return requestBuilder . build ( ) ;
}
@ -152,7 +152,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
TbKafkaProducerTemplate . TbKafkaProducerTemplateBuilder < TbProtoQueueMsg < ToCoreMsg > > requestBuilder = TbKafkaProducerTemplate . builder ( ) ;
requestBuilder . settings ( kafkaSettings ) ;
requestBuilder . clientId ( "monolith-core-" + serviceInfoProvider . getServiceId ( ) ) ;
requestBuilder . defaultTopic ( coreSettings . getTopic ( ) ) ;
requestBuilder . defaultTopic ( topicService . buildTopicName ( coreSettings . getTopic ( ) ) ) ;
requestBuilder . admin ( coreAdmin ) ;
return requestBuilder . build ( ) ;
}
@ -162,7 +162,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
TbKafkaProducerTemplate . TbKafkaProducerTemplateBuilder < TbProtoQueueMsg < ToCoreNotificationMsg > > requestBuilder = TbKafkaProducerTemplate . builder ( ) ;
requestBuilder . settings ( kafkaSettings ) ;
requestBuilder . clientId ( "monolith-core-notifications-" + serviceInfoProvider . getServiceId ( ) ) ;
requestBuilder . defaultTopic ( coreSettings . getTopic ( ) ) ;
requestBuilder . defaultTopic ( topicService . buildTopicName ( coreSettings . getTopic ( ) ) ) ;
requestBuilder . admin ( notificationAdmin ) ;
return requestBuilder . build ( ) ;
}
@ -171,9 +171,9 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
public TbQueueConsumer < TbProtoQueueMsg < TransportProtos . ToVersionControlServiceMsg > > createToVersionControlMsgConsumer ( ) {
TbKafkaConsumerTemplate . TbKafkaConsumerTemplateBuilder < TbProtoQueueMsg < TransportProtos . ToVersionControlServiceMsg > > consumerBuilder = TbKafkaConsumerTemplate . builder ( ) ;
consumerBuilder . settings ( kafkaSettings ) ;
consumerBuilder . topic ( vcSettings . getTopic ( ) ) ;
consumerBuilder . topic ( topicService . buildTopicName ( vcSettings . getTopic ( ) ) ) ;
consumerBuilder . clientId ( "monolith-vc-consumer-" + serviceInfoProvider . getServiceId ( ) ) ;
consumerBuilder . groupId ( "monolith-vc-node" ) ;
consumerBuilder . groupId ( topicService . buildTopicName ( "monolith-vc-node" ) ) ;
consumerBuilder . decoder ( msg - > new TbProtoQueueMsg < > ( msg . getKey ( ) , TransportProtos . ToVersionControlServiceMsg . parseFrom ( msg . getData ( ) ) , msg . getHeaders ( ) ) ) ;
consumerBuilder . admin ( vcAdmin ) ;
consumerBuilder . statsService ( consumerStatsService ) ;
@ -185,9 +185,9 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
String queueName = configuration . getName ( ) ;
TbKafkaConsumerTemplate . TbKafkaConsumerTemplateBuilder < TbProtoQueueMsg < ToRuleEngineMsg > > consumerBuilder = TbKafkaConsumerTemplate . builder ( ) ;
consumerBuilder . settings ( kafkaSettings ) ;
consumerBuilder . topic ( configuration . getTopic ( ) ) ;
consumerBuilder . topic ( topicService . buildTopicName ( configuration . getTopic ( ) ) ) ;
consumerBuilder . clientId ( "re-" + queueName + "-consumer-" + serviceInfoProvider . getServiceId ( ) + "-" + consumerCount . incrementAndGet ( ) ) ;
consumerBuilder . groupId ( "re-" + queueName + ( configuration . getTenantId ( ) . isSysTenantId ( ) ? "" : ( "-" + configuration . getTenantId ( ) ) ) + "-consumer" ) ;
consumerBuilder . groupId ( topicService . buildTopicName ( "re-" + queueName + ( configuration . getTenantId ( ) . isSysTenantId ( ) ? "" : ( "-" + configuration . getTenantId ( ) ) ) + "-consumer" ) ) ;
consumerBuilder . decoder ( msg - > new TbProtoQueueMsg < > ( msg . getKey ( ) , ToRuleEngineMsg . parseFrom ( msg . getData ( ) ) , msg . getHeaders ( ) ) ) ;
consumerBuilder . admin ( ruleEngineAdmin ) ;
consumerBuilder . statsService ( consumerStatsService ) ;
@ -198,9 +198,9 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
public TbQueueConsumer < TbProtoQueueMsg < ToRuleEngineNotificationMsg > > createToRuleEngineNotificationsMsgConsumer ( ) {
TbKafkaConsumerTemplate . TbKafkaConsumerTemplateBuilder < TbProtoQueueMsg < ToRuleEngineNotificationMsg > > consumerBuilder = TbKafkaConsumerTemplate . builder ( ) ;
consumerBuilder . settings ( kafkaSettings ) ;
consumerBuilder . topic ( notificationsT opicService. getNotificationsTopic ( ServiceType . TB_RULE_ENGINE , serviceInfoProvider . getServiceId ( ) ) . getFullTopicName ( ) ) ;
consumerBuilder . topic ( topicService . buildTopicName ( t opicService. getNotificationsTopic ( ServiceType . TB_RULE_ENGINE , serviceInfoProvider . getServiceId ( ) ) . getFullTopicName ( ) ) ) ;
consumerBuilder . clientId ( "monolith-rule-engine-notifications-consumer-" + serviceInfoProvider . getServiceId ( ) ) ;
consumerBuilder . groupId ( "monolith-rule-engine-notifications-consumer-" + serviceInfoProvider . getServiceId ( ) ) ;
consumerBuilder . groupId ( topicService . buildTopicName ( "monolith-rule-engine-notifications-consumer-" + serviceInfoProvider . getServiceId ( ) ) ) ;
consumerBuilder . decoder ( msg - > new TbProtoQueueMsg < > ( msg . getKey ( ) , ToRuleEngineNotificationMsg . parseFrom ( msg . getData ( ) ) , msg . getHeaders ( ) ) ) ;
consumerBuilder . admin ( notificationAdmin ) ;
consumerBuilder . statsService ( consumerStatsService ) ;
@ -211,9 +211,9 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
public TbQueueConsumer < TbProtoQueueMsg < ToCoreMsg > > createToCoreMsgConsumer ( ) {
TbKafkaConsumerTemplate . TbKafkaConsumerTemplateBuilder < TbProtoQueueMsg < ToCoreMsg > > consumerBuilder = TbKafkaConsumerTemplate . builder ( ) ;
consumerBuilder . settings ( kafkaSettings ) ;
consumerBuilder . topic ( coreSettings . getTopic ( ) ) ;
consumerBuilder . topic ( topicService . buildTopicName ( coreSettings . getTopic ( ) ) ) ;
consumerBuilder . clientId ( "monolith-core-consumer-" + serviceInfoProvider . getServiceId ( ) ) ;
consumerBuilder . groupId ( "monolith-core-consumer" ) ;
consumerBuilder . groupId ( topicService . buildTopicName ( "monolith-core-consumer" ) ) ;
consumerBuilder . decoder ( msg - > new TbProtoQueueMsg < > ( msg . getKey ( ) , ToCoreMsg . parseFrom ( msg . getData ( ) ) , msg . getHeaders ( ) ) ) ;
consumerBuilder . admin ( coreAdmin ) ;
consumerBuilder . statsService ( consumerStatsService ) ;
@ -224,9 +224,9 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
public TbQueueConsumer < TbProtoQueueMsg < ToCoreNotificationMsg > > createToCoreNotificationsMsgConsumer ( ) {
TbKafkaConsumerTemplate . TbKafkaConsumerTemplateBuilder < TbProtoQueueMsg < ToCoreNotificationMsg > > consumerBuilder = TbKafkaConsumerTemplate . builder ( ) ;
consumerBuilder . settings ( kafkaSettings ) ;
consumerBuilder . topic ( notificationsT opicService. getNotificationsTopic ( ServiceType . TB_CORE , serviceInfoProvider . getServiceId ( ) ) . getFullTopicName ( ) ) ;
consumerBuilder . topic ( topicService . buildTopicName ( t opicService. getNotificationsTopic ( ServiceType . TB_CORE , serviceInfoProvider . getServiceId ( ) ) . getFullTopicName ( ) ) ) ;
consumerBuilder . clientId ( "monolith-core-notifications-consumer-" + serviceInfoProvider . getServiceId ( ) ) ;
consumerBuilder . groupId ( "monolith-core-notifications-consumer-" + serviceInfoProvider . getServiceId ( ) ) ;
consumerBuilder . groupId ( topicService . buildTopicName ( "monolith-core-notifications-consumer-" + serviceInfoProvider . getServiceId ( ) ) ) ;
consumerBuilder . decoder ( msg - > new TbProtoQueueMsg < > ( msg . getKey ( ) , ToCoreNotificationMsg . parseFrom ( msg . getData ( ) ) , msg . getHeaders ( ) ) ) ;
consumerBuilder . admin ( notificationAdmin ) ;
consumerBuilder . statsService ( consumerStatsService ) ;
@ -237,9 +237,9 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
public TbQueueConsumer < TbProtoQueueMsg < TransportApiRequestMsg > > createTransportApiRequestConsumer ( ) {
TbKafkaConsumerTemplate . TbKafkaConsumerTemplateBuilder < TbProtoQueueMsg < TransportApiRequestMsg > > consumerBuilder = TbKafkaConsumerTemplate . builder ( ) ;
consumerBuilder . settings ( kafkaSettings ) ;
consumerBuilder . topic ( transportApiSettings . getRequestsTopic ( ) ) ;
consumerBuilder . topic ( topicService . buildTopicName ( t ransportApiSettings . getRequestsTopic ( ) ) ) ;
consumerBuilder . clientId ( "monolith-transport-api-consumer-" + serviceInfoProvider . getServiceId ( ) ) ;
consumerBuilder . groupId ( "monolith-transport-api-consumer" ) ;
consumerBuilder . groupId ( topicService . buildTopicName ( "monolith-transport-api-consumer" ) ) ;
consumerBuilder . decoder ( msg - > new TbProtoQueueMsg < > ( msg . getKey ( ) , TransportApiRequestMsg . parseFrom ( msg . getData ( ) ) , msg . getHeaders ( ) ) ) ;
consumerBuilder . admin ( transportApiRequestAdmin ) ;
consumerBuilder . statsService ( consumerStatsService ) ;
@ -251,7 +251,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
TbKafkaProducerTemplate . TbKafkaProducerTemplateBuilder < TbProtoQueueMsg < TransportApiResponseMsg > > requestBuilder = TbKafkaProducerTemplate . builder ( ) ;
requestBuilder . settings ( kafkaSettings ) ;
requestBuilder . clientId ( "monolith-transport-api-producer-" + serviceInfoProvider . getServiceId ( ) ) ;
requestBuilder . defaultTopic ( transportApiSettings . getResponsesTopic ( ) ) ;
requestBuilder . defaultTopic ( topicService . buildTopicName ( t ransportApiSettings . getResponsesTopic ( ) ) ) ;
requestBuilder . admin ( transportApiResponseAdmin ) ;
return requestBuilder . build ( ) ;
}
@ -294,9 +294,9 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
public TbQueueConsumer < TbProtoQueueMsg < ToUsageStatsServiceMsg > > createToUsageStatsServiceMsgConsumer ( ) {
TbKafkaConsumerTemplate . TbKafkaConsumerTemplateBuilder < TbProtoQueueMsg < ToUsageStatsServiceMsg > > consumerBuilder = TbKafkaConsumerTemplate . builder ( ) ;
consumerBuilder . settings ( kafkaSettings ) ;
consumerBuilder . topic ( coreSettings . getUsageStatsTopic ( ) ) ;
consumerBuilder . topic ( topicService . buildTopicName ( coreSettings . getUsageStatsTopic ( ) ) ) ;
consumerBuilder . clientId ( "monolith-us-consumer-" + serviceInfoProvider . getServiceId ( ) ) ;
consumerBuilder . groupId ( "monolith-us-consumer" ) ;
consumerBuilder . groupId ( topicService . buildTopicName ( "monolith-us-consumer" ) ) ;
consumerBuilder . decoder ( msg - > new TbProtoQueueMsg < > ( msg . getKey ( ) , ToUsageStatsServiceMsg . parseFrom ( msg . getData ( ) ) , msg . getHeaders ( ) ) ) ;
consumerBuilder . admin ( coreAdmin ) ;
consumerBuilder . statsService ( consumerStatsService ) ;
@ -307,9 +307,9 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
public TbQueueConsumer < TbProtoQueueMsg < ToOtaPackageStateServiceMsg > > createToOtaPackageStateServiceMsgConsumer ( ) {
TbKafkaConsumerTemplate . TbKafkaConsumerTemplateBuilder < TbProtoQueueMsg < ToOtaPackageStateServiceMsg > > consumerBuilder = TbKafkaConsumerTemplate . builder ( ) ;
consumerBuilder . settings ( kafkaSettings ) ;
consumerBuilder . topic ( coreSettings . getOtaPackageTopic ( ) ) ;
consumerBuilder . topic ( topicService . buildTopicName ( coreSettings . getOtaPackageTopic ( ) ) ) ;
consumerBuilder . clientId ( "monolith-ota-consumer-" + serviceInfoProvider . getServiceId ( ) ) ;
consumerBuilder . groupId ( "monolith-ota-consumer" ) ;
consumerBuilder . groupId ( topicService . buildTopicName ( "monolith-ota-consumer" ) ) ;
consumerBuilder . decoder ( msg - > new TbProtoQueueMsg < > ( msg . getKey ( ) , ToOtaPackageStateServiceMsg . parseFrom ( msg . getData ( ) ) , msg . getHeaders ( ) ) ) ;
consumerBuilder . admin ( fwUpdatesAdmin ) ;
consumerBuilder . statsService ( consumerStatsService ) ;
@ -321,7 +321,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
TbKafkaProducerTemplate . TbKafkaProducerTemplateBuilder < TbProtoQueueMsg < ToOtaPackageStateServiceMsg > > requestBuilder = TbKafkaProducerTemplate . builder ( ) ;
requestBuilder . settings ( kafkaSettings ) ;
requestBuilder . clientId ( "monolith-ota-producer-" + serviceInfoProvider . getServiceId ( ) ) ;
requestBuilder . defaultTopic ( coreSettings . getOtaPackageTopic ( ) ) ;
requestBuilder . defaultTopic ( topicService . buildTopicName ( coreSettings . getOtaPackageTopic ( ) ) ) ;
requestBuilder . admin ( fwUpdatesAdmin ) ;
return requestBuilder . build ( ) ;
}
@ -331,7 +331,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
TbKafkaProducerTemplate . TbKafkaProducerTemplateBuilder < TbProtoQueueMsg < ToUsageStatsServiceMsg > > requestBuilder = TbKafkaProducerTemplate . builder ( ) ;
requestBuilder . settings ( kafkaSettings ) ;
requestBuilder . clientId ( "monolith-us-producer-" + serviceInfoProvider . getServiceId ( ) ) ;
requestBuilder . defaultTopic ( coreSettings . getUsageStatsTopic ( ) ) ;
requestBuilder . defaultTopic ( topicService . buildTopicName ( coreSettings . getUsageStatsTopic ( ) ) ) ;
requestBuilder . admin ( coreAdmin ) ;
return requestBuilder . build ( ) ;
}
@ -341,7 +341,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
TbKafkaProducerTemplate . TbKafkaProducerTemplateBuilder < TbProtoQueueMsg < TransportProtos . ToVersionControlServiceMsg > > requestBuilder = TbKafkaProducerTemplate . builder ( ) ;
requestBuilder . settings ( kafkaSettings ) ;
requestBuilder . clientId ( "monolith-vc-producer-" + serviceInfoProvider . getServiceId ( ) ) ;
requestBuilder . defaultTopic ( vcSettings . getTopic ( ) ) ;
requestBuilder . defaultTopic ( topicService . buildTopicName ( vcSettings . getTopic ( ) ) ) ;
requestBuilder . admin ( vcAdmin ) ;
return requestBuilder . build ( ) ;
}