diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/EndStageNotification.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/EndStageNotification.java index 8a27ff254..9766e1b64 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/EndStageNotification.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/EndStageNotification.java @@ -1,5 +1,8 @@ package ru.spcex.clearing.session.stage.impl; +import java.time.Instant; +import java.util.Collection; +import java.util.stream.Collectors; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; @@ -8,6 +11,8 @@ import org.springframework.context.annotation.Scope; import org.springframework.stereotype.Service; import ru.clearing.classes.statics.data.misc.Session; import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.component.GroupRgsKey; +import static ru.spcex.clearing.error.ClearingErrorInternal.SessionGeneralError; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SdfClearingRequest; @@ -24,13 +29,6 @@ import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.utils.enumeration.EnumMessage; import ru.spcex.platform.utils.enumeration.IMessageResolver; -import java.time.Instant; -import java.util.Collection; -import java.util.Objects; -import java.util.stream.Collectors; - -import static ru.spcex.clearing.error.ClearingErrorInternal.SessionGeneralError; - @Service @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) public class EndStageNotification implements ISessionStage { @@ -83,25 +81,24 @@ public class EndStageNotification implements ISessionStage { } Collection forRegistries = selectRegistry(); - Collection groups = forRegistries.stream() - .map(Registry::getGroupId) - .filter(Objects::nonNull) + Collection groups = forRegistries.stream() + .map(registry -> new GroupRgsKey(registry.getGroupId(), registry.getMarket())) .distinct().collect(Collectors.toList()); log.debug("Sending notifications for {} groups ({} registers) on section {}", groups, forRegistries.size(), section); - for (Long groupId : groups) { - log.trace("For registry group {} send notification", - groupId); + for (GroupRgsKey groupKey : groups) { + log.trace("For registry group/market {}/{} send notification", + groupKey.groupId(), groupKey.market()); StageResult sResult; if (Section.FOND.equalsByKey(section)) { - sResult = notificationDF14(groupId); + sResult = notificationDF14(groupKey.groupId()); } else if (Section.MKR.equalsByKey(section)) { - sResult = notificationDF05(groupId); + sResult = notificationDF05(groupKey.groupId()); } else { throw new IllegalArgumentException("Unsupported section " + section); } if (sResult.getError() != null) { - log.warn("When sending groupId={} has error: {}", groupId, msgResolver.resolve(sResult.getError())); + log.warn("When sending groupId={} has error: {}", groupKey, msgResolver.resolve(sResult.getError())); } } StageResult> res = new StageResult<>(null, true);