company-service http://jira.mfd.msk:8088/browse/CLS-272 исправил использование BiDirectionQueueExchanger

This commit is contained in:
AKurakin 2023-05-04 16:29:52 +03:00
parent f1a803aacf
commit 864bbcf125
4 changed files with 5 additions and 5 deletions

View file

@ -55,7 +55,7 @@ public class BiDirectionQueueExchanger<TOut extends BaseRequest<?>> extends Queu
/**
* Синхронно-ассинхронный обмен сообщениями
*
* @param kafkaQueue
* @param kafkaQueue уникальный kafka Consumer (Spring prototype). Нельзя переиспользовать существующие.
* @param kafkaProducer
* @param outQueue отправляет в очередь
* @param inQueue слушает очередь/топик, ожидает ответов. Пример: Consts.CONTINUE_CLEARING

View file

@ -70,7 +70,7 @@ public class ClientCodeService extends QueueConsumer implements InitializingBean
BiDirectionQueueExchanger<BaseRequest<TradingClearingRegistryNewRequest>> accountServiceExchanger;
@Autowired
public ClientCodeService(Consumer<String, Object> kafkaQueue,
public ClientCodeService(Consumer<String, Object> kafkaQueue1, Consumer<String, Object> kafkaQueue2,
Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider,
ValidationHelper validationHelper,
@ -79,7 +79,7 @@ public class ClientCodeService extends QueueConsumer implements InitializingBean
@Qualifier("clientCodeNewRequestValidator") Function<ClientCodeNewRequest, IValidator> clientCodeNewRequestValidator,
@Qualifier("clientCodeUpdateRequestValidator") Function<ClientCodeUpdateRequest, IValidator> clientCodeUpdateRequestValidator,
@Qualifier("clientCodeDeleteRequestValidator") Function<CommonDeleteRequest, IValidator> clientCodeDeleteRequestValidator) {
super(kafkaQueue, kafkaProducer);
super(kafkaQueue1, kafkaProducer);
this.kafkaProducer = kafkaProducer;
this.idGenerator = imdgProvider.getImdgIdGenerator();
this.clientCodeMap = imdgProvider.getImdg(IMDGDistributedNames.Map_ClientCode, ClientCode.class);
@ -92,7 +92,7 @@ public class ClientCodeService extends QueueConsumer implements InitializingBean
this.messageResolver = messageResolver;
this.userRoleVerification = userRoleVerification;
this.accountServiceExchanger = new BiDirectionQueueExchanger<>(kafkaQueue, kafkaProducer,
this.accountServiceExchanger = new BiDirectionQueueExchanger<>(kafkaQueue2, kafkaProducer,
Consts.DESTINATION_TRADING_CLEARING_REGISTRY_NEW,
Consts.REQUEST_INFO_UPDATE, true, // REQUEST_INFO_UPDATE - стандартная очередь, для результатов всех реквестов. DESTINATION_TRADING_CLEARING_REGISTRY_REPLY,
60000

View file

@ -45,7 +45,6 @@ import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparato
import static ru.spcex.clearing.test.TestUtils.*;
import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId;
//todo тесты случайно ломаются из-за BiDirectionQueueExchanger. Когда там выключаю initReplyListener() то тесты стабильно работают.
@ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = {
ClientCodeService.class,

View file

@ -74,6 +74,7 @@ public class QueueConsumer implements AutoCloseable {
public void init() {
inputExecutor.submit(() -> {
try {
log.debug("Using kafka consumer {} for subscribe on \"{}\"", consumer, callbacks.keySet());
if (supportStartOffsetTimeWindow) {
consumer.subscribe(callbacks.keySet(), new OffsetChanger(consumer, callbacks.keySet()));
} else {