|
@@ -328,6 +328,7 @@ public class SimulationProjectServiceImpl implements SimulationProjectService {
|
|
|
kafkaParam.setScenePackageId(po.getScene());
|
|
|
kafkaParam.setMaxSimulationTime(po.getMaxSimulationTime());
|
|
|
kafkaParam.setParallelism(Integer.valueOf(po.getParallelism()));
|
|
|
+ kafkaParam.setType(DictConstants.PROJECT_TYPE_MANUAL);
|
|
|
KafkaParameter kafkaParameter = new KafkaParameter();
|
|
|
kafkaParameter.setTopic(ProjectConstants.RUN_TASK_TOPIC);
|
|
|
String data = JsonUtil.beanToJson(kafkaParam);
|
|
@@ -338,6 +339,7 @@ public class SimulationProjectServiceImpl implements SimulationProjectService {
|
|
|
private void projectStopToKafka(SimulationManualProjectPo po) throws JsonProcessingException{
|
|
|
SimulationManualProjectKafkaParam kafkaParam = new SimulationManualProjectKafkaParam();
|
|
|
kafkaParam.setProjectId(po.getId());
|
|
|
+ kafkaParam.setType(DictConstants.PROJECT_TYPE_MANUAL);
|
|
|
KafkaParameter kafkaParameter = new KafkaParameter();
|
|
|
kafkaParameter.setTopic(ProjectConstants.STOP_TASK_TOPPIC);
|
|
|
String data = JsonUtil.beanToJson(kafkaParam);
|
|
@@ -350,8 +352,9 @@ public class SimulationProjectServiceImpl implements SimulationProjectService {
|
|
|
private void autoProjectStopToKafka(SimulationAutomaticSubProjectPo po) throws JsonProcessingException{
|
|
|
SimulationManualProjectKafkaParam kafkaParam = new SimulationManualProjectKafkaParam();
|
|
|
kafkaParam.setProjectId(po.getId());
|
|
|
+ kafkaParam.setType(DictConstants.PROJECT_TYPE_AUTO_SUB);
|
|
|
KafkaParameter kafkaParameter = new KafkaParameter();
|
|
|
- kafkaParameter.setTopic(ProjectConstants.STOP_AUTO_PROJECT);
|
|
|
+ kafkaParameter.setTopic(ProjectConstants.STOP_TASK_TOPPIC);
|
|
|
String data = JsonUtil.beanToJson(kafkaParam);
|
|
|
kafkaParameter.setData(data);
|
|
|
log.info("推送自动项目中止消息到kafka:"+data);
|
|
@@ -4030,8 +4033,9 @@ public class SimulationProjectServiceImpl implements SimulationProjectService {
|
|
|
kafkaParam.setScenePackageId(po.getScene());
|
|
|
kafkaParam.setMaxSimulationTime(po.getMaxSimulationTime());
|
|
|
kafkaParam.setParallelism(Integer.valueOf(po.getParallelism()));
|
|
|
+ kafkaParam.setType(DictConstants.PROJECT_TYPE_AUTO_SUB);
|
|
|
KafkaParameter kafkaParameter = new KafkaParameter();
|
|
|
- kafkaParameter.setTopic(ProjectConstants.AUTO_PROJECT);
|
|
|
+ kafkaParameter.setTopic(ProjectConstants.RUN_TASK_TOPIC);
|
|
|
String data = JsonUtil.beanToJson(kafkaParam);
|
|
|
kafkaParameter.setData(data);
|
|
|
log.info("自动运行项目推送消息到kafka:"+data);
|