refactoring SDF*executors

This commit is contained in:
etreschenkov 2023-05-24 13:41:27 +03:00
parent d3069db920
commit 699587c8d4
4 changed files with 18 additions and 17 deletions

View file

@ -61,6 +61,7 @@ public class StatementService extends QueueConsumer implements InitializingBean
.forDestination(Consts.STATEMENT_PROCESS, callbacks::put); .forDestination(Consts.STATEMENT_PROCESS, callbacks::put);
init(); init();
} }
private void process(BaseRequest<StatementRequest> systemRequest) { private void process(BaseRequest<StatementRequest> systemRequest) {
StatementRequest statementRequest = systemRequest.getRequestPayload(); StatementRequest statementRequest = systemRequest.getRequestPayload();
Collection<? extends SpcexObjectBase> sdfGroup; Collection<? extends SpcexObjectBase> sdfGroup;
@ -75,19 +76,16 @@ public class StatementService extends QueueConsumer implements InitializingBean
.map(part -> sdfImdg.getSingleObjectByID(part.getSdfId())) .map(part -> sdfImdg.getSingleObjectByID(part.getSdfId()))
.collect(Collectors.toList()); .collect(Collectors.toList());
} }
// sdf01Group = sdf01Group
// .stream()
// .sorted(Comparator.comparing(SpcexObjectBase::getId))
// .collect(Collectors.toList());
AbstractExecutor service = executorsMap.get(table); AbstractExecutor service = executorsMap.get(table);
Result res = service.execute(sdfGroup, statementRequest); Result res = service.execute(sdfGroup, statementRequest);
if (res.getAccountRequests().size() == 0) { if (res.getAccountRequests().size() != 0) {
kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests()));
} else if (service.isNeedToSendCommandToExport()) {
ExportToFileRequest exportRequest = new ExportToFileRequest(); ExportToFileRequest exportRequest = new ExportToFileRequest();
exportRequest.setSdfGroupId(res.getGenerationId()); exportRequest.setSdfGroupId(res.getGenerationId());
exportRequest.setNameOfTable(service.exportTableName()); exportRequest.setNameOfTable(service.exportTableName());
kafkaSender.sendRequestToQueue(Consts.EXPORT_PROCESS, exportRequest); kafkaSender.sendRequestToQueue(Consts.EXPORT_PROCESS, exportRequest);
} else {
kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests()));
} }
} }
@ -97,12 +95,4 @@ public class StatementService extends QueueConsumer implements InitializingBean
r.setAccounts(accountRequests); r.setAccounts(accountRequests);
return r; return r;
} }
private Imdg<?> defineImdgByTableName(SdfTable table){
switch (table){
case SDF_01 -> imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class);
case SDF_57 -> imdgProvider.getImdg(IMDGDistributedNames.Map_SDf57, SDf57.class);
}
return null;
}
} }

View file

@ -8,4 +8,5 @@ import java.util.Collection;
public abstract class AbstractExecutor<T> { public abstract class AbstractExecutor<T> {
public abstract Result execute(Collection<T> sdf, StatementRequest statementRequest); public abstract Result execute(Collection<T> sdf, StatementRequest statementRequest);
public abstract String exportTableName(); public abstract String exportTableName();
public abstract boolean isNeedToSendCommandToExport();
} }

View file

@ -64,6 +64,11 @@ public class Sdf01Executor extends AbstractExecutor<SDf01> {
return "DF-02"; return "DF-02";
} }
@Override
public boolean isNeedToSendCommandToExport() {
return true;
}
public Result execute(Collection<SDf01> sdf, StatementRequest statementRequest) { public Result execute(Collection<SDf01> sdf, StatementRequest statementRequest) {
Result result = new Result(); Result result = new Result();
Long generationIdForGroup = imdgProvider.getImdgIdGenerator().nextId(); Long generationIdForGroup = imdgProvider.getImdgIdGenerator().nextId();

View file

@ -65,8 +65,13 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
@Override @Override
public String exportTableName() { public String exportTableName() {
return "DF-57"; return null;
} //no need... }
@Override
public boolean isNeedToSendCommandToExport() {
return false;
}//no need...
//V - Изменение statement по sDf57 //V - Изменение statement по sDf57
// //