From 109240166e2d2fb7dc5e713cd81f69da30ae90a6 Mon Sep 17 00:00:00 2001 From: akulikov Date: Wed, 5 Jul 2023 15:17:24 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-358 send reports to gateway --- .../clearing/reports/config/KafkaConfig.java | 35 +++++++++++ .../reports/services/ReportService.java | 61 +++++++++++++------ .../platform/messaging/domain/Consts.java | 2 + .../domain/cud/gateway/ReportPart.java | 27 ++++++++ .../domain/cud/gateway/SendReportRequest.java | 40 ++++++++++++ 5 files changed, 145 insertions(+), 20 deletions(-) create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/gateway/ReportPart.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/gateway/SendReportRequest.java diff --git a/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/config/KafkaConfig.java b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/config/KafkaConfig.java index 3a7f95b94..8246741f0 100644 --- a/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/config/KafkaConfig.java +++ b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/config/KafkaConfig.java @@ -7,9 +7,18 @@ import org.springframework.beans.factory.config.ConfigurableBeanFactory; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Scope; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.core.ProducerFactory; +import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; +import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings; +import ru.spcex.clearing.platform.messaging.service.RequestInfo; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.reports.config.element.ReportsServiceSettings; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; @Configuration public class KafkaConfig { @@ -25,4 +34,30 @@ public class KafkaConfig { public Producer createProducer(ReportsServiceSettings settings) { return KafkaProducerFactory.producer(settings.getKafkaProducer()); } + + @Bean + public ProducerFactory pf(ReportsServiceSettings settings) { + KafkaProducerSettings kafkaSettings = settings.getKafkaProducer(); + return KafkaProducerFactory.producerFactory(kafkaSettings); + } + + @Bean + public KafkaTemplate kafkaTemplate(ProducerFactory pf) { + return new KafkaTemplate<>(pf); + } + + @Autowired + @Bean + public KafkaSender kafkaSender(KafkaTemplate kafkaTemplate, ImdgProvider imdgProvider) { + ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); + return KafkaSender + .setup() + .setKafkaTemplate(kafkaTemplate) + .idGenerator(imdgIdGenerator::nextId) + .imdgProvider(s -> { + Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); + return imdg::insert; + }) + .build(); + } } diff --git a/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/ReportService.java b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/ReportService.java index c0b5ee568..468ff4430 100644 --- a/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/ReportService.java +++ b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/ReportService.java @@ -9,6 +9,8 @@ import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.stereotype.Service; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.gateway.ReportPart; +import ru.spcex.clearing.platform.messaging.domain.cud.gateway.SendReportRequest; import ru.spcex.clearing.platform.messaging.domain.cud.reports.ReportRequest; import ru.spcex.clearing.platform.messaging.domain.cud.reports.ReportRequestWithPeriod; import ru.spcex.clearing.platform.messaging.domain.cud.reports.ReportRequestWithSessionId; @@ -16,6 +18,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.reports.ReportRequestWith import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.reports.config.element.ReportsServiceSettings; import ru.spcex.clearing.reports.reports.*; import ru.spcex.clearing.util.security.UserRoleVerification; @@ -27,10 +30,7 @@ import ru.spcex.platform.utils.validation.IValidator; import java.io.File; import java.time.LocalDate; -import java.util.ArrayList; -import java.util.Collection; -import java.util.List; -import java.util.Map; +import java.util.*; import java.util.function.Function; @Service @@ -38,6 +38,7 @@ public class ReportService extends QueueConsumer implements InitializingBean { private final Logger log = LoggerFactory.getLogger(getClass()); private final UserRoleVerification userRoleVerification; private final ValidationHelper validationHelper; + private final KafkaSender kafkaSender; private final Function reportRequestValidator; private final Function reportRequestWithSessionIdValidator; private final Function reportRequestWithSessionIdListValidator; @@ -54,6 +55,7 @@ public class ReportService extends QueueConsumer implements InitializingBean { public ReportService(Consumer kafkaQueue, Producer kafkaResponseQueue, + KafkaSender kafkaSender, UserRoleVerification userRoleVerification, ValidationHelper validationHelper, ReportsServiceSettings reportsServiceSettings, @@ -82,6 +84,7 @@ public class ReportService extends QueueConsumer implements InitializingBean { super(kafkaQueue, kafkaResponseQueue); this.userRoleVerification = userRoleVerification; this.validationHelper = validationHelper; + this.kafkaSender = kafkaSender; this.reportRequestValidator = reportRequestValidator; this.reportRequestWithSessionIdValidator = reportRequestWithSessionIdValidator; this.reportRequestWithSessionIdListValidator = reportRequestWithSessionIdListValidator; @@ -283,7 +286,7 @@ public class ReportService extends QueueConsumer implements InitializingBean { logUnknownProperties(userRequest); int cntErrors = 0; - List outFilenames = new ArrayList<>(); + Map outFilenames = new HashMap<>(); Collection>> buildersWithoutParams = reportBuildersWithoutParams.values(); for (List> builderList : buildersWithoutParams) { @@ -294,7 +297,7 @@ public class ReportService extends QueueConsumer implements InitializingBean { cntErrors++; continue; } - outFilenames.add(outFile.getAbsolutePath()); + outFilenames.put(outFile.getName(), builder.getReportType()); } } @@ -320,7 +323,7 @@ public class ReportService extends QueueConsumer implements InitializingBean { cntErrors++; continue; } - outFilenames.add(outFile.getAbsolutePath()); + outFilenames.put(outFile.getName(), builder.getReportType()); } } @@ -335,7 +338,7 @@ public class ReportService extends QueueConsumer implements InitializingBean { cntErrors++; continue; } - outFilenames.add(outFile.getAbsolutePath()); + outFilenames.put(outFile.getName(), builder.getReportType()); } } @@ -354,13 +357,14 @@ public class ReportService extends QueueConsumer implements InitializingBean { cntErrors++; continue; } - outFilenames.add(outFile.getAbsolutePath()); + outFilenames.put(outFile.getName(), builder.getReportType()); } } - if (cntErrors == 0) log.debug("All reports successfully created. Output files: {}", String.join(", ", outFilenames)); - else log.debug("Some reports ({}) create failed. Output files: {}", cntErrors, String.join(", ", outFilenames)); + if (cntErrors == 0) log.debug("All reports successfully created. Output files: {}", String.join(", ", outFilenames.keySet())); + else log.debug("Some reports ({}) create failed. Output files: {}", cntErrors, String.join(", ", outFilenames.keySet())); + sendReportsToGateway(sessionId, outFilenames, "REPORT_KS_TMP"); return null; } @@ -377,7 +381,7 @@ public class ReportService extends QueueConsumer implements InitializingBean { logUnknownProperties(userRequest); int cntErrors = 0; - List outFilenames = new ArrayList<>(); + Map outFilenames = new HashMap<>(); if (sessionId != null) { @@ -399,14 +403,15 @@ public class ReportService extends QueueConsumer implements InitializingBean { cntErrors++; continue; } - outFilenames.add(outFile.getAbsolutePath()); + outFilenames.put(outFile.getName(), builder.getReportType()); } } } - if (cntErrors == 0) log.debug("All reports successfully created. Output files: {}", String.join(", ", outFilenames)); - else log.debug("Some reports ({}) create failed. Output files: {}", cntErrors, String.join(", ", outFilenames)); + if (cntErrors == 0) log.debug("All reports successfully created. Output files: {}", String.join(", ", outFilenames.keySet())); + else log.debug("Some reports ({}) create failed. Output files: {}", cntErrors, String.join(", ", outFilenames.keySet())); + sendReportsToGateway(sessionId, outFilenames, null); return null; } @@ -424,7 +429,7 @@ public class ReportService extends QueueConsumer implements InitializingBean { logUnknownProperties(userRequest); int cntErrors = 0; - List outFilenames = new ArrayList<>(); + Map outFilenames = new HashMap<>(); Collection>> buildersWithoutParams = reportBuildersWithoutParams.values(); for (List> builderList : buildersWithoutParams) { @@ -446,7 +451,7 @@ public class ReportService extends QueueConsumer implements InitializingBean { cntErrors++; continue; } - outFilenames.add(outFile.getAbsolutePath()); + outFilenames.put(outFile.getName(), builder.getReportType()); } } @@ -469,17 +474,33 @@ public class ReportService extends QueueConsumer implements InitializingBean { cntErrors++; continue; } - outFilenames.add(outFile.getAbsolutePath()); + outFilenames.put(outFile.getName(), builder.getReportType()); } } } - if (cntErrors == 0) log.debug("All reports successfully created. Output files: {}", String.join(", ", outFilenames)); - else log.debug("Some reports ({}) create failed. Output files: {}", cntErrors, String.join(", ", outFilenames)); + if (cntErrors == 0) log.debug("All reports successfully created. Output files: {}", String.join(", ", outFilenames.keySet())); + else log.debug("Some reports ({}) create failed. Output files: {}", cntErrors, String.join(", ", outFilenames.keySet())); + sendReportsToGateway(sessionId, outFilenames, "REPORT_KS_FINAL"); return null; } + private void sendReportsToGateway(Long sessionId, Map reportsWithReportType, String type) { + SendReportRequest sendReportRequest = new SendReportRequest(); + sendReportRequest.setSessionId(sessionId); + sendReportRequest.setType(type); + List reports = new ArrayList<>(reportsWithReportType.size()); + for (Map.Entry entry : reportsWithReportType.entrySet()) { + ReportPart reportPart = new ReportPart(); + reportPart.setReportType(entry.getValue().getKey()); + reportPart.setFileName(entry.getKey()); + reports.add(reportPart); + } + sendReportRequest.setReports(reports); + kafkaSender.sendRequestToQueue(Consts.REPORTS_TO_GATEWAY, sendReportRequest); + } + private void logUnknownProperties(BaseRequest request) { Map unknownProperties = request.getUnknownProperties(); for (Map.Entry unknownProperty : unknownProperties.entrySet()) { diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java index d3f8f6ee4..439e82098 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java @@ -182,6 +182,8 @@ public interface Consts { String REGISTRY_CLEAR_MEMBER_REGISTER_NEW = "registry-clear-member-register-new"; String REGISTRY_CLEAR_MEMBER_REGISTER_CHANGE_NEW = "registry-clear-member-register-change-new"; + String REPORTS_TO_GATEWAY = "reports-to-gateway"; + String REGISTRY_NEW = "registry-new"; } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/gateway/ReportPart.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/gateway/ReportPart.java new file mode 100644 index 000000000..ffea38146 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/gateway/ReportPart.java @@ -0,0 +1,27 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.gateway; + +import com.fasterxml.jackson.annotation.JsonProperty; + +public class ReportPart { + @JsonProperty + public String reportType; + + @JsonProperty + public String fileName; + + public String getReportType() { + return reportType; + } + + public void setReportType(String reportType) { + this.reportType = reportType; + } + + public String getFileName() { + return fileName; + } + + public void setFileName(String fileName) { + this.fileName = fileName; + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/gateway/SendReportRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/gateway/SendReportRequest.java new file mode 100644 index 000000000..07708f1e9 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/gateway/SendReportRequest.java @@ -0,0 +1,40 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.gateway; + +import com.fasterxml.jackson.annotation.JsonProperty; + +import java.util.List; + +public class SendReportRequest { + @JsonProperty + public Long sessionId; + + @JsonProperty + public String type; + + @JsonProperty + public List reports; + + public Long getSessionId() { + return sessionId; + } + + public void setSessionId(Long sessionId) { + this.sessionId = sessionId; + } + + public String getType() { + return type; + } + + public void setType(String type) { + this.type = type; + } + + public List getReports() { + return reports; + } + + public void setReports(List reports) { + this.reports = reports; + } +}