Fix for Transport API subscriptions
This commit is contained in:
parent
acf900e1de
commit
676868b4d3
@ -200,7 +200,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
|
|||||||
consumerBuilder.settings(kafkaSettings);
|
consumerBuilder.settings(kafkaSettings);
|
||||||
consumerBuilder.topic(transportApiSettings.getRequestsTopic());
|
consumerBuilder.topic(transportApiSettings.getRequestsTopic());
|
||||||
consumerBuilder.clientId("monolith-transport-api-consumer-" + serviceInfoProvider.getServiceId());
|
consumerBuilder.clientId("monolith-transport-api-consumer-" + serviceInfoProvider.getServiceId());
|
||||||
consumerBuilder.groupId("monolith-transport-api-consumer-" + serviceInfoProvider.getServiceId());
|
consumerBuilder.groupId("monolith-transport-api-consumer");
|
||||||
consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportApiRequestMsg.parseFrom(msg.getData()), msg.getHeaders()));
|
consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportApiRequestMsg.parseFrom(msg.getData()), msg.getHeaders()));
|
||||||
consumerBuilder.admin(transportApiAdmin);
|
consumerBuilder.admin(transportApiAdmin);
|
||||||
return consumerBuilder.build();
|
return consumerBuilder.build();
|
||||||
|
|||||||
@ -170,7 +170,7 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory {
|
|||||||
consumerBuilder.settings(kafkaSettings);
|
consumerBuilder.settings(kafkaSettings);
|
||||||
consumerBuilder.topic(transportApiSettings.getRequestsTopic());
|
consumerBuilder.topic(transportApiSettings.getRequestsTopic());
|
||||||
consumerBuilder.clientId("tb-core-transport-api-consumer-" + serviceInfoProvider.getServiceId());
|
consumerBuilder.clientId("tb-core-transport-api-consumer-" + serviceInfoProvider.getServiceId());
|
||||||
consumerBuilder.groupId("tb-core-transport-api-consumer-" + serviceInfoProvider.getServiceId());
|
consumerBuilder.groupId("tb-core-transport-api-consumer");
|
||||||
consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportApiRequestMsg.parseFrom(msg.getData()), msg.getHeaders()));
|
consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportApiRequestMsg.parseFrom(msg.getData()), msg.getHeaders()));
|
||||||
consumerBuilder.admin(transportApiAdmin);
|
consumerBuilder.admin(transportApiAdmin);
|
||||||
return consumerBuilder.build();
|
return consumerBuilder.build();
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user