This commit is contained in:
parent
2e46a484ff
commit
15cad890c9
2 changed files with 20 additions and 20 deletions
|
|
@ -12,27 +12,26 @@ import ru.spcex.platform.imdg.api.Imdg;
|
|||
import ru.spcex.platform.imdg.api.ImdgId;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
|
||||
import java.util.function.Supplier;
|
||||
|
||||
@Configuration
|
||||
public class KafkaConfig {
|
||||
|
||||
@Autowired
|
||||
@Bean
|
||||
public Producer<String, Object> createProducer(ImportDBFServiceSettings settings) {
|
||||
return KafkaProducerFactory.producer(settings.getKafkaProducer());
|
||||
}
|
||||
|
||||
@Autowired
|
||||
@Bean
|
||||
public KafkaSender kafkaSender(Producer<String, Object> kafkaProducer, ImdgProvider imdgProvider) {
|
||||
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
|
||||
return KafkaSender
|
||||
.setup()
|
||||
.producer(kafkaProducer)
|
||||
.idGenerator(imdgIdGenerator::nextId)
|
||||
.imdgProvider(s -> {
|
||||
Imdg<RequestInfo> imdg = imdgProvider.getImdg(s, RequestInfo.class);
|
||||
return imdg::insert;
|
||||
})
|
||||
.build();
|
||||
public Supplier<KafkaSender> kafkaSender(ImportDBFServiceSettings settings, ImdgProvider imdgProvider) {
|
||||
return () -> {
|
||||
Producer<String, Object> kafkaProducer = KafkaProducerFactory.producer(settings.getKafkaProducer());
|
||||
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
|
||||
return KafkaSender
|
||||
.setup()
|
||||
.producer(kafkaProducer)
|
||||
.idGenerator(imdgIdGenerator::nextId)
|
||||
.imdgProvider(s -> {
|
||||
Imdg<RequestInfo> imdg = imdgProvider.getImdg(s, RequestInfo.class);
|
||||
return imdg::insert;
|
||||
})
|
||||
.build();
|
||||
};
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -18,6 +18,7 @@ import java.io.IOException;
|
|||
import java.io.InputStream;
|
||||
import java.nio.charset.Charset;
|
||||
import java.util.Map;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
/**
|
||||
* Заливка проверенных данных в базу
|
||||
|
|
@ -27,11 +28,11 @@ public class ImportToDB extends Stage {
|
|||
private final AProperties properties;
|
||||
private final HazelcastService hazelcastService;
|
||||
private final Map<ETable, AbstractTable> mappingEnumTableObjectTable;
|
||||
private final KafkaSender kafka;
|
||||
private final Supplier<KafkaSender> kafka;
|
||||
|
||||
public ImportToDB(@Qualifier("dbfImporterProperties") AProperties properties,
|
||||
HazelcastService hazelcastService, @Qualifier("mapOfTable") Map<ETable, AbstractTable> mappingEnumTableObjectTable,
|
||||
KafkaSender kafka) {
|
||||
Supplier<KafkaSender> kafka) {
|
||||
this.properties = properties;
|
||||
this.hazelcastService = hazelcastService;
|
||||
this.mappingEnumTableObjectTable = mappingEnumTableObjectTable;
|
||||
|
|
@ -70,6 +71,6 @@ public class ImportToDB extends Stage {
|
|||
private void sendStatementRequest(Long fileId) {
|
||||
StatementRequest statementRequest = new StatementRequest();
|
||||
statementRequest.setSdf01GroupId(fileId);
|
||||
kafka.sendRequestToQueue(Consts.STATEMENT_PROCESS, statementRequest);
|
||||
kafka.get().sendRequestToQueue(Consts.STATEMENT_PROCESS, statementRequest);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue