http://jira.mfd.msk:8088/browse/CLS-189 fix SDF 03/11 отправку сообщений на создание
This commit is contained in:
parent
0a1857ead8
commit
fde8174e86
1 changed files with 31 additions and 8 deletions
|
|
@ -69,11 +69,14 @@ public class SdfCreatorBySTLDPayment {
|
||||||
log.info("Group category {}, index {}", entry.getKey().category, entry.getKey().getSortIndex());
|
log.info("Group category {}, index {}", entry.getKey().category, entry.getKey().getSortIndex());
|
||||||
GroupWrapper group = entry.getKey();
|
GroupWrapper group = entry.getKey();
|
||||||
Collection<List<PaymentInstruction>> bySenderGroups = entry.getValue();
|
Collection<List<PaymentInstruction>> bySenderGroups = entry.getValue();
|
||||||
|
Set<String> queueToSend = new HashSet<>();
|
||||||
for (List<PaymentInstruction> bySenderGroup : bySenderGroups) {
|
for (List<PaymentInstruction> bySenderGroup : bySenderGroups) {
|
||||||
Long senderIdSpecificToGroup = bySenderGroup.iterator().next().getSenderId();
|
Long senderIdSpecificToGroup = bySenderGroup.iterator().next().getSenderId();
|
||||||
log.info("Group senderId {}", senderIdSpecificToGroup);
|
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 -> {
|
allEntries.forEach(entry -> {
|
||||||
//category group
|
//category group
|
||||||
|
|
@ -120,24 +123,44 @@ public class SdfCreatorBySTLDPayment {
|
||||||
}).toList();
|
}).toList();
|
||||||
}
|
}
|
||||||
|
|
||||||
private void saveSdfAnSendToKafka(Long senderId, ClearingCategory category, List<PaymentInstruction> batch, Long generationId) {
|
/**
|
||||||
SdfClearingRequest kafkaMessage = new SdfClearingRequest();
|
*
|
||||||
kafkaMessage.setGroupId(generationId);
|
* @param senderId
|
||||||
Long requestId = null;
|
* @param category
|
||||||
|
* @param batch
|
||||||
|
* @param generationId
|
||||||
|
* @return kafka queue: Consts.SDF03_PROCESS/Consts.SDF11_PROCESS
|
||||||
|
*/
|
||||||
|
private String saveSdfBeforeSendToKafka(Long senderId, ClearingCategory category, List<PaymentInstruction> batch, Long generationId) {
|
||||||
|
String queue = null;
|
||||||
if (Sender.One.equalsById(senderId) || category.equals(ClearingCategory.I) || category.equals(ClearingCategory.V)) {
|
if (Sender.One.equalsById(senderId) || category.equals(ClearingCategory.I) || category.equals(ClearingCategory.V)) {
|
||||||
batch.stream()
|
batch.stream()
|
||||||
.map(paymentInstruction -> new Sdf03Builder(companyImdg).buildSdf03(paymentInstruction, generationId))
|
.map(paymentInstruction -> new Sdf03Builder(companyImdg).buildSdf03(paymentInstruction, generationId))
|
||||||
.forEach(sdf03Imdg::insert);
|
.forEach(sdf03Imdg::insert);
|
||||||
requestId = kafkaSender.sendRequestToQueue(Consts.SDF03_PROCESS, kafkaMessage);
|
queue = Consts.SDF03_PROCESS;
|
||||||
} else if (category.equals(ClearingCategory.B)) {
|
} else if (category.equals(ClearingCategory.B)) {
|
||||||
batch.stream()
|
batch.stream()
|
||||||
.map(paymentInstruction -> new Sdf11Builder(companyImdg).buildSdf11(paymentInstruction, generationId))
|
.map(paymentInstruction -> new Sdf11Builder(companyImdg).buildSdf11(paymentInstruction, generationId))
|
||||||
.forEach(sdf11Imdg::insert);
|
.forEach(sdf11Imdg::insert);
|
||||||
requestId = kafkaSender.sendRequestToQueue(Consts.SDF11_PROCESS, kafkaMessage);
|
queue = Consts.SDF11_PROCESS;
|
||||||
} else {
|
} else {
|
||||||
log.error("Unknown action (sdf03/sdf11) senderId {}, generationId {}, category {}", senderId, generationId, category);
|
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<String> 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<GroupWrapper> {
|
public static class GroupWrapper implements Comparable<GroupWrapper> {
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue