Commit 82ecd236 by huangfusuper

公共参数传递及聚合

parent f2b1247d
......@@ -321,9 +321,6 @@ public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
waitingRecord.setExtendedConfiguration(JSON.toJSONString(flowExtendedConfiguration));
waitingRecord.setFlowVersionName(flow.getVersionName());
String runId = IDGenerationStrategy.runIdGenerationStrategy(serverPort);
//将公共参数保存到redis
String redisKey = String.format(RedisKeyNameEnum.REDIS_PUBLIC_PARAM_KEY.getKeyName(), runId);
stringRedisTemplate.opsForValue().set(redisKey, publicParam);
waitingRecord.setRunId(runId);
waitingRecord.setFlowNodeCount(nodeCount);
......
package com.byit.annotations;
import java.lang.annotation.*;
/**
* 自定义类排序
*
* @author huangfu
* @date 2020年12月14日15:10:15
*/
@Documented
@Retention(RetentionPolicy.RUNTIME)
@Target(ElementType.TYPE)
public @interface MythRankOrder {
/**
* 级别
* @return 该类的级别
*/
int value() default -1;
}
package com.byit.dto;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
/**
* 策略包包裹
*
* @author huangfu
* @date 2020年12月14日15:14:56
*/
@Data
@AllArgsConstructor
@NoArgsConstructor
public class BeanStrategyPackage<T> {
private String beanName;
private T bean;
}
......@@ -8,7 +8,7 @@ public enum RedisKeyNameEnum {
/**
* 调度中心公共参数
*/
REDIS_PUBLIC_PARAM_KEY("myth:param:public:%s"),
REDIS_PUBLIC_PARAM_KEY("myth:param:public:%s%s"),
;
private final String keyName;
......
package com.byit.service.impl;
import com.byit.dto.BeanStrategyPackage;
import com.byit.enums.NodeNameEnum;
import com.byit.enums.RunRecordingEnum;
import com.byit.enums.ScheduleTypeEnum;
......@@ -12,6 +13,7 @@ import com.byit.model.Node;
import com.byit.model.RunRecording;
import com.byit.service.*;
import com.byit.strategy.InstanceRunTheLifeCycleCallback;
import com.byit.util.ClassSortUtil;
import com.byit.util.IDGenerationStrategy;
import com.byit.util.SpringUtil;
import lombok.extern.slf4j.Slf4j;
......@@ -106,12 +108,13 @@ public class RunNodeServiceImpl implements RunNodeServer, ApplicationEventPublis
Map<String, InstanceRunTheLifeCycleCallback> beansOfType = SpringUtil.getBeansOfType(InstanceRunTheLifeCycleCallback.class);
Set<Map.Entry<String, InstanceRunTheLifeCycleCallback>> entries = beansOfType.entrySet();
//回调生命周期
try {
for (Map.Entry<String, InstanceRunTheLifeCycleCallback> entry : entries) {
log.info("-------------【开始回调周期{}前置】---------------", entry.getKey());
RunRecording runRecording = entry.getValue().postProcessAfterInitialization(build);
//数据排序
List<BeanStrategyPackage<InstanceRunTheLifeCycleCallback>> beanStrategyPackages = ClassSortUtil.objectSort(beansOfType);
for (BeanStrategyPackage<InstanceRunTheLifeCycleCallback> beanStrategyPackage : beanStrategyPackages) {
log.info("-------------【开始回调周期{}前置】---------------", beanStrategyPackage.getBeanName());
RunRecording runRecording = beanStrategyPackage.getBean().postProcessAfterInitialization(build);
if (runRecording != null) {
build = runRecording;
}
......
package com.byit.service.mapservice.impl;
import cn.hutool.core.collection.CollectionUtil;
import com.byit.dto.BeanStrategyPackage;
import com.byit.dto.RunRecordingWrapped;
import com.byit.enums.*;
import com.byit.enums.task.RunResultEnum;
......@@ -10,6 +11,7 @@ import com.byit.model.*;
import com.byit.service.*;
import com.byit.service.mapservice.RunRecordingAndJobTaskService;
import com.byit.strategy.InstanceRunTheLifeCycleCallback;
import com.byit.util.ClassSortUtil;
import com.byit.util.SpringUtil;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
......@@ -269,16 +271,15 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
//开始执行生命周期 开始的
Map<String, InstanceRunTheLifeCycleCallback> beansOfType = SpringUtil.getBeansOfType(InstanceRunTheLifeCycleCallback.class);
Set<Map.Entry<String, InstanceRunTheLifeCycleCallback>> entries = beansOfType.entrySet();
try {
List<BeanStrategyPackage<InstanceRunTheLifeCycleCallback>> beanStrategyPackages = ClassSortUtil.objectSort(beansOfType);
//回调生命周期
for (Map.Entry<String, InstanceRunTheLifeCycleCallback> entry : entries) {
log.info("-------------【开始回调周期{}前置】---------------",entry.getKey());
RunRecording runRecordingLi = entry.getValue().postProcessAfterInitialization(runRecording);
for (BeanStrategyPackage<InstanceRunTheLifeCycleCallback> beanStrategyPackage : beanStrategyPackages) {
log.info("-------------【开始回调周期{}前置】---------------",beanStrategyPackage.getBeanName());
RunRecording runRecordingLi = beanStrategyPackage.getBean().postProcessAfterInitialization(runRecording);
if(runRecordingLi != null) {
runRecording = runRecordingLi;
}
}
}catch (Exception e) {
e.printStackTrace();
......
......@@ -2,6 +2,7 @@ package com.byit.service.mapservice.impl;
import cn.hutool.core.collection.CollectionUtil;
import com.byit.conf.RunRecordingThreadPool;
import com.byit.dto.BeanStrategyPackage;
import com.byit.dto.RunRecordingWrapped;
import com.byit.enums.EmailEnum;
import com.byit.enums.FlowPropertyEnum;
......@@ -12,6 +13,7 @@ import com.byit.model.*;
import com.byit.service.*;
import com.byit.service.mapservice.RunRecordingAndLogService;
import com.byit.strategy.InstanceRunTheLifeCycleCallback;
import com.byit.util.ClassSortUtil;
import com.byit.util.SpringUtil;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.BeanUtils;
......@@ -155,11 +157,11 @@ public class RunRecordingAndLogServiceImpl implements RunRecordingAndLogService,
RunRecordingThreadPool.RUN_RECORDING_THREAD_POOL.submit(() ->{
Map<String, InstanceRunTheLifeCycleCallback> beansOfType = SpringUtil.getBeansOfType(InstanceRunTheLifeCycleCallback.class);
Set<Map.Entry<String, InstanceRunTheLifeCycleCallback>> entries = beansOfType.entrySet();
List<BeanStrategyPackage<InstanceRunTheLifeCycleCallback>> beanStrategyPackages = ClassSortUtil.objectSort(beansOfType);
//回调生命周期
for (Map.Entry<String, InstanceRunTheLifeCycleCallback> entry : entries) {
log.info("-------------【开始回调周期{}后置】---------------",entry.getKey());
entry.getValue().postProcessBeforeInitialization(runRecording);
for (BeanStrategyPackage<InstanceRunTheLifeCycleCallback> beanStrategyPackage : beanStrategyPackages) {
log.info("-------------【开始回调周期{}后置】---------------",beanStrategyPackage.getBeanName());
beanStrategyPackage.getBean().postProcessBeforeInitialization(runRecording);
}
});
......
package com.byit.strategy.impl;
import cn.hutool.core.collection.CollectionUtil;
import com.alibaba.fastjson.JSON;
import com.byit.annotations.MythRankOrder;
import com.byit.call.TaskServer;
import com.byit.dto.common.CallbackDto;
import com.byit.dto.executor.ScriptParamAndPlaceholderDto;
import com.byit.dto.plugin.FlowExtendedConfiguration;
import com.byit.dto.web.ReturnResult;
import com.byit.enums.CallMethodEnum;
......@@ -23,7 +22,6 @@ import org.apache.commons.lang3.StringUtils;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;
import java.util.HashMap;
import java.util.Map;
/**
......@@ -33,9 +31,9 @@ import java.util.Map;
*/
@Component
@Slf4j
@MythRankOrder(1)
public class RpcCallbackInstanceRunTheLifeCycleCallback implements InstanceRunTheLifeCycleCallback {
public static final String EXTENDED_PARAMS_S = "extended:params:%s";
@TaskClient(timeout = 30000, callMethodEnum = CallMethodEnum.SYNCHRONIZE)
private TaskServer taskServer;
......@@ -65,54 +63,25 @@ public class RpcCallbackInstanceRunTheLifeCycleCallback implements InstanceRunTh
return runRecording;
}
CallbackDto callbackDto = new CallbackDto();
callbackDto.setData(flowExtendedConfiguration.getPublicParam());
PluginRpcRequestPacket param = buildPluginRpcRequestPacket(runRecording, callbackDto, rpcServerKey);
//给用户全部的公共参数
PluginRpcResponsePacket call = taskServer.call(param);
if (!call.isStatus()) {
throw new RuntimeException(call.getMsg());
}
Object result = call.getResult();
//获取客户方返回修改后的公共参数
ReturnResult<String> returnResult = (ReturnResult<String>) result;
//将对应的参数转换为公共参数 修改后的公共参数
Map<String, String> lifeCyclePublicParam = JSON.parseObject(returnResult.getContent(), Map.class);
//放置到对应的redis里面
stringRedisTemplate.opsForValue().set(String.format(EXTENDED_PARAMS_S, runRecording.getRunId()), JSON.toJSONString(lifeCyclePublicParam));
String publicParam = flowExtendedConfiguration.getPublicParam();
Map<String, String> publicParamMap = JSON.parseObject(publicParam, Map.class);
publicParamMap.putAll(lifeCyclePublicParam);
String publicParamMapValue = JSON.toJSONString(publicParamMap);
flowExtendedConfiguration.setPublicParam(publicParamMapValue);
String format = String.format(RedisKeyNameEnum.REDIS_PUBLIC_PARAM_KEY.getKeyName(), runRecording.getRunId());
//放置到对应的redis里面
stringRedisTemplate.opsForValue().set(format, publicParamMapValue);
runRecording.setExtendedConfiguration(JSON.toJSONString(flowExtendedConfiguration));
//修改后的公共参数重新保存到实例里面以及redis里面
flowExtendedConfiguration.setPublicParam(JSON.toJSONString(lifeCyclePublicParam));
//保存到redis
stringRedisTemplate.opsForValue().set(String.format(RedisKeyNameEnum.REDIS_PUBLIC_PARAM_KEY.getKeyName(), runRecording.getRunId(), flowExtendedConfiguration.getFlowEmbedLogo()), JSON.toJSONString(flowExtendedConfiguration));
return runRecording;
}
/**
* 构建 PluginRpcRequestPacket
*
* @param runRecording 执行实例
* @return {@link PluginRpcRequestPacket}
*/
private PluginRpcRequestPacket buildPluginRpcRequestPacket(RunRecording runRecording, CallbackDto callbackDto, String serverKey) {
PluginRpcRequestPacket rpcRequestPacket = new PluginRpcRequestPacket();
rpcRequestPacket.setCallMethodEnum(CallMethodEnum.SYNCHRONIZE);
rpcRequestPacket.setJobName(serverKey);
CommunicationParam param = new CommunicationParam();
param.setRunId(runRecording.getRunId());
param.setHasMakeUp(String.valueOf(runRecording.getScheduleType()));
callbackDto.setFlowName(runRecording.getFlowName());
callbackDto.setVersionName(runRecording.getFlowVersionName());
Workspace workspaceServiceOneById = workspaceService.findOneById(runRecording.getWorkspaceId());
callbackDto.setWorkspaceName(workspaceServiceOneById.getWorkspaceName());
callbackDto.setRunId(runRecording.getRunId());
param.setBody(JSON.toJSONString(callbackDto));
rpcRequestPacket.setParam(param);
return rpcRequestPacket;
}
/**
* 后置处理器
......@@ -121,7 +90,7 @@ public class RpcCallbackInstanceRunTheLifeCycleCallback implements InstanceRunTh
*/
@Override
public void postProcessBeforeInitialization(RunRecording runRecording) {
//获取执行扩展参数
String extendedConfiguration = runRecording.getExtendedConfiguration();
FlowExtendedConfiguration flowExtendedConfiguration = JSON.parseObject(extendedConfiguration, FlowExtendedConfiguration.class);
if (flowExtendedConfiguration == null) {
......@@ -131,46 +100,49 @@ public class RpcCallbackInstanceRunTheLifeCycleCallback implements InstanceRunTh
if (StringUtils.isBlank(rpcEndServerKey)) {
return;
}
String extendedParam = stringRedisTemplate.opsForValue().get(String.format(EXTENDED_PARAMS_S, runRecording.getRunId()));
ScriptParamAndPlaceholderDto scriptParamAndPlaceholderDto = JSON.parseObject(extendedParam, ScriptParamAndPlaceholderDto.class);
//获取本实例的公共参数
String publicParam = flowExtendedConfiguration.getPublicParam();
Map<String, String> callData = new HashMap<>(2);
if (StringUtils.isNoneBlank(publicParam)) {
Map<String, String> stringStringMap = (Map<String, String>) JSON.parse(publicParam);
if (CollectionUtil.isNotEmpty(stringStringMap)) {
callData.putAll(stringStringMap);
}
}
if (scriptParamAndPlaceholderDto != null) {
Map<String, String> param = scriptParamAndPlaceholderDto.getParam();
Map<String, String> placeholder = scriptParamAndPlaceholderDto.getPlaceholder();
if (CollectionUtil.isNotEmpty(param)) {
callData.putAll(param);
}
if (CollectionUtil.isNotEmpty(placeholder)) {
callData.putAll(placeholder);
}
}
CallbackDto callbackDto = new CallbackDto();
if (RunRecordingEnum.RUN_FLOW_SUCCESS.getCode().equals(runRecording.getFlowRunResult()) || RunRecordingEnum.RUN_FLOW_RE_SUCCESS.getCode().equals(runRecording.getFlowRunResult())) {
callbackDto.setRunRecordingStatus(true);
}
callbackDto.setData(JSON.toJSONString(callData));
callbackDto.setData(publicParam);
PluginRpcRequestPacket rpcRequestPacket = buildPluginRpcRequestPacket(runRecording, callbackDto, rpcEndServerKey);
try {
PluginRpcResponsePacket call = taskServer.call(rpcRequestPacket);
log.info("回调{}通讯成功!", call);
} catch (Exception e) {
log.info("回调通讯失败!");
} finally {
stringRedisTemplate.delete(String.format(EXTENDED_PARAMS_S, runRecording.getRunId()));
}
}
/**
* 构建 PluginRpcRequestPacket
*
* @param runRecording 执行实例
* @return {@link PluginRpcRequestPacket}
*/
private PluginRpcRequestPacket buildPluginRpcRequestPacket(RunRecording runRecording, CallbackDto callbackDto, String serverKey) {
PluginRpcRequestPacket rpcRequestPacket = new PluginRpcRequestPacket();
rpcRequestPacket.setCallMethodEnum(CallMethodEnum.SYNCHRONIZE);
rpcRequestPacket.setJobName(serverKey);
CommunicationParam param = new CommunicationParam();
param.setRunId(runRecording.getRunId());
param.setHasMakeUp(String.valueOf(runRecording.getScheduleType()));
callbackDto.setFlowName(runRecording.getFlowName());
callbackDto.setVersionName(runRecording.getFlowVersionName());
Workspace workspaceServiceOneById = workspaceService.findOneById(runRecording.getWorkspaceId());
callbackDto.setWorkspaceName(workspaceServiceOneById.getWorkspaceName());
callbackDto.setRunId(runRecording.getRunId());
param.setBody(JSON.toJSONString(callbackDto));
rpcRequestPacket.setParam(param);
return rpcRequestPacket;
}
}
package com.byit.strategy.impl;
import com.alibaba.fastjson.JSON;
import com.byit.annotations.MythRankOrder;
import com.byit.dto.plugin.FlowExtendedConfiguration;
import com.byit.enums.RedisKeyNameEnum;
import com.byit.model.RunRecording;
import com.byit.strategy.InstanceRunTheLifeCycleCallback;
......@@ -15,6 +18,7 @@ import org.springframework.stereotype.Component;
*/
@Component
@Slf4j
@MythRankOrder(1000)
public class RunRecordingRunRedisRemoveTheLifeCycleCallback implements InstanceRunTheLifeCycleCallback {
private final StringRedisTemplate stringRedisTemplate;
......@@ -29,10 +33,12 @@ public class RunRecordingRunRedisRemoveTheLifeCycleCallback implements InstanceR
*/
@Override
public void postProcessBeforeInitialization(RunRecording runRecording) {
String runId = runRecording.getRunId();
String extendedConfiguration = runRecording.getExtendedConfiguration();
FlowExtendedConfiguration flowExtendedConfiguration = JSON.parseObject(extendedConfiguration, FlowExtendedConfiguration.class);
//删除本次实例的公共参数
String runId = runRecording.getRunId();
String publicRedisName = String.format(RedisKeyNameEnum.REDIS_PUBLIC_PARAM_KEY.getKeyName(), runId);
String publicRedisName = String.format(RedisKeyNameEnum.REDIS_PUBLIC_PARAM_KEY.getKeyName(), runId,flowExtendedConfiguration.getFlowEmbedLogo());
stringRedisTemplate.delete(publicRedisName);
}
}
package com.byit.strategy.impl;
import cn.hutool.core.collection.CollectionUtil;
import com.alibaba.fastjson.JSON;
import com.byit.annotations.MythRankOrder;
import com.byit.dto.plugin.FlowExtendedConfiguration;
import com.byit.enums.RedisKeyNameEnum;
import com.byit.model.RunRecording;
import com.byit.strategy.InstanceRunTheLifeCycleCallback;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;
import java.util.Map;
/**
* 工作流参数汇总
* key的生成策略为: myth:param:public:runId:mark
*
* @author huangfu
* @date 2020年12月14日16:28:35
*/
@Component
@Slf4j
@MythRankOrder(5)
public class WorkflowCommonParameterAggregation implements InstanceRunTheLifeCycleCallback {
private final StringRedisTemplate stringRedisTemplate;
public WorkflowCommonParameterAggregation(StringRedisTemplate stringRedisTemplate) {
this.stringRedisTemplate = stringRedisTemplate;
}
@Override
public RunRecording postProcessAfterInitialization(RunRecording runRecording) {
log.info("-------{}实例参数聚合-------", runRecording);
String extendedConfiguration = runRecording.getExtendedConfiguration();
//获取扩展配置
FlowExtendedConfiguration flowExtendedConfiguration = JSON.parseObject(extendedConfiguration, FlowExtendedConfiguration.class);
//获取工作流层级
String flowEmbedLogo = flowExtendedConfiguration.getFlowEmbedLogo();
String publicRedisKey = String.format(RedisKeyNameEnum.REDIS_PUBLIC_PARAM_KEY.getKeyName(), runRecording.getRunId(), flowEmbedLogo);
String redisKeyPrefix = String.format("myth:param:public:%s", runRecording.getRunId());
String publicParamStr = flowExtendedConfiguration.getPublicParam();
Map<String, String> publicMap = JSON.parseObject(publicParamStr, Map.class);
//获取工作流的层级关系
String leven = publicRedisKey.replace(redisKeyPrefix, "");
while (StringUtils.isNoneBlank(leven)) {
polymerizationPublicParam(publicMap, leven, runRecording.getRunId());
leven = leven.substring(0, leven.lastIndexOf(":"));
}
//转换聚合后的公共参数
String publicParam = JSON.toJSONString(publicMap);
//将聚合后的参数放置到扩展配置
flowExtendedConfiguration.setPublicParam(publicParam);
String jsonString = JSON.toJSONString(flowExtendedConfiguration);
//将扩展配置更新进redis
stringRedisTemplate.opsForValue().set(publicRedisKey, jsonString);
//更新该实例的扩展配置
runRecording.setExtendedConfiguration(jsonString);
return runRecording;
}
/**
* 聚合参数
*
* @param publicParam 公共参数
*/
private void polymerizationPublicParam(Map<String, String> publicParam, String leven, String runId) {
String redisKey = String.format(RedisKeyNameEnum.REDIS_PUBLIC_PARAM_KEY.getKeyName(), runId, leven);
String thisRedisKeyData = stringRedisTemplate.opsForValue().get(redisKey);
Map<String, String> thisRedisKeyMap = JSON.parseObject(thisRedisKeyData, Map.class);
if (CollectionUtil.isNotEmpty(thisRedisKeyMap)) {
thisRedisKeyMap.forEach((key, value) -> {
if (!publicParam.containsKey(key)) {
publicParam.put(key, value);
}
});
}
}
}
package com.byit.thread.helper;
import com.byit.conf.RunRecordingThreadPool;
import com.byit.dto.BeanStrategyPackage;
import com.byit.enums.EmailEnum;
import com.byit.enums.NodeRunStatusPropertyEnum;
import com.byit.enums.RunRecordingEnum;
......@@ -13,6 +14,7 @@ import com.byit.service.JobTaskRunLogService;
import com.byit.service.RunRecordingService;
import com.byit.strategy.InstanceRunTheLifeCycleCallback;
import com.byit.thread.BaseThreadRunHelper;
import com.byit.util.ClassSortUtil;
import com.byit.util.SpringUtil;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
......@@ -81,11 +83,11 @@ public class ClosingExampleThreadRunHelper extends BaseThreadRunHelper implement
RunRecordingThreadPool.RUN_RECORDING_THREAD_POOL.submit(() ->{
Map<String, InstanceRunTheLifeCycleCallback> beansOfType = SpringUtil.getBeansOfType(InstanceRunTheLifeCycleCallback.class);
Set<Map.Entry<String, InstanceRunTheLifeCycleCallback>> entries = beansOfType.entrySet();
List<BeanStrategyPackage<InstanceRunTheLifeCycleCallback>> beanStrategyPackages = ClassSortUtil.objectSort(beansOfType);
//回调生命周期
for (Map.Entry<String, InstanceRunTheLifeCycleCallback> entry : entries) {
log.info("-------------【开始回调周期{}后置】---------------",entry.getKey());
entry.getValue().postProcessBeforeInitialization(runRecording);
for (BeanStrategyPackage<InstanceRunTheLifeCycleCallback> beanStrategyPackage : beanStrategyPackages) {
log.info("-------------【开始回调周期{}后置】---------------",beanStrategyPackage.getBeanName());
beanStrategyPackage.getBean().postProcessBeforeInitialization(runRecording);
}
});
......
package com.byit.util;
import com.byit.annotations.MythRankOrder;
import com.byit.dto.BeanStrategyPackage;
import lombok.Data;
import java.util.*;
/**
* 类排序
*
* @author huangfu
* @date 2020年12月14日15:11:55
*/
public class ClassSortUtil {
/**
* 对Spring返回的bean进行数据排序
* @param beanMap beanMap
* @param <T> 泛型
* @return 排序后的集合
*/
public static <T> List<BeanStrategyPackage<T>> objectSort(Map<String,T> beanMap){
List<BeanStrategyPackage<T>> beanStrategyPackages = new ArrayList<>(8);
beanMap.forEach((beanName,bean) ->{
if(bean.getClass().isAnnotationPresent(MythRankOrder.class)){
throw new RuntimeException(String.format("%s,不存在排序注解,请联系调度官方团队!",beanName));
}
beanStrategyPackages.add(new BeanStrategyPackage<>(beanName,bean));
});
beanStrategyPackages.sort((o1, o2) -> {
Object bean1 = o1.getBean();
Object bean2 = o2.getBean();
MythRankOrder mythRankOrder1 = bean1.getClass().getAnnotation(MythRankOrder.class);
MythRankOrder mythRankOrder2 = bean2.getClass().getAnnotation(MythRankOrder.class);
int value1 = mythRankOrder1.value();
int value2 = mythRankOrder2.value();
return Integer.compare(value1, value2);
});
return beanStrategyPackages;
}
}
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment