This commit is contained in:
parent
700e2442cc
commit
1442ab6c50
3 changed files with 58 additions and 3 deletions
|
|
@ -0,0 +1,39 @@
|
|||
package ru.spcex.clearing.csv.importer.config;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.kafka.core.KafkaTemplate;
|
||||
import org.springframework.kafka.core.ProducerFactory;
|
||||
import ru.spcex.clearing.csv.importer.config.settings.CsvImporterSettings;
|
||||
import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory;
|
||||
import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings;
|
||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||
import ru.spcex.platform.imdg.api.ImdgId;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
|
||||
@Configuration
|
||||
public class KafkaConfig {
|
||||
@Bean
|
||||
public ProducerFactory<String, Object> pf(CsvImporterSettings settings) {
|
||||
KafkaProducerSettings kafkaSettings = settings.getKafkaProducer();
|
||||
return KafkaProducerFactory.producerFactory(kafkaSettings);
|
||||
}
|
||||
|
||||
@Bean("kafkaTemplate")
|
||||
public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> pf) {
|
||||
return new KafkaTemplate<>(pf);
|
||||
}
|
||||
|
||||
@Autowired
|
||||
@Bean("kafkaSenderWithoutRequestInfo")
|
||||
public KafkaSender kafkaSenderWithoutRequestInfo(KafkaTemplate<String, Object> kafkaTemplate, ImdgProvider imdgProvider) {
|
||||
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
|
||||
return KafkaSender
|
||||
.setup()
|
||||
.setKafkaTemplate(kafkaTemplate)
|
||||
.idGenerator(imdgIdGenerator::nextId)
|
||||
.saveRequestInfo(false)
|
||||
.build();
|
||||
}
|
||||
}
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package ru.spcex.clearing.csv.importer.component;
|
||||
package ru.spcex.clearing.csv.importer.service;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
|
|
@ -10,12 +10,16 @@ import org.springframework.scheduling.annotation.EnableScheduling;
|
|||
import org.springframework.scheduling.annotation.Scheduled;
|
||||
import org.springframework.stereotype.Component;
|
||||
import ru.clearing.classes.statics.data.misc.SOrders;
|
||||
import ru.spcex.clearing.csv.importer.component.FileService;
|
||||
import ru.spcex.clearing.csv.importer.component.tro.data.SOrder;
|
||||
import ru.spcex.clearing.csv.importer.component.tro.parser.TroParserType;
|
||||
import ru.spcex.clearing.csv.importer.component.tro.processor.TroProcessor;
|
||||
import ru.spcex.clearing.csv.importer.component.tro.processor.TroProcessorProvider;
|
||||
import ru.spcex.clearing.csv.importer.config.settings.CsvImporterSettings;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.order.OrderCurrencyStatusUpdateRequest;
|
||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||
import ru.spcex.platform.imdg.api.Imdg;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
import ru.spcex.platform.utils.log.ExceptionUtils;
|
||||
|
|
@ -29,13 +33,15 @@ public class SOrdersScheduler {
|
|||
private final FileService fileService;
|
||||
private final TroProcessor<SOrder> processor;
|
||||
private final Imdg<SOrders> sOrderImdg;
|
||||
private final KafkaSender kafkaSender;
|
||||
|
||||
@Autowired
|
||||
public SOrdersScheduler(CsvImporterSettings settings, FileService fileService, TroProcessorProvider provider, ImdgProvider imdgProvider) {
|
||||
public SOrdersScheduler(CsvImporterSettings settings, FileService fileService, TroProcessorProvider provider, ImdgProvider imdgProvider, KafkaSender kafkaSender) {
|
||||
this.importDir = Path.of(settings.getSrcDir());
|
||||
this.fileService = fileService;
|
||||
this.processor = provider.createProcessor(TroParserType.SOrder, SOrder.class);
|
||||
this.sOrderImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SOrders, SOrders.class);
|
||||
this.kafkaSender = kafkaSender;
|
||||
}
|
||||
|
||||
@Scheduled(fixedDelayString = "${csv-importer.poll-period:5000}")
|
||||
|
|
@ -74,10 +80,20 @@ public class SOrdersScheduler {
|
|||
sOrderImdg.insert(order);
|
||||
}
|
||||
}
|
||||
sendToDefaultManagement(sOrders);
|
||||
fileService.moveSuccess(path);
|
||||
} catch (Throwable e) {
|
||||
log.error("error for {}: {}", path, ExceptionUtils.getStackTrace(e));
|
||||
fileService.moveError(path);
|
||||
}
|
||||
}
|
||||
|
||||
private void sendToDefaultManagement(List<SOrder> sOrders) {
|
||||
String destination = Consts.ORDER_CURRENCY_STATUS_UPDATE_IMPORTER;
|
||||
OrderCurrencyStatusUpdateRequest req = new OrderCurrencyStatusUpdateRequest();
|
||||
sOrders.forEach(order -> req.addId(order.getTransId()));
|
||||
log.info("sending request to topic {}...", destination);
|
||||
kafkaSender.sendRequestToQueue(destination, req);
|
||||
log.info("kafka request sent");
|
||||
}
|
||||
}
|
||||
|
|
@ -27,6 +27,6 @@ class TroProcessorTest {
|
|||
);
|
||||
TroProcessor<SOrder> processor = provider.createProcessor(TroParserType.SOrder, SOrder.class);
|
||||
List<SOrder> sOrders = processor.readFile(path);
|
||||
Assertions.assertTrue(sOrders.size() > 0, "read SOrder from tro file");
|
||||
Assertions.assertFalse(sOrders.isEmpty(), "read SOrder's from tro file");
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue