scheduler set up

This commit is contained in:
aalehin 2022-12-01 12:27:57 +03:00
parent ca4483fe6f
commit 6d14c43994
10 changed files with 424 additions and 75 deletions

View file

@ -7,73 +7,75 @@ import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.misc.KeyRate;
import ru.clearing.classes.statics.data.scheduler.Scheduler;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.KeyRateNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.KeyRateUpdateRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.SchedulerNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.SchedulerUpdateRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.platform.enumeration.Status;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
@Service
public class SchedulerService extends QueueConsumer implements InitializingBean {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<KeyRate> keyRateMap;
private final Imdg<Scheduler> schedulerMap;
@Autowired
public SchedulerService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider) {
super(kafkaQueue, kafkaProducer);
this.keyRateMap = imdgProvider.getImdg(IMDGDistributedNames.Map_KeyRate, KeyRate.class);
this.schedulerMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Scheduler, Scheduler.class);
}
@Override
public void afterPropertiesSet() {
callback(KeyRateNewRequest.class)
.setConsumer(this::newKeyRate)
.forDestination(Consts.DESTINATION_KEY_RATE_NEW, callbacks::put);
callback(KeyRateUpdateRequest.class)
.setConsumer(this::updateKeyRate)
.forDestination(Consts.DESTINATION_KEY_RATE_UPDATE, callbacks::put);
callback(SchedulerNewRequest.class)
.setConsumer(this::newScheduler)
.forDestination(Consts.DESTINATION_SCHEDULER_NEW, callbacks::put);
callback(SchedulerUpdateRequest.class)
.setConsumer(this::updateScheduler)
.forDestination(Consts.DESTINATION_SCHEDULER_UPDATE, callbacks::put);
callback(CommonDeleteRequest.class)
.setConsumer(this::deleteKeyRate)
.forDestination(Consts.DESTINATION_KEY_RATE_DELETE, callbacks::put);
.setConsumer(this::deleteScheduler)
.forDestination(Consts.DESTINATION_SCHEDULER_DELETE, callbacks::put);
init();
}
private void newKeyRate(BaseRequest<KeyRateNewRequest> userRequest) {
KeyRateNewRequest req = userRequest.getRequestPayload();
log.debug("MoneyMarketSecurityNewRequest received");
KeyRate keyRate = new KeyRate();
keyRate.setEndDate(req.getEndDate());
keyRate.setDocument(req.getDocument());
keyRate.setRate(req.getKeyRate());
keyRate.setStartDate(req.getStartDate());
keyRate.setWorkflowStatus(Status.Active.getKey());
keyRateMap.insert(keyRate);
log.debug("successfully processed, new id {}", keyRate.getId());
private void newScheduler(BaseRequest<SchedulerNewRequest> userRequest) {
SchedulerNewRequest req = userRequest.getRequestPayload();
log.debug("SchedulerNewRequest received");
Scheduler scheduler = new Scheduler();
scheduler.setTask(req.getTask());
scheduler.setTaskTime(req.getTaskTime());
scheduler.setClearingDate(req.getClearingDate());
scheduler.setMarket(req.getMarket());
scheduler.setTaskStatus(req.getTaskStatus());
scheduler.setSecurityId(req.getSecurityId());
schedulerMap.insert(scheduler);
log.debug("successfully processed, new id {}", scheduler.getId());
}
private void updateKeyRate(BaseRequest<KeyRateUpdateRequest> userRequest) {
KeyRateUpdateRequest req = userRequest.getRequestPayload();
private void updateScheduler(BaseRequest<SchedulerUpdateRequest> userRequest) {
SchedulerUpdateRequest req = userRequest.getRequestPayload();
log.debug("MoneyMarketSecurityUpdateRequest received id = {}", req.getId());
KeyRate keyRate = keyRateMap.getSingleObjectByID(req.getId());
keyRate.setEndDate(req.getEndDate());
keyRate.setDocument(req.getDocument());
keyRate.setRate(req.getKeyRate());
keyRate.setStartDate(req.getStartDate());
keyRateMap.update(keyRate);
Scheduler scheduler = schedulerMap.getSingleObjectByID(req.getId());
scheduler.setTask(req.getTask());
scheduler.setTaskTime(req.getTaskTime());
scheduler.setClearingDate(req.getClearingDate());
scheduler.setMarket(req.getMarket());
scheduler.setTaskStatus(req.getTaskStatus());
scheduler.setSecurityId(req.getSecurityId());
schedulerMap.update(scheduler);
}
private void deleteKeyRate(BaseRequest<CommonDeleteRequest> userRequest) {
private void deleteScheduler(BaseRequest<CommonDeleteRequest> userRequest) {
CommonDeleteRequest req = userRequest.getRequestPayload();
log.debug("CommonDeleteRequest received id = {}", req.getId());
KeyRate keyRate = keyRateMap.getSingleObjectByID(req.getId());
keyRateMap.delete(keyRate);
Scheduler scheduler = schedulerMap.getSingleObjectByID(req.getId());
schedulerMap.delete(scheduler);
}
}

View file

@ -0,0 +1,75 @@
package ru.spcex.clearing.scheduler.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.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.scheduler.Timetable;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.TimetableNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.TimetableUpdateRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
@Service
public class TimetableService extends QueueConsumer implements InitializingBean {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<Timetable> timetableMap;
@Autowired
public TimetableService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider) {
super(kafkaQueue, kafkaProducer);
this.timetableMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Timetable, Timetable.class);
}
@Override
public void afterPropertiesSet() {
callback(TimetableNewRequest.class)
.setConsumer(this::newTimetable)
.forDestination(Consts.DESTINATION_TIMETABLE_NEW, callbacks::put);
callback(TimetableUpdateRequest.class)
.setConsumer(this::updateTimetable)
.forDestination(Consts.DESTINATION_TIMETABLE_UPDATE, callbacks::put);
callback(CommonDeleteRequest.class)
.setConsumer(this::deleteTimetable)
.forDestination(Consts.DESTINATION_TIMETABLE_DELETE, callbacks::put);
init();
}
private void newTimetable(BaseRequest<TimetableNewRequest> userRequest) {
TimetableNewRequest req = userRequest.getRequestPayload();
log.debug("TimetableNewRequest received");
Timetable timetable = new Timetable();
timetable.setTask(req.getTask());
timetable.setTaskTime(req.getTaskTime());
timetable.setTaskStatus(req.getTaskStatus());
timetableMap.insert(timetable);
log.debug("successfully processed, new id {}", timetable.getId());
}
private void updateTimetable(BaseRequest<TimetableUpdateRequest> userRequest) {
TimetableUpdateRequest req = userRequest.getRequestPayload();
log.debug("TimetableNewRequest received");
Timetable timetable = timetableMap.getSingleObjectByID(req.getId());
timetable.setTask(req.getTask());
timetable.setTaskTime(req.getTaskTime());
timetable.setTaskStatus(req.getTaskStatus());
timetableMap.update(timetable);
log.debug("successfully processed, new id {}", timetable.getId());
}
private void deleteTimetable(BaseRequest<CommonDeleteRequest> userRequest) {
CommonDeleteRequest req = userRequest.getRequestPayload();
log.debug("CommonDeleteRequest received id = {}", req.getId());
Timetable timetable = timetableMap.getSingleObjectByID(req.getId());
timetableMap.delete(timetable);
}
}

View file

@ -0,0 +1,75 @@
package ru.spcex.clearing.scheduler.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.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.scheduler.TradingCalendar;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.TradingCalendarNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.TradingCalendarUpdateRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
@Service
public class TradingCalendarService extends QueueConsumer implements InitializingBean {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<TradingCalendar> tradingCalendarMap;
@Autowired
public TradingCalendarService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider) {
super(kafkaQueue, kafkaProducer);
this.tradingCalendarMap = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingCalendar, TradingCalendar.class);
}
@Override
public void afterPropertiesSet() {
callback(TradingCalendarNewRequest.class)
.setConsumer(this::newTradingCalendar)
.forDestination(Consts.DESTINATION_TRADING_CALENDAR_NEW, callbacks::put);
callback(TradingCalendarUpdateRequest.class)
.setConsumer(this::updateTradingCalendar)
.forDestination(Consts.DESTINATION_TRADING_CALENDAR_UPDATE, callbacks::put);
callback(CommonDeleteRequest.class)
.setConsumer(this::deleteTradingCalendar)
.forDestination(Consts.DESTINATION_TRADING_CALENDAR_DELETE, callbacks::put);
init();
}
private void newTradingCalendar(BaseRequest<TradingCalendarNewRequest> userRequest) {
TradingCalendarNewRequest req = userRequest.getRequestPayload();
log.debug("TradingCalendarNewRequest received");
TradingCalendar tradingCalendar = new TradingCalendar();
tradingCalendar.setClearingDate(req.getClearingDate());
tradingCalendar.setCompanyId(req.getCompanyId());
tradingCalendar.setTradingStatus(req.getTradingStatus());
tradingCalendarMap.insert(tradingCalendar);
log.debug("successfully processed, new id {}", tradingCalendar.getId());
}
private void updateTradingCalendar(BaseRequest<TradingCalendarUpdateRequest> userRequest) {
TradingCalendarUpdateRequest req = userRequest.getRequestPayload();
log.debug("TradingCalendarUpdateRequest received");
TradingCalendar tradingCalendar = tradingCalendarMap.getSingleObjectByID(req.getId());
tradingCalendar.setClearingDate(req.getClearingDate());
tradingCalendar.setCompanyId(req.getCompanyId());
tradingCalendar.setTradingStatus(req.getTradingStatus());
tradingCalendarMap.update(tradingCalendar);
log.debug("successfully processed, new id {}", tradingCalendar.getId());
}
private void deleteTradingCalendar(BaseRequest<CommonDeleteRequest> userRequest) {
CommonDeleteRequest req = userRequest.getRequestPayload();
log.debug("CommonDeleteRequest received id = {}", req.getId());
TradingCalendar tradingCalendar = tradingCalendarMap.getSingleObjectByID(req.getId());
tradingCalendarMap.delete(tradingCalendar);
}
}

View file

@ -9,6 +9,14 @@ public interface Consts {
String DESTINATION_KEY_RATE_UPDATE = "key-rate-update";
String DESTINATION_KEY_RATE_DELETE = "key-rate-delete";
String DESTINATION_TIMETABLE_NEW = "timetable-new";
String DESTINATION_TIMETABLE_UPDATE = "timetable-update";
String DESTINATION_TIMETABLE_DELETE = "timetable-delete";
String DESTINATION_TRADING_CALENDAR_NEW = "timetable-new";
String DESTINATION_TRADING_CALENDAR_UPDATE = "timetable-update";
String DESTINATION_TRADING_CALENDAR_DELETE = "timetable-delete";
String DESTINATION_SCHEDULER_NEW = "scheduler-new";
String DESTINATION_SCHEDULER_UPDATE = "scheduler-update";
String DESTINATION_SCHEDULER_DELETE = "scheduler-delete";

View file

@ -4,23 +4,26 @@ import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.databind.annotation.JsonDeserialize;
import com.fasterxml.jackson.databind.annotation.JsonSerialize;
import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalDateDeserializer;
import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalTimeDeserializer;
import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalDateSerializer;
import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalTimeSerializer;
import java.time.LocalDate;
import java.time.LocalTime;
public class SchedulerNewRequest {
@JsonProperty
@JsonSerialize(using = LocalDateSerializer.class)
@JsonDeserialize(using = LocalDateDeserializer.class)
public LocalTime taskTime;
@JsonProperty
public String task;
@JsonSerialize(using = LocalTimeSerializer.class)
@JsonDeserialize(using = LocalTimeDeserializer.class)
@JsonProperty
public LocalTime taskTime;
@JsonSerialize(using = LocalDateSerializer.class)
@JsonDeserialize(using = LocalDateDeserializer.class)
@JsonProperty
public LocalDate clearingDate;
@JsonProperty
@ -32,13 +35,12 @@ public class SchedulerNewRequest {
@JsonProperty
public Long securityId;
public SchedulerNewRequest(LocalTime taskTime, String task, LocalDate clearingDate, String market, String taskStatus, Long securityId) {
this.taskTime = taskTime;
public String getTask() {
return task;
}
public void setTask(String task) {
this.task = task;
this.clearingDate = clearingDate;
this.market = market;
this.taskStatus = taskStatus;
this.securityId = securityId;
}
public LocalTime getTaskTime() {
@ -49,14 +51,6 @@ public class SchedulerNewRequest {
this.taskTime = taskTime;
}
public String getTask() {
return task;
}
public void setTask(String task) {
this.task = task;
}
public LocalDate getClearingDate() {
return clearingDate;
}

View file

@ -12,30 +12,33 @@ import java.time.LocalDate;
import java.time.LocalTime;
public class SchedulerUpdateRequest {
@JsonProperty
private String task;
@JsonProperty
public Long id;
@JsonProperty
public String task;
@JsonSerialize(using = LocalTimeSerializer.class)
@JsonDeserialize(using = LocalTimeDeserializer.class)
@JsonProperty
private LocalTime taskTime;
public LocalTime taskTime;
@JsonSerialize(using = LocalDateSerializer.class)
@JsonDeserialize(using = LocalDateDeserializer.class)
@JsonProperty
private LocalDate clearingDate;
public LocalDate clearingDate;
@JsonProperty
private String market;
public String market;
@JsonProperty
private String taskStatus;
public String taskStatus;
@JsonProperty
private Long securityId;
public Long securityId;
@JsonProperty
private Long id;
public Long getId() {
return id;
}
public void setId(Long id) {
this.id = id;
}
public String getTask() {
return task;
@ -84,12 +87,4 @@ public class SchedulerUpdateRequest {
public void setSecurityId(Long securityId) {
this.securityId = securityId;
}
public Long getId() {
return id;
}
public void setId(Long id) {
this.id = id;
}
}

View file

@ -0,0 +1,45 @@
package ru.spcex.clearing.platform.messaging.domain.cud.schedule;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.databind.annotation.JsonDeserialize;
import com.fasterxml.jackson.databind.annotation.JsonSerialize;
import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalTimeDeserializer;
import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalTimeSerializer;
import java.time.LocalTime;
public class TimetableNewRequest {
@JsonProperty
public String task;
@JsonSerialize(using = LocalTimeSerializer.class)
@JsonDeserialize(using = LocalTimeDeserializer.class)
@JsonProperty
public LocalTime taskTime;
@JsonProperty
public String taskStatus;
public String getTask() {
return task;
}
public void setTask(String task) {
this.task = task;
}
public LocalTime getTaskTime() {
return taskTime;
}
public void setTaskTime(LocalTime taskTime) {
this.taskTime = taskTime;
}
public String getTaskStatus() {
return taskStatus;
}
public void setTaskStatus(String taskStatus) {
this.taskStatus = taskStatus;
}
}

View file

@ -0,0 +1,56 @@
package ru.spcex.clearing.platform.messaging.domain.cud.schedule;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.databind.annotation.JsonDeserialize;
import com.fasterxml.jackson.databind.annotation.JsonSerialize;
import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalTimeDeserializer;
import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalTimeSerializer;
import java.time.LocalTime;
public class TimetableUpdateRequest {
@JsonProperty
public Long id;
@JsonProperty
public String task;
@JsonSerialize(using = LocalTimeSerializer.class)
@JsonDeserialize(using = LocalTimeDeserializer.class)
@JsonProperty
public LocalTime taskTime;
@JsonProperty
public String taskStatus;
public String getTask() {
return task;
}
public void setTask(String task) {
this.task = task;
}
public LocalTime getTaskTime() {
return taskTime;
}
public void setTaskTime(LocalTime taskTime) {
this.taskTime = taskTime;
}
public String getTaskStatus() {
return taskStatus;
}
public void setTaskStatus(String taskStatus) {
this.taskStatus = taskStatus;
}
public Long getId() {
return id;
}
public void setId(Long id) {
this.id = id;
}
}

View file

@ -0,0 +1,44 @@
package ru.spcex.clearing.platform.messaging.domain.cud.schedule;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.databind.annotation.JsonDeserialize;
import com.fasterxml.jackson.databind.annotation.JsonSerialize;
import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalDateDeserializer;
import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalDateSerializer;
import java.time.LocalDate;
public class TradingCalendarNewRequest {
@JsonSerialize(using = LocalDateSerializer.class)
@JsonDeserialize(using = LocalDateDeserializer.class)
@JsonProperty
public LocalDate clearingDate;
@JsonProperty
public Long companyId;
@JsonProperty
public String tradingStatus;
public LocalDate getClearingDate() {
return clearingDate;
}
public void setClearingDate(LocalDate clearingDate) {
this.clearingDate = clearingDate;
}
public Long getCompanyId() {
return companyId;
}
public void setCompanyId(Long companyId) {
this.companyId = companyId;
}
public String getTradingStatus() {
return tradingStatus;
}
public void setTradingStatus(String tradingStatus) {
this.tradingStatus = tradingStatus;
}
}

View file

@ -0,0 +1,55 @@
package ru.spcex.clearing.platform.messaging.domain.cud.schedule;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.databind.annotation.JsonDeserialize;
import com.fasterxml.jackson.databind.annotation.JsonSerialize;
import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalDateDeserializer;
import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalDateSerializer;
import java.time.LocalDate;
public class TradingCalendarUpdateRequest {
@JsonProperty
public Long id;
@JsonSerialize(using = LocalDateSerializer.class)
@JsonDeserialize(using = LocalDateDeserializer.class)
@JsonProperty
public LocalDate clearingDate;
@JsonProperty
public Long companyId;
@JsonProperty
public String tradingStatus;
public Long getId() {
return id;
}
public void setId(Long id) {
this.id = id;
}
public LocalDate getClearingDate() {
return clearingDate;
}
public void setClearingDate(LocalDate clearingDate) {
this.clearingDate = clearingDate;
}
public Long getCompanyId() {
return companyId;
}
public void setCompanyId(Long companyId) {
this.companyId = companyId;
}
public String getTradingStatus() {
return tradingStatus;
}
public void setTradingStatus(String tradingStatus) {
this.tradingStatus = tradingStatus;
}
}