Commit 1b905ef7270ca73e328e7d7ed936c489a0c69d3a
Committed by
Andrew Shvayka
1 parent
4de258f2
refactored
Showing
4 changed files
with
9 additions
and
9 deletions
... | ... | @@ -165,7 +165,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi |
165 | 165 | consumerBuilder.settings(kafkaSettings); |
166 | 166 | consumerBuilder.topic(partitionService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceInfoProvider.getServiceId()).getFullTopicName()); |
167 | 167 | consumerBuilder.clientId("monolith-rule-engine-notifications-consumer-" + serviceInfoProvider.getServiceId()); |
168 | - consumerBuilder.groupId("monolith-rule-engine-notifications-consumer"); | |
168 | + consumerBuilder.groupId("monolith-rule-engine-notifications-consumer-" + serviceInfoProvider.getServiceId()); | |
169 | 169 | consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineNotificationMsg.parseFrom(msg.getData()), msg.getHeaders())); |
170 | 170 | consumerBuilder.admin(notificationAdmin); |
171 | 171 | return consumerBuilder.build(); |
... | ... | @@ -189,7 +189,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi |
189 | 189 | consumerBuilder.settings(kafkaSettings); |
190 | 190 | consumerBuilder.topic(partitionService.getNotificationsTopic(ServiceType.TB_CORE, serviceInfoProvider.getServiceId()).getFullTopicName()); |
191 | 191 | consumerBuilder.clientId("monolith-core-notifications-consumer-" + serviceInfoProvider.getServiceId()); |
192 | - consumerBuilder.groupId("monolith-core-notifications-consumer"); | |
192 | + consumerBuilder.groupId("monolith-core-notifications-consumer-" + serviceInfoProvider.getServiceId()); | |
193 | 193 | consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreNotificationMsg.parseFrom(msg.getData()), msg.getHeaders())); |
194 | 194 | consumerBuilder.admin(notificationAdmin); |
195 | 195 | return consumerBuilder.build(); |
... | ... | @@ -230,7 +230,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi |
230 | 230 | responseBuilder.settings(kafkaSettings); |
231 | 231 | responseBuilder.topic(jsInvokeSettings.getResponseTopic() + "." + serviceInfoProvider.getServiceId()); |
232 | 232 | responseBuilder.clientId("js-" + serviceInfoProvider.getServiceId()); |
233 | - responseBuilder.groupId("rule-engine-node"); | |
233 | + responseBuilder.groupId("rule-engine-node-" + serviceInfoProvider.getServiceId()); | |
234 | 234 | responseBuilder.decoder(msg -> { |
235 | 235 | JsInvokeProtos.RemoteJsResponse.Builder builder = JsInvokeProtos.RemoteJsResponse.newBuilder(); |
236 | 236 | JsonFormat.parser().ignoringUnknownFields().merge(new String(msg.getData(), StandardCharsets.UTF_8), builder); | ... | ... |
... | ... | @@ -159,7 +159,7 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory { |
159 | 159 | consumerBuilder.settings(kafkaSettings); |
160 | 160 | consumerBuilder.topic(partitionService.getNotificationsTopic(ServiceType.TB_CORE, serviceInfoProvider.getServiceId()).getFullTopicName()); |
161 | 161 | consumerBuilder.clientId("tb-core-notifications-consumer-" + serviceInfoProvider.getServiceId()); |
162 | - consumerBuilder.groupId("tb-core-notifications-node"); | |
162 | + consumerBuilder.groupId("tb-core-notifications-node-" + serviceInfoProvider.getServiceId()); | |
163 | 163 | consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreNotificationMsg.parseFrom(msg.getData()), msg.getHeaders())); |
164 | 164 | consumerBuilder.admin(notificationAdmin); |
165 | 165 | return consumerBuilder.build(); |
... | ... | @@ -200,7 +200,7 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory { |
200 | 200 | responseBuilder.settings(kafkaSettings); |
201 | 201 | responseBuilder.topic(jsInvokeSettings.getResponseTopic() + "." + serviceInfoProvider.getServiceId()); |
202 | 202 | responseBuilder.clientId("js-" + serviceInfoProvider.getServiceId()); |
203 | - responseBuilder.groupId("rule-engine-node"); | |
203 | + responseBuilder.groupId("rule-engine-node-" + serviceInfoProvider.getServiceId()); | |
204 | 204 | responseBuilder.decoder(msg -> { |
205 | 205 | JsInvokeProtos.RemoteJsResponse.Builder builder = JsInvokeProtos.RemoteJsResponse.newBuilder(); |
206 | 206 | JsonFormat.parser().ignoringUnknownFields().merge(new String(msg.getData(), StandardCharsets.UTF_8), builder); | ... | ... |
... | ... | @@ -154,7 +154,7 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { |
154 | 154 | consumerBuilder.settings(kafkaSettings); |
155 | 155 | consumerBuilder.topic(partitionService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceInfoProvider.getServiceId()).getFullTopicName()); |
156 | 156 | consumerBuilder.clientId("tb-rule-engine-notifications-consumer-" + serviceInfoProvider.getServiceId()); |
157 | - consumerBuilder.groupId("tb-rule-engine-notifications-node"); | |
157 | + consumerBuilder.groupId("tb-rule-engine-notifications-node-" + serviceInfoProvider.getServiceId()); | |
158 | 158 | consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineNotificationMsg.parseFrom(msg.getData()), msg.getHeaders())); |
159 | 159 | consumerBuilder.admin(notificationAdmin); |
160 | 160 | return consumerBuilder.build(); |
... | ... | @@ -173,7 +173,7 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { |
173 | 173 | responseBuilder.settings(kafkaSettings); |
174 | 174 | responseBuilder.topic(jsInvokeSettings.getResponseTopic() + "." + serviceInfoProvider.getServiceId()); |
175 | 175 | responseBuilder.clientId("js-" + serviceInfoProvider.getServiceId()); |
176 | - responseBuilder.groupId("rule-engine-node"); | |
176 | + responseBuilder.groupId("rule-engine-node-" + serviceInfoProvider.getServiceId()); | |
177 | 177 | responseBuilder.decoder(msg -> { |
178 | 178 | JsInvokeProtos.RemoteJsResponse.Builder builder = JsInvokeProtos.RemoteJsResponse.newBuilder(); |
179 | 179 | JsonFormat.parser().ignoringUnknownFields().merge(new String(msg.getData(), StandardCharsets.UTF_8), builder); | ... | ... |
... | ... | @@ -92,7 +92,7 @@ public class KafkaTbTransportQueueFactory implements TbTransportQueueFactory { |
92 | 92 | responseBuilder.settings(kafkaSettings); |
93 | 93 | responseBuilder.topic(transportApiSettings.getResponsesTopic() + "." + serviceInfoProvider.getServiceId()); |
94 | 94 | responseBuilder.clientId("transport-api-response-" + serviceInfoProvider.getServiceId()); |
95 | - responseBuilder.groupId("transport-node"); | |
95 | + responseBuilder.groupId("transport-node-" + serviceInfoProvider.getServiceId()); | |
96 | 96 | responseBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportApiResponseMsg.parseFrom(msg.getData()), msg.getHeaders())); |
97 | 97 | responseBuilder.admin(transportApiAdmin); |
98 | 98 | |
... | ... | @@ -133,7 +133,7 @@ public class KafkaTbTransportQueueFactory implements TbTransportQueueFactory { |
133 | 133 | responseBuilder.settings(kafkaSettings); |
134 | 134 | responseBuilder.topic(transportNotificationSettings.getNotificationsTopic() + "." + serviceInfoProvider.getServiceId()); |
135 | 135 | responseBuilder.clientId("transport-api-notifications-" + serviceInfoProvider.getServiceId()); |
136 | - responseBuilder.groupId("transport-node"); | |
136 | + responseBuilder.groupId("transport-node-" + serviceInfoProvider.getServiceId()); | |
137 | 137 | responseBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToTransportMsg.parseFrom(msg.getData()), msg.getHeaders())); |
138 | 138 | responseBuilder.admin(notificationAdmin); |
139 | 139 | return responseBuilder.build(); | ... | ... |