From 95f7741111dd9df2efe289bc319103be71b92bf4 Mon Sep 17 00:00:00 2001 From: Ivan Nikolaev-Axenov Date: Thu, 4 Jul 2024 16:21:43 +0300 Subject: [PATCH] Snapshot maker created, new task code added to Task.java --- clearing-parent/pom.xml | 1 + clearing-parent/snapshot-maker/pom.xml | 91 ++++++++++ .../maker/SnapshotMakerApplication.java | 11 ++ .../snapshot/maker/config/KafkaConfig.java | 32 ++++ .../config/settings/DatabaseSettings.java | 31 ++++ .../config/settings/LocationSettings.java | 67 +++++++ .../settings/SnapshotMakerSettings.java | 49 +++++ .../maker/service/CommandService.java | 52 ++++++ .../maker/service/SessionService.java | 167 ++++++++++++++++++ .../src/main/resources/application.properties | 30 ++++ .../ru/spcex/platform/enumeration/Task.java | 2 + pom.xml | 1 + 12 files changed, 534 insertions(+) create mode 100644 clearing-parent/snapshot-maker/pom.xml create mode 100644 clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/SnapshotMakerApplication.java create mode 100644 clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/config/KafkaConfig.java create mode 100644 clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/config/settings/DatabaseSettings.java create mode 100644 clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/config/settings/LocationSettings.java create mode 100644 clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/config/settings/SnapshotMakerSettings.java create mode 100644 clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/service/CommandService.java create mode 100644 clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/service/SessionService.java create mode 100644 clearing-parent/snapshot-maker/src/main/resources/application.properties diff --git a/clearing-parent/pom.xml b/clearing-parent/pom.xml index c5b1fafd3..a159e39c7 100644 --- a/clearing-parent/pom.xml +++ b/clearing-parent/pom.xml @@ -43,6 +43,7 @@ swt-importer gateway-api imdg-hist + snapshot-maker diff --git a/clearing-parent/snapshot-maker/pom.xml b/clearing-parent/snapshot-maker/pom.xml new file mode 100644 index 000000000..1db3214c6 --- /dev/null +++ b/clearing-parent/snapshot-maker/pom.xml @@ -0,0 +1,91 @@ + + + 4.0.0 + + snapshot-maker + snapshot-maker + Database snapshot maker + SPCEX-3.12.5 + + + ru.spcex.clearing + clearing-parent + SPCEX-3.12.5 + + + + 17 + 17 + + 5.1.0 + 2.5.0 + 2.16.1 + + + + + + org.springframework.boot + spring-boot-starter-web + + + + + ru.spcex.platform + platform-messaging + + + + + commons-io + commons-io + ${commons-io.version} + + + ru.spcex.clearing + classes + + + ru.spcex.platform + platform-enum + + + + + + org.springframework.boot + spring-boot-starter-test + test + + + + + + + src/main/resources + + application.properties + + false + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + repackage + + + + + ${project.artifactId} + + + + + \ No newline at end of file diff --git a/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/SnapshotMakerApplication.java b/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/SnapshotMakerApplication.java new file mode 100644 index 000000000..66edf48a7 --- /dev/null +++ b/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/SnapshotMakerApplication.java @@ -0,0 +1,11 @@ +package ru.spcex.clearing.snapshot.maker; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class SnapshotMakerApplication { + public static void main(String[] args) { + SpringApplication.run(SnapshotMakerApplication.class, args); + } +} diff --git a/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/config/KafkaConfig.java b/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/config/KafkaConfig.java new file mode 100644 index 000000000..18756143e --- /dev/null +++ b/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/config/KafkaConfig.java @@ -0,0 +1,32 @@ +package ru.spcex.clearing.snapshot.maker.config; + +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; +import org.springframework.beans.factory.annotation.Autowired; +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 ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; +import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; +import ru.spcex.clearing.snapshot.maker.config.settings.SnapshotMakerSettings; + +@Configuration +public class KafkaConfig { + @Autowired + @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) + @Bean + public Consumer createConsumer(SnapshotMakerSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); + } + + @Autowired + @Bean + public Producer createProducer(SnapshotMakerSettings settings) { + if (settings.getKafkaProducer() != null) { + return KafkaProducerFactory.producer(settings.getKafkaProducer()); + } else { + return null; + } + } +} diff --git a/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/config/settings/DatabaseSettings.java b/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/config/settings/DatabaseSettings.java new file mode 100644 index 000000000..5cc1d6a4d --- /dev/null +++ b/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/config/settings/DatabaseSettings.java @@ -0,0 +1,31 @@ +package ru.spcex.clearing.snapshot.maker.config.settings; + +public class DatabaseSettings { + private String jdbcUrl; + private String username; + private String password; + + public String getJdbcUrl() { + return jdbcUrl; + } + + public void setJdbcUrl(String jdbcUrl) { + this.jdbcUrl = jdbcUrl; + } + + public String getUsername() { + return username; + } + + public void setUsername(String username) { + this.username = username; + } + + public String getPassword() { + return password; + } + + public void setPassword(String password) { + this.password = password; + } +} diff --git a/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/config/settings/LocationSettings.java b/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/config/settings/LocationSettings.java new file mode 100644 index 000000000..78bf0ecf2 --- /dev/null +++ b/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/config/settings/LocationSettings.java @@ -0,0 +1,67 @@ +package ru.spcex.clearing.snapshot.maker.config.settings; + +public class LocationSettings { + private String localStore; + private String dbfImporterFolder; + private String dbfExporterFolder; + private String swtImporterFolder; + private String swtExporterFolder; + private String xmlImporterFolder; + private String xmlExporterFolder; + + public String getLocalStore() { + return localStore; + } + + public void setLocalStore(String localStore) { + this.localStore = localStore; + } + + public String getDbfImporterFolder() { + return dbfImporterFolder; + } + + public void setDbfImporterFolder(String dbfImporterFolder) { + this.dbfImporterFolder = dbfImporterFolder; + } + + public String getDbfExporterFolder() { + return dbfExporterFolder; + } + + public void setDbfExporterFolder(String dbfExporterFolder) { + this.dbfExporterFolder = dbfExporterFolder; + } + + public String getSwtImporterFolder() { + return swtImporterFolder; + } + + public void setSwtImporterFolder(String swtImporterFolder) { + this.swtImporterFolder = swtImporterFolder; + } + + public String getSwtExporterFolder() { + return swtExporterFolder; + } + + public void setSwtExporterFolder(String swtExporterFolder) { + this.swtExporterFolder = swtExporterFolder; + } + + public String getXmlImporterFolder() { + return xmlImporterFolder; + } + + public void setXmlImporterFolder(String xmlImporterFolder) { + this.xmlImporterFolder = xmlImporterFolder; + } + + public String getXmlExporterFolder() { + return xmlExporterFolder; + } + + public void setXmlExporterFolder(String xmlExporterFolder) { + this.xmlExporterFolder = xmlExporterFolder; + } +} diff --git a/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/config/settings/SnapshotMakerSettings.java b/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/config/settings/SnapshotMakerSettings.java new file mode 100644 index 000000000..c3c241601 --- /dev/null +++ b/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/config/settings/SnapshotMakerSettings.java @@ -0,0 +1,49 @@ +package ru.spcex.clearing.snapshot.maker.config.settings; + +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.context.annotation.PropertySource; +import org.springframework.stereotype.Component; +import ru.spcex.clearing.platform.messaging.config.element.KafkaConsumerSettings; +import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings; + +@Component +@PropertySource("file:${spring.config.location}/application.properties") +@ConfigurationProperties("snapshot-maker") +public class SnapshotMakerSettings { + private DatabaseSettings database; + private KafkaConsumerSettings kafkaConsumer; + private KafkaProducerSettings kafkaProducer; + private LocationSettings location; + + public DatabaseSettings getDatabase() { + return database; + } + + public void setDatabase(DatabaseSettings database) { + this.database = database; + } + + public KafkaConsumerSettings getKafkaConsumer() { + return kafkaConsumer; + } + + public void setKafkaConsumer(KafkaConsumerSettings kafkaConsumer) { + this.kafkaConsumer = kafkaConsumer; + } + + public KafkaProducerSettings getKafkaProducer() { + return kafkaProducer; + } + + public void setKafkaProducer(KafkaProducerSettings kafkaProducer) { + this.kafkaProducer = kafkaProducer; + } + + public LocationSettings getLocation() { + return location; + } + + public void setLocation(LocationSettings location) { + this.location = location; + } +} diff --git a/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/service/CommandService.java b/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/service/CommandService.java new file mode 100644 index 000000000..221be4499 --- /dev/null +++ b/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/service/CommandService.java @@ -0,0 +1,52 @@ +package ru.spcex.clearing.snapshot.maker.service; + +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.InitializingBean; +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.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.Status; + +@Service +public class CommandService extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final SessionService sessionService; + + public CommandService(Consumer kafkaQueue, + Producer kafkaResponseQueue, + SessionService sessionService) { + super(kafkaQueue, kafkaResponseQueue); + this.sessionService = sessionService; + } + + @Override + public void afterPropertiesSet() throws Exception { + callback(LauncherCommandRequest.class) + .setFunction(this::process) + .forDestination(Consts.LAUNCHER_NEW, callbacks::put); + + init(); + } + + private RequestInfoUpdate process(BaseRequest request) { + RequestInfoUpdate requestInfoUpdate = new RequestInfoUpdate(); + + if (request.getRequestPayload().getTaskName().equals("STSS")) { + requestInfoUpdate = sessionService.startSession(); + } else if (request.getRequestPayload().getTaskName().equals("SPSS")) { + requestInfoUpdate = sessionService.stopSession(); + } else { + log.info("Unknown task: {}", request.getRequestPayload().getTaskName()); + requestInfoUpdate.setStatus(Status.Error); + requestInfoUpdate.setMessage("Unknown task: " + request.getRequestPayload().getTaskName()); + } + + return requestInfoUpdate; + } +} diff --git a/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/service/SessionService.java b/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/service/SessionService.java new file mode 100644 index 000000000..e1cd71142 --- /dev/null +++ b/clearing-parent/snapshot-maker/src/main/java/ru/spcex/clearing/snapshot/maker/service/SessionService.java @@ -0,0 +1,167 @@ +package ru.spcex.clearing.snapshot.maker.service; + +import java.io.BufferedReader; +import java.io.IOException; +import java.io.InputStreamReader; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.Paths; +import java.time.LocalDateTime; +import java.time.format.DateTimeFormatter; +import java.util.List; +import java.util.regex.Matcher; +import java.util.regex.Pattern; +import org.apache.commons.io.FileUtils; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.stereotype.Service; +import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; +import ru.spcex.clearing.platform.messaging.service.Status; +import ru.spcex.clearing.snapshot.maker.config.settings.SnapshotMakerSettings; + +@Service +public class SessionService { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final SnapshotMakerSettings settings; + private final DateTimeFormatter dtf = DateTimeFormatter.ofPattern("yyyy-MM-dd'T'HH-mm-ss"); + private final static Pattern JDBC_URL_PATTERN = Pattern.compile("(?<=//)(.+?)(?=:):(.+?)(?=/)/(.+?)(?=\\?|$)"); + private Path currentPath; + + public SessionService(SnapshotMakerSettings settings) { + this.settings = settings; + } + + public RequestInfoUpdate startSession() { + RequestInfoUpdate requestInfoUpdate = new RequestInfoUpdate(); + + if (currentPath != null) { + log.info("You need to stop session first!"); + requestInfoUpdate.setStatus(Status.Error); + requestInfoUpdate.setMessage("You need to stop session first!"); + return requestInfoUpdate; + } + + log.info("Starting making snapshot"); + + try { + this.currentPath = Files.createDirectories(Paths.get(settings.getLocation().getLocalStore()).resolve(LocalDateTime.now().format(dtf))); + log.info("Session folder created: {}", this.currentPath); + } catch (IOException e) { + log.info("Can't create directory, error={}", e.getMessage()); + requestInfoUpdate.setStatus(Status.Error); + requestInfoUpdate.setMessage("Can't create directory, error = " + e.getMessage()); + return requestInfoUpdate; + } + dumpDatabase(requestInfoUpdate); + + if (requestInfoUpdate.getStatus() == Status.Error) { + return requestInfoUpdate; + } + + requestInfoUpdate.setStatus(Status.Success); + requestInfoUpdate.setMessage("Snapshot created successfully in folder " + this.currentPath.toAbsolutePath()); + return requestInfoUpdate; + } + + public RequestInfoUpdate stopSession() { + RequestInfoUpdate requestInfoUpdate = new RequestInfoUpdate(); + + if (currentPath == null) { + log.info("You need to start session first!"); + requestInfoUpdate.setStatus(Status.Error); + requestInfoUpdate.setMessage("You need to start session first!"); + return requestInfoUpdate; + } + + log.info("Stopping session"); + + copyDirectories(requestInfoUpdate); + this.currentPath = null; + + if (requestInfoUpdate.getStatus() == Status.Error) { + return requestInfoUpdate; + } + + requestInfoUpdate.setStatus(Status.Success); + requestInfoUpdate.setMessage("Session snapshot stopped successfully!"); + return requestInfoUpdate; + } + + private void dumpDatabase(RequestInfoUpdate requestInfoUpdate) { + Matcher matcher = JDBC_URL_PATTERN.matcher(settings.getDatabase().getJdbcUrl()); + if (!matcher.find()) { + requestInfoUpdate.setStatus(Status.Error); + requestInfoUpdate.setMessage("Database URL is not valid!"); + return; + } + + Path dumpLocation = this.currentPath.resolve("database_dump.sql"); + + ProcessBuilder pb = new ProcessBuilder(); + pb.environment().put("PGPASSWORD", settings.getDatabase().getPassword()); + + if (System.getProperty("os.name").contains("Windows")) { + pb.command("cmd", "/c", "pg_dump -U " + settings.getDatabase().getUsername() + " -h " + matcher.group(1) + " -p " + matcher.group(2) + " -d " + matcher.group(3) + " --column-inserts -f " + dumpLocation.toAbsolutePath()); + } else { + pb.command("/bin/sh", "-c", "pg_dump -U " + settings.getDatabase().getUsername() + " -h " + matcher.group(1) + " -p " + matcher.group(2) + " -d " + matcher.group(3) + " --column-inserts -f " + dumpLocation.toAbsolutePath()); + } + + try { + log.info("Executing command {}", pb.command()); + Process process = pb.start(); + StringBuilder out = new StringBuilder(); + BufferedReader reader = new BufferedReader(new InputStreamReader(process.getInputStream())); + + String line; + while ((line = reader.readLine()) != null) { + out.append(line).append(System.lineSeparator()); + } + int statusCode = process.waitFor(); + log.info("Status code={}, output {}", statusCode, out); + } catch (IOException | InterruptedException e) { + log.info("Can't dump database, error={}", e.getMessage()); + requestInfoUpdate.setStatus(Status.Error); + requestInfoUpdate.setMessage("Can't dump database, error = " + e.getMessage()); + } + } + + private void copyDirectories(RequestInfoUpdate requestInfoUpdate) { + List directoriesToCopy = List.of( + Paths.get(settings.getLocation().getDbfImporterFolder()), + Paths.get(settings.getLocation().getDbfExporterFolder()), + Paths.get(settings.getLocation().getSwtImporterFolder()), + Paths.get(settings.getLocation().getSwtExporterFolder()), + Paths.get(settings.getLocation().getXmlImporterFolder()), + Paths.get(settings.getLocation().getXmlExporterFolder()) + ); + + for (Path directory : directoriesToCopy) { + if (!Files.exists(directory) || !Files.isDirectory(directory)) { + log.info("{} is not a directory or does not exist.", directory.toAbsolutePath()); + requestInfoUpdate.setStatus(Status.Error); + requestInfoUpdate.setMessage(directory.toAbsolutePath() + " is not a directory or does not exist."); + return; + } + + Path target; + try { + target = Files.createDirectories(this.currentPath.resolve(directory.getParent().getFileName().toString()).resolve(directory.getFileName().toString())); + } catch (IOException e) { + log.info("Can't create directory, error={}", e.getMessage()); + requestInfoUpdate.setStatus(Status.Error); + requestInfoUpdate.setMessage("Can't create directory, error = " + e.getMessage()); + return; + } + + try { + log.info("Copying {} to {}", directory.toAbsolutePath(), target.toAbsolutePath()); + FileUtils.copyDirectory(directory.toFile(), target.toFile()); + } catch (IOException e) { + log.info("Can't copy directory {}, error={}", directory, e.getMessage()); + requestInfoUpdate.setStatus(Status.Error); + requestInfoUpdate.setMessage("Can't copy directory " + directory + ", error = " + e.getMessage()); + return; + } + } + } +} diff --git a/clearing-parent/snapshot-maker/src/main/resources/application.properties b/clearing-parent/snapshot-maker/src/main/resources/application.properties new file mode 100644 index 000000000..9bd5b6b5e --- /dev/null +++ b/clearing-parent/snapshot-maker/src/main/resources/application.properties @@ -0,0 +1,30 @@ +# Database settings +snapshot-maker.database.jdbcUrl=jdbc:postgresql://localhost:5433/clearing +snapshot-maker.database.username=clearing +snapshot-maker.database.password=Aa111111 + +# Locations setting +snapshot-maker.location.local-store=/mnt/c/Users/ivan/Desktop/snapshot-test/out +snapshot-maker.location.dbf-importer-folder=/mnt/c/Users/ivan/Desktop/snapshot-test/in/dbf/importer +snapshot-maker.location.dbf-exporter-folder=/mnt/c/Users/ivan/Desktop/snapshot-test/in/dbf/exporter +snapshot-maker.location.swt-importer-folder=/mnt/c/Users/ivan/Desktop/snapshot-test/in/swt/importer +snapshot-maker.location.swt-exporter-folder=/mnt/c/Users/ivan/Desktop/snapshot-test/in/swt/exporter +snapshot-maker.location.xml-importer-folder=/mnt/c/Users/ivan/Desktop/snapshot-test/in/xml/importer +snapshot-maker.location.xml-exporter-folder=/mnt/c/Users/ivan/Desktop/snapshot-test/in/xml/exporter + +# Kafka producer settings +snapshot-maker.kafka-producer.bootstrap-servers=localhost:9092 +snapshot-maker.kafka-producer.acks=all +snapshot-maker.kafka-producer.retries=0 +snapshot-maker.kafka-producer.batch-size=16384 +snapshot-maker.kafka-producer.linger-ms=1 +snapshot-maker.kafka-producer.buffer-memory=33554432 + +# Kafka consumer settings +snapshot-maker.kafka-consumer.bootstrap-servers=localhost:9092 +snapshot-maker.kafka-consumer.group-id=dev-group-clearing-service +snapshot-maker.kafka-consumer.enable-auto-commit=false +snapshot-maker.kafka-consumer.session-timeout-ms=30000 +snapshot-maker.kafka-consumer.auto-offset-reset=latest +snapshot-maker.kafka-consumer.linger-ms=1 +snapshot-maker.kafka-consumer.buffer-memory=33554432 \ No newline at end of file diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java index 50988fa24..3ef59e05a 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java @@ -59,6 +59,8 @@ public enum Task implements IEnumKey { finishBadSessions("CCLR"),//Завершение неудачных клиринговых сессий sendLim_LIMC("LIMC"), // Выгрузка в Торговую систему остатков по валюте (отправка lim)" makeFiles_MTCR("MTCR"), // Формирование файлов с МТКР + startSessionSnapshot("STSS"), // Старт формирования снэпшота сессии + stopSessionSnapshot("SPSS"), // Завершение формирования снэпшота сессии ; private final String key; diff --git a/pom.xml b/pom.xml index 15897e247..460a0a3c4 100644 --- a/pom.xml +++ b/pom.xml @@ -53,6 +53,7 @@ ${folder_root_clearing}/clearing-parent/registry-service ${folder_root_clearing}/clearing-parent/scheduler-service ${folder_root_clearing}/clearing-parent/gateway-api + ${folder_root_clearing}/clearing-parent/snapshot-maker 3.12.4 1.7.33