diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/SdfCreatorBySTLDPayment.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/SdfCreatorBySTLDPayment.java index 58904a0f7..9438e3c93 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/SdfCreatorBySTLDPayment.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/SdfCreatorBySTLDPayment.java @@ -69,11 +69,14 @@ public class SdfCreatorBySTLDPayment { log.info("Group category {}, index {}", entry.getKey().category, entry.getKey().getSortIndex()); GroupWrapper group = entry.getKey(); Collection> bySenderGroups = entry.getValue(); + Set queueToSend = new HashSet<>(); for (List bySenderGroup : bySenderGroups) { Long senderIdSpecificToGroup = bySenderGroup.iterator().next().getSenderId(); log.info("Group senderId {}", senderIdSpecificToGroup); - saveSdfAnSendToKafka(senderIdSpecificToGroup, group.category, bySenderGroup, generationId); + queueToSend.add(saveSdfBeforeSendToKafka(senderIdSpecificToGroup, group.category, bySenderGroup, generationId)); } + log.info("Send command target queue: {}", queueToSend); + sendToKafkaSDFRequests(queueToSend, generationId); }); allEntries.forEach(entry -> { //category group @@ -120,24 +123,44 @@ public class SdfCreatorBySTLDPayment { }).toList(); } - private void saveSdfAnSendToKafka(Long senderId, ClearingCategory category, List batch, Long generationId) { - SdfClearingRequest kafkaMessage = new SdfClearingRequest(); - kafkaMessage.setGroupId(generationId); - Long requestId = null; + /** + * + * @param senderId + * @param category + * @param batch + * @param generationId + * @return kafka queue: Consts.SDF03_PROCESS/Consts.SDF11_PROCESS + */ + private String saveSdfBeforeSendToKafka(Long senderId, ClearingCategory category, List batch, Long generationId) { + String queue = null; if (Sender.One.equalsById(senderId) || category.equals(ClearingCategory.I) || category.equals(ClearingCategory.V)) { batch.stream() .map(paymentInstruction -> new Sdf03Builder(companyImdg).buildSdf03(paymentInstruction, generationId)) .forEach(sdf03Imdg::insert); - requestId = kafkaSender.sendRequestToQueue(Consts.SDF03_PROCESS, kafkaMessage); + queue = Consts.SDF03_PROCESS; } else if (category.equals(ClearingCategory.B)) { batch.stream() .map(paymentInstruction -> new Sdf11Builder(companyImdg).buildSdf11(paymentInstruction, generationId)) .forEach(sdf11Imdg::insert); - requestId = kafkaSender.sendRequestToQueue(Consts.SDF11_PROCESS, kafkaMessage); + queue = Consts.SDF11_PROCESS; } else { log.error("Unknown action (sdf03/sdf11) senderId {}, generationId {}, category {}", senderId, generationId, category); } - log.debug("Send to kafka command, CategoryD={}, requestId={}", category, requestId); + log.debug("Prepare send to kafka command, CategoryD={}, destination queue={}", category, queue); + return queue; + } + + private void sendToKafkaSDFRequests(Collection queues, Long generationId) { + for (String kQueue : queues) { + if (kQueue == null) { + log.trace("Ignore null queue destination"); + continue; + } + SdfClearingRequest kafkaMessage = new SdfClearingRequest(); + kafkaMessage.setGroupId(generationId); + Long requestId = kafkaSender.sendRequestToQueue(Consts.SDF03_PROCESS, kafkaMessage); + log.debug("Send to kafka command, queue={}, requestId={}", kQueue, requestId); + } } public static class GroupWrapper implements Comparable {