dolphinscheduler-筆記2

springboot集成dolphinscheduler

說明

為了避免對DolphinScheduler產生過度依賴,實踐中通常不會全面采用其內置的所有任務節點類型。相反,會選擇性地利用DolphinScheduler的HTTP任務節點功能,以此作為工作流執行管理的橋梁,對接并驅動自有項目的業務流程。這種策略不僅確保了流程編排的靈活性與擴展性,還有效減少了對外部調度系統的深度綁定,從而在提升項目自洽能力的同時,保持了良好的系統間解耦。

簡而言之,我們傾向于僅采納DolphinScheduler中的HTTP任務節點,作為調度機制的一部分,來促進我們內部項目工作流的自動化執行。這樣做既能享受DolphinScheduler帶來的調度便利,又避免了全盤接受其所有組件所帶來的潛在風險,實現了更為穩健、可控的項目管理方案。

代碼實現

為了優化與DolphinScheduler的集成,以下是三個關鍵配置類的概述,它們旨在通過初始化接口實現項目及租戶信息的同步通知。值得注意的是,為了確保數據一致性和高效通信,你的Spring Boot應用所使用的數據庫應與DolphinScheduler共享同一數據源。這一策略不僅簡化了數據管理,還促進了實時狀態更新,增強了系統的整體協調性。

簡而言之,我們精心設計了三組配置規則,允許我們的Spring Boot項目無縫對接DolphinScheduler平臺。通過這些配置,項目和租戶的動態變化能夠及時反映到DolphinScheduler中,前提是兩者共用一個數據庫實例。這種架構決策不僅優化了資源分配,還促進了跨系統間的緊密協作,為后續的業務拓展奠定了堅實的基礎。

package cn.com.lyb.data.dev.init;import cn.com.lyb.common.security.annotation.InnerRequest;
import cn.com.lyb.core.exception.BizException;
import cn.com.lyb.core.web.request.Response;
import cn.com.lyb.core.web.request.ResultWrap;
import cn.com.lyb.data.dev.web.config.DolphinschedulerConfig;
import io.swagger.annotations.Api;
import io.swagger.annotations.ApiOperation;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;import javax.sql.DataSource;
import java.sql.Connection;
import java.sql.SQLException;
import java.sql.Statement;
import java.util.Arrays;
import java.util.List;@Api(tags = "初始化dolphinscheduler數據庫信息")
@RestController
@RequestMapping("/data-dev")
public class InitializePlugin {@Autowiredprivate DataSource dataSource;@Autowiredprivate DolphinschedulerConfig dolphinschedulerConfig;@GetMapping("/init")@ApiOperation("部署全新的環境可以用此接口,否則會報錯")@InnerRequestpublic Response init(){Connection connection = null;try{connection = dataSource.getConnection();// 項目表(xgov如果集成項目的話,這個sql不能再執行,配置也需要改)String projectSql = "INSERT INTO `dolphinscheduler`.`t_ds_project` (`id`, `name`, `code`, `description`, `user_id`, `flag`, `create_time`, `update_time`) VALUES (1, 'lyb', '" + dolphinschedulerConfig.getProjectCode() + "', '', 1, 1, '2024-06-13 02:49:43', '2024-06-13 02:49:43');";String tokenSql = "INSERT INTO `dolphinscheduler`.`t_ds_access_token` (`id`, `user_id`, `token`, `expire_time`, `create_time`, `update_time`) VALUES (1, 1, '"+dolphinschedulerConfig.getDsdToken()+"', '2039-12-30 10:51:26', '2024-06-13 02:50:37', '2024-06-18 10:00:13');";String tenantSql = "INSERT INTO `dolphinscheduler`.`t_ds_tenant` (`id`, `tenant_code`, `description`, `queue_id`, `create_time`, `update_time`) VALUES (1, 'default', '', 1, '2024-06-13 02:50:20', '2024-06-13 02:50:20');";String userSql = "UPDATE `dolphinscheduler`.`t_ds_user` SET `user_name` = 'admin', `user_password` = '470b9934942620215ad1cb3ac2d48497', `user_type` = 0, `email` = 'xxx@qq.com', `phone` = '', `tenant_id` = 1, `create_time` = '2024-06-12 10:23:37', `update_time` = '2024-06-18 09:59:52', `queue` = '', `state` = 1, `time_zone` = 'Asia/Shanghai' WHERE `id` = 1;";List<String> sqlList = Arrays.asList(projectSql, tenantSql, tokenSql, userSql);connection.setAutoCommit(false);Statement statement = connection.createStatement();for (String sql : sqlList) {statement.addBatch(sql);}statement.executeBatch();connection.commit();statement.close();}catch (Exception e){throw new BizException("初始化報錯");}finally {if(connection != null) {try {connection.close();} catch (SQLException e) {e.printStackTrace();}}}return ResultWrap.ok();}
}
package cn.com.lyb.data.dev.web.config;import org.apache.http.client.config.RequestConfig;
import org.apache.http.impl.client.CloseableHttpClient;
import org.apache.http.impl.client.HttpClients;
import org.apache.http.impl.conn.PoolingHttpClientConnectionManager;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.http.client.HttpComponentsClientHttpRequestFactory;
import org.springframework.web.client.RestTemplate;@Configuration
public class RestTemplateConfig {@Value("${xgov.template.connectTimeout}")private int connectTimeout;@Value("${xgov.template.socketTimeout}")private int socketTimeout;@Beanpublic RestTemplate restTemplate() {return new RestTemplate(httpRequestFactory());}@Beanpublic HttpComponentsClientHttpRequestFactory httpRequestFactory() {PoolingHttpClientConnectionManager connectionManager = new PoolingHttpClientConnectionManager();connectionManager.setMaxTotal(200); // 最大連接數connectionManager.setDefaultMaxPerRoute(20); // 每個路由默認的最大連接數RequestConfig requestConfig = RequestConfig.custom().setConnectTimeout(connectTimeout) // 連接超時時間.setSocketTimeout(socketTimeout) // 讀取超時時間.setConnectionRequestTimeout(5000) // 從連接池獲取連接的超時時間.build();CloseableHttpClient httpClient = HttpClients.custom().setConnectionManager(connectionManager).setDefaultRequestConfig(requestConfig).build();return new HttpComponentsClientHttpRequestFactory(httpClient);}
}
package cn.com.dev.data.dev.web.config;import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Configuration;@Configuration
public class DolphinschedulerConfig {@Value("${lyb.dolphinscheduler.server.username}")private String dsdUsername;@Value("${lyb.dolphinscheduler.server.token}")private String dsdToken;@Value("${lyb.dolphinscheduler.server.url}")private String dsdUrl;@Value("${lyb.dolphinscheduler.server.porjectCode}")private String projectCode;@Value("${lyb.dolphinscheduler.server.tenantCode}")private String tenantCode;public String getDsdUsername() {return dsdUsername;}public void setDsdUsername(String dsdUsername) {this.dsdUsername = dsdUsername;}public String getDsdToken() {return dsdToken;}public void setDsdToken(String dsdToken) {this.dsdToken = dsdToken;}public String getDsdUrl() {return dsdUrl;}public void setDsdUrl(String dsdUrl) {this.dsdUrl = dsdUrl;}public String getProjectCode() {return projectCode;}public void setProjectCode(String projectCode) {this.projectCode = projectCode;}public String getTenantCode() {return tenantCode;}public void setTenantCode(String tenantCode) {this.tenantCode = tenantCode;}}

構建一個調用類,該類全面集成了與DolphinScheduler接口的交互邏輯,為我們的應用提供了一層抽象。對于涉及具體業務邏輯的數據封裝細節,此處將不再贅述,旨在保持代碼的清晰度與通用性。

簡言之,我們設計了一個專門的類來處理所有與DolphinScheduler的API調用,確保了業務核心邏輯的獨立性和可維護性。這一封裝策略使得代碼庫更加整潔,同時也提升了開發效率和系統的整體健壯性。

通過這種方式,我們不僅隔離了與外部服務的直接交互,還簡化了業務邏輯的實現,使其更加專注于核心功能,而非調度系統的細節。這樣的架構設計,有助于團隊成員快速理解系統架構,同時也便于未來的功能擴展和系統維護。

package cn.com.lyb.data.dev.web.dolphinscheduler.service;import cn.com.lyb.common.redis.service.RedisService;
import cn.com.lyb.core.exception.BizException;
import cn.com.lyb.data.dev.enums.ProcessExecutionTypeEnum;
import cn.com.lyb.data.dev.enums.TaskExecutionStatus;
import cn.com.lyb.data.dev.enums.WorkflowExecutionStatus;
import cn.com.lyb.data.dev.web.config.DolphinschedulerConfig;
import cn.com.lyb.data.dev.workflow.entity.delphinscheduler.TaskDefinition;
import cn.com.lyb.data.dev.workflow.entity.vo.GanttTaskVO;
import cn.com.lyb.data.dev.workflow.entity.vo.ProcessInstanceVO;
import cn.com.lyb.data.dev.workflow.entity.vo.ResponseTaskLog;
import cn.com.lyb.data.dev.workflow.entity.vo.TaskInstanceVO;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject;
import com.github.pagehelper.PageInfo;
import org.apache.commons.lang3.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.HttpEntity;
import org.springframework.http.HttpHeaders;
import org.springframework.http.HttpMethod;
import org.springframework.http.ResponseEntity;
import org.springframework.stereotype.Service;
import org.springframework.util.LinkedMultiValueMap;
import org.springframework.util.MultiValueMap;
import org.springframework.web.client.RestTemplate;import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.TimeUnit;@Service
public class DolphinschedulerService {private static final Logger logger = LoggerFactory.getLogger(DolphinschedulerService.class);private static final String CONTENT_TYPE = "application/x-www-form-urlencoded; charset=utf-8";private static final String X_REQUESTED_WITH = "XMLHttpRequest";@Autowiredprivate RestTemplate restTemplate;@Autowiredprivate DolphinschedulerConfig dolphinschedulerConfig;private static final Boolean DSD_SUCCESS = true;@Autowiredprivate RedisService redisService;private static final String DSD_SESSION_KEY = "DSD_SESSION_KEY";private static final String SUCCESS = "success";private static final String MSG = "msg";// 此種方法適用登錄后獲取SESSION設置到header里面private static final String SESSION_ID = "sessionId";private static final String TOKEN = "token";private static final String DATA = "data";/*** 登錄,返回 sessionId*/public String login() {String sessionValue = redisService.getCacheObject(DSD_SESSION_KEY);if (StringUtils.isNotBlank(sessionValue)) {return sessionValue;}//SSLUtil.turnOffSslChecking();HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);LinkedMultiValueMap<String, String> linkedMultiValueMap = new LinkedMultiValueMap<String, String>();linkedMultiValueMap.add("userName", dolphinschedulerConfig.getDsdUsername());//linkedMultiValueMap.add("userPassword", dolphinschedulerConfig.getDsdPassword());HttpEntity<MultiValueMap<String, String>> httpEntity = new HttpEntity<>(linkedMultiValueMap, hs);String url = dolphinschedulerConfig.getDsdUrl() + "/login";JSONObject resultJSON = doPostForObject(url, httpEntity);JSONObject data = (JSONObject) resultJSON.get(DATA);String sessionId = data.get(SESSION_ID).toString();redisService.setCacheObject(DSD_SESSION_KEY, sessionId, 23L, TimeUnit.HOURS);return sessionId;}/*** 創建項目** @return*/public void createProject(String projectName, String description) {// 如果是https登錄可以使用該方法//SSLUtil.turnOffSslChecking();HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());LinkedMultiValueMap<String, String> linkedMultiValueMap = new LinkedMultiValueMap<String, String>();linkedMultiValueMap.add("projectName", projectName);linkedMultiValueMap.add("description", description);linkedMultiValueMap.add("userName", dolphinschedulerConfig.getDsdUsername());HttpEntity<MultiValueMap<String, String>> httpEntity = new HttpEntity<>(linkedMultiValueMap, hs);logger.info("Azkaban請求信息:" + httpEntity.toString());String url = dolphinschedulerConfig.getDsdUrl() + "/projects";doPostForObject(url, httpEntity);}/*** 項目列表** @param pageNo* @param pageSize* @param searchVal*/public Object projectsPage(Integer pageNo, Integer pageSize, String searchVal) {//SSLUtil.turnOffSslChecking();HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());HttpEntity<String> entity = new HttpEntity<String>(hs);String url;if (StringUtils.isNotBlank(searchVal)) {url = dolphinschedulerConfig.getDsdUrl() + "/projects?pageNo=" + pageNo + "&pageSize=" + pageSize + "&searchVal=" + searchVal;} else {url = dolphinschedulerConfig.getDsdUrl() + "/projects?pageNo=" + pageNo + "&pageSize=" + pageSize;}JSONObject resultJSON = doGetForObject(url, entity);return resultJSON.get(DATA);}/*** 修改項目** @param code* @param projectName* @param description* @return*/public Object updateProjects(String code, String projectName, String description) {//SSLUtil.turnOffSslChecking();HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());LinkedMultiValueMap<String, String> linkedMultiValueMap = new LinkedMultiValueMap<String, String>();linkedMultiValueMap.add("projectName", projectName);linkedMultiValueMap.add("description", description);linkedMultiValueMap.add("userName", dolphinschedulerConfig.getDsdUsername());HttpEntity<MultiValueMap<String, String>> httpEntity = new HttpEntity<>(linkedMultiValueMap, hs);String url = dolphinschedulerConfig.getDsdUrl() + "/projects/" + code;JSONObject resultJSON = doPutForObject(url, httpEntity);return resultJSON.get(DATA);}/*** 刪除項目** @param code* @return*/public Object delProjects(String code) {//SSLUtil.turnOffSslChecking();HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());HttpEntity<MultiValueMap<String, String>> httpEntity = new HttpEntity<>(hs);String url = dolphinschedulerConfig.getDsdUrl() + "/projects/" + code;JSONObject resultJSON = doDeleteForObject(url, httpEntity);return resultJSON.get(DATA);}public JSONObject doPostForObject(String url, HttpEntity httpEntity) {logger.info("調用url:{}", url);try {ResponseEntity<String> exchange = restTemplate.exchange(url, HttpMethod.POST, httpEntity, String.class);String result = exchange.getBody();logger.info("post類型接口調用返回信息:{}", result);return getJsonObject(result);} catch (BizException be) {throw be;} catch (Exception e) {logger.error("post類型接口調用失敗:{}", e);throw new BizException(e.getMessage());}}public JSONObject doGetForObject(String url, HttpEntity httpEntity) {logger.info("調用url:{}", url);try {ResponseEntity<String> exchange = restTemplate.exchange(url, HttpMethod.GET, httpEntity, String.class);String result = exchange.getBody();logger.info("get類型接口調用返回信息:{}", result);return getJsonObject(result);} catch (BizException be) {throw be;} catch (Exception e) {logger.error("get類型接口調用失敗:{}", e);throw new BizException(e.getMessage());}}public JSONObject doPutForObject(String url, HttpEntity httpEntity) {logger.info("調用url:{}", url);try {ResponseEntity<String> exchange = restTemplate.exchange(url, HttpMethod.PUT, httpEntity, String.class);String result = exchange.getBody();logger.info("put類型接口調用返回信息:{}", result);return getJsonObject(result);} catch (BizException be) {throw be;} catch (Exception e) {logger.error("put類型接口調用失敗:{}", e);throw new BizException(e.getMessage());}}public JSONObject doDeleteForObject(String url, HttpEntity httpEntity) {logger.info("調用url:{}", url);try {ResponseEntity<String> exchange = restTemplate.exchange(url, HttpMethod.DELETE, httpEntity, String.class);String result = exchange.getBody();logger.info("delete類型接口調用返回信息:{}", result);return getJsonObject(result);} catch (BizException be) {throw be;} catch (Exception e) {logger.error("delete類型接口調用失敗:{}", e);throw new BizException(e.getMessage());}}private static JSONObject getJsonObject(String result) {JSONObject resultJSON = JSON.parseObject(result);if (!DSD_SUCCESS.equals(resultJSON.get(SUCCESS))) {logger.error("調用結果返回異常:{}" + result);Integer code = (Integer) resultJSON.get("code");if(code.intValue() == 50019){throw new BizException("流程節點間存在循環依賴");}else if(code.intValue() == 50036){throw new BizException("工作流任務關系參數錯誤");} else {throw new BizException(resultJSON.get(MSG).toString());}}return resultJSON;}/*** 創建工作流** @param name* @param description* @param globalParams* @param locations* @param timeout* @param taskRelationJson* @param taskDefinitionJson* @param otherParamsJson* @param executionType* @return 3.2.0 版本*/public Long createWorkFlow320(String name, String description, String globalParams, String locations, int timeout,String taskRelationJson, String taskDefinitionJson, String otherParamsJson,ProcessExecutionTypeEnum executionType) {//SSLUtil.turnOffSslChecking();HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());LinkedMultiValueMap<String, Object> linkedMultiValueMap = new LinkedMultiValueMap<String, Object>();linkedMultiValueMap.add("description", description);linkedMultiValueMap.add("name", name);linkedMultiValueMap.add("taskDefinitionJson", taskDefinitionJson);linkedMultiValueMap.add("taskRelationJson", taskRelationJson);linkedMultiValueMap.add("timeout", timeout);linkedMultiValueMap.add("executionType", executionType);linkedMultiValueMap.add("otherParamsJson", otherParamsJson);linkedMultiValueMap.add("globalParams", globalParams);HttpEntity<MultiValueMap<String, Object>> httpEntity = new HttpEntity<>(linkedMultiValueMap, hs);String url = dolphinschedulerConfig.getDsdUrl() + "/projects/" + dolphinschedulerConfig.getProjectCode() + "/process-definition";JSONObject resultJSON = doPostForObject(url, httpEntity);JSONObject data = (JSONObject) resultJSON.get(DATA);Long code = (Long) data.get("code");return code;}/*** 查詢工作流列表** @param pageNo* @param pageSize* @param searchVal* @return*/public Object selectFlowPage(Integer pageNo, Integer pageSize, String searchVal) {//SSLUtil.turnOffSslChecking();HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());HttpEntity<String> entity = new HttpEntity<String>(hs);String url;if (StringUtils.isNotBlank(searchVal)) {url = dolphinschedulerConfig.getDsdUrl() + "/projects/" + dolphinschedulerConfig.getProjectCode() + "/process-definition?pageNo=" + pageNo + "&pageSize=" + pageSize + "&searchVal=" + searchVal;} else {url = dolphinschedulerConfig.getDsdUrl() + "/projects/" + dolphinschedulerConfig.getProjectCode() + "/process-definition?pageNo=" + pageNo + "&pageSize=" + pageSize;}JSONObject resultJSON = doGetForObject(url, entity);return resultJSON.get(DATA);}public Object selectOneFlow(String code) {//SSLUtil.turnOffSslChecking();HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());HttpEntity<String> entity = new HttpEntity<String>(hs);String url = dolphinschedulerConfig.getDsdUrl() + "/projects/" + dolphinschedulerConfig.getProjectCode() + "/process-definition/" + code;JSONObject resultJSON = doGetForObject(url, entity);return resultJSON.get(DATA);}public Long createWorkFlow(String name, String description, String locations,String taskDefinitionJson, String taskRelationJson,String executionType) {//SSLUtil.turnOffSslChecking();HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());LinkedMultiValueMap<String, Object> linkedMultiValueMap = new LinkedMultiValueMap<String, Object>();linkedMultiValueMap.add("description", description);linkedMultiValueMap.add("name", name);linkedMultiValueMap.add("taskDefinitionJson", taskDefinitionJson);linkedMultiValueMap.add("taskRelationJson", taskRelationJson);linkedMultiValueMap.add("timeout", 0);linkedMultiValueMap.add("executionType", executionType);linkedMultiValueMap.add("tenantCode", dolphinschedulerConfig.getTenantCode());linkedMultiValueMap.add("locations", locations);HttpEntity<MultiValueMap<String, Object>> httpEntity = new HttpEntity<>(linkedMultiValueMap, hs);String url = dolphinschedulerConfig.getDsdUrl() + "/projects/" + dolphinschedulerConfig.getProjectCode() + "/process-definition";JSONObject resultJSON = doPostForObject(url, httpEntity);JSONObject data = (JSONObject) resultJSON.get(DATA);Long code = (Long) data.get("code");return code;}public void updateReleaseState(String name, String releaseState, Long code) {//SSLUtil.turnOffSslChecking();HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());LinkedMultiValueMap<String, Object> linkedMultiValueMap = new LinkedMultiValueMap<String, Object>();linkedMultiValueMap.add("name", name);linkedMultiValueMap.add("releaseState", releaseState);HttpEntity<MultiValueMap<String, Object>> httpEntity = new HttpEntity<>(linkedMultiValueMap, hs);String url = dolphinschedulerConfig.getDsdUrl() + "/projects/" + dolphinschedulerConfig.getProjectCode() + "/process-definition/" + code + "/release";doPostForObject(url, httpEntity);}public List<TaskDefinition> getTaskByWorkflowCode(Long dsdCode) {//SSLUtil.turnOffSslChecking();HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());HttpEntity<MultiValueMap<String, Object>> httpEntity = new HttpEntity<>(hs);String url = dolphinschedulerConfig.getDsdUrl() + "/projects/" + dolphinschedulerConfig.getProjectCode() + "/process-definition/" + dsdCode;JSONObject resultJSON = doGetForObject(url, httpEntity);JSONObject data = (JSONObject) resultJSON.get(DATA);JSONArray jsonArray = (JSONArray) data.get("taskDefinitionList");return jsonArray.toJavaList(TaskDefinition.class);}/*** 刪除工作流** @param codes*/public void delWorkflow(String codes) {HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());LinkedMultiValueMap<String, Object> linkedMultiValueMap = new LinkedMultiValueMap<String, Object>();linkedMultiValueMap.add("codes", codes);HttpEntity<MultiValueMap<String, Object>> httpEntity = new HttpEntity<>(linkedMultiValueMap, hs);String url = dolphinschedulerConfig.getDsdUrl() + "/projects/" + dolphinschedulerConfig.getProjectCode() + "/process-definition/batch-delete";doPostForObject(url, httpEntity);}public void updateWorkFlow(Long dsdCode, String name, String description, String locations, String taskDefinitionJson, String taskRelationJson, String executionType) {HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());LinkedMultiValueMap<String, Object> linkedMultiValueMap = new LinkedMultiValueMap<String, Object>();linkedMultiValueMap.add("description", description);linkedMultiValueMap.add("name", name);linkedMultiValueMap.add("taskDefinitionJson", taskDefinitionJson);linkedMultiValueMap.add("taskRelationJson", taskRelationJson);linkedMultiValueMap.add("timeout", 0);linkedMultiValueMap.add("executionType", executionType);linkedMultiValueMap.add("tenantCode", dolphinschedulerConfig.getTenantCode());linkedMultiValueMap.add("locations", locations);HttpEntity<MultiValueMap<String, Object>> httpEntity = new HttpEntity<>(linkedMultiValueMap, hs);String url = dolphinschedulerConfig.getDsdUrl() + "/projects/" + dolphinschedulerConfig.getProjectCode() + "/process-definition/" + dsdCode;doPutForObject(url, httpEntity);}/*** 運行工作流** @param code*/public void runWorkflow(Long code) {HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());LinkedMultiValueMap<String, Object> linkedMultiValueMap = new LinkedMultiValueMap<String, Object>();linkedMultiValueMap.add("processDefinitionCode", code);linkedMultiValueMap.add("failureStrategy", "CONTINUE");linkedMultiValueMap.add("warningType", "NONE");linkedMultiValueMap.add("scheduleTime", getStringDate());HttpEntity<MultiValueMap<String, Object>> httpEntity = new HttpEntity<>(linkedMultiValueMap, hs);String url = dolphinschedulerConfig.getDsdUrl() + "/projects/" + dolphinschedulerConfig.getProjectCode() + "/executors/start-process-instance";doPostForObject(url, httpEntity);}public String getStringDate() {LocalDateTime currentDateTime = LocalDateTime.now();LocalDateTime startDate = currentDateTime.withHour(0).withMinute(0).withSecond(0).withNano(0);LocalDateTime endDate = startDate;DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");String formattedStartDate = startDate.format(formatter);JSONObject jsonObject = new JSONObject();jsonObject.put("complementStartDate", formattedStartDate);jsonObject.put("complementEndDate", formattedStartDate);return jsonObject.toString();}/*** 獲取任務日志** @param id* @return*/public ResponseTaskLog getLog(Integer id, Integer limit, Integer skipLineNum) {HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());HttpEntity<String> entity = new HttpEntity<String>(hs);String url = dolphinschedulerConfig.getDsdUrl() + "/log/detail?taskInstanceId=" + id + "&limit=" + limit + "&skipLineNum=" + skipLineNum;JSONObject resultJSON = doGetForObject(url, entity);return JSON.parseObject(resultJSON.get(DATA).toString(), ResponseTaskLog.class);}/*** 重跑任務** @param processInstanceId*/public void operation(Integer processInstanceId, String executeType) {HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());LinkedMultiValueMap<String, Object> linkedMultiValueMap = new LinkedMultiValueMap<String, Object>();// 1 REPEAT_RUNNING 重跑 2 STOP 停止 3 RECOVER_SUSPENDED_PROCESS 恢復運行 4 PAUSE 暫停linkedMultiValueMap.add("processInstanceId", processInstanceId);switch (executeType) {case "1":addExecutionDetails(linkedMultiValueMap, 1, "REPEAT_RUNNING", "run");break;case "2":linkedMultiValueMap.add("executeType", "STOP");break;case "3":addExecutionDetails(linkedMultiValueMap, 0, "RECOVER_SUSPENDED_PROCESS", "suspend");break;case "4":linkedMultiValueMap.add("executeType", "PAUSE");break;default:throw new BizException("暫不支持該操作");}HttpEntity<MultiValueMap<String, Object>> httpEntity = new HttpEntity<>(linkedMultiValueMap, hs);String url = dolphinschedulerConfig.getDsdUrl() + "/projects/" + dolphinschedulerConfig.getProjectCode() + "/executors/execute";doPostForObject(url, httpEntity);}public void addExecutionDetails(MultiValueMap<String, Object> map, int index, String executeType, String buttonType) {map.add("index", String.valueOf(index));map.add("executeType", executeType);if (buttonType != null) {map.add("buttonType", buttonType);}}public PageInfo<ProcessInstanceVO> processInstances(Long dsdWorkflowCode, String searchVal, Integer pageNum, Integer pageSize,String startDate, String endDate, String stateType) {HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());HttpEntity<String> entity = new HttpEntity<String>(hs);StringBuilder urlBuilder = new StringBuilder(dolphinschedulerConfig.getDsdUrl()).append("/projects/").append(dolphinschedulerConfig.getProjectCode()).append("/process-instances?pageNo=").append(pageNum).append("&pageSize=").append(pageSize).append("&call=").append("1");// 這個必須加,不然刪除工作流后,實例會不見if(null != dsdWorkflowCode){urlBuilder.append("&processDefineCode=").append(dsdWorkflowCode);}if(StringUtils.isNotBlank(searchVal)){urlBuilder.append("&searchVal=").append(searchVal);}supplementaryParameters(startDate, endDate, stateType, urlBuilder);JSONObject resultJSON = doGetForObject(urlBuilder.toString(), entity);JSONObject data = (JSONObject) resultJSON.get(DATA);JSONArray jsonArray = (JSONArray) data.get("totalList");List<ProcessInstanceVO> javaList = jsonArray.toJavaList(ProcessInstanceVO.class);Integer total = (Integer) data.get("total");PageInfo<ProcessInstanceVO> res = new PageInfo<>();res.setList(javaList);res.setTotal(total);return res;}private void supplementaryParameters(String startDate, String endDate, String stateType, StringBuilder urlBuilder) {if (stateType != null && !stateType.isEmpty()) {urlBuilder.append("&stateType=").append(stateType);}if (startDate != null && !startDate.isEmpty()) {urlBuilder.append("&startDate=").append(startDate);}if (endDate != null && !endDate.isEmpty()) {urlBuilder.append("&endDate=").append(endDate);}}public PageInfo<TaskInstanceVO> taskInstances(Integer processInstanceId, String startDate, String endDate,String stateType, int pageNum, int pageSize) {HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());HttpEntity<String> entity = new HttpEntity<String>(hs);StringBuilder urlBuilder = new StringBuilder(dolphinschedulerConfig.getDsdUrl()).append("/projects/").append(dolphinschedulerConfig.getProjectCode()).append("/task-instances?pageNo=").append(pageNum).append("&pageSize=").append(pageSize).append("&processInstanceId=").append(processInstanceId).append("&taskExecuteType=").append("BATCH");supplementaryParameters(startDate, endDate, stateType, urlBuilder);String url = urlBuilder.toString();JSONObject resultJSON = doGetForObject(url, entity);JSONObject data = (JSONObject) resultJSON.get(DATA);JSONArray jsonArray = (JSONArray) data.get("totalList");List<TaskInstanceVO> javaList = jsonArray.toJavaList(TaskInstanceVO.class);Integer total = (Integer) data.get("total");PageInfo<TaskInstanceVO> res = new PageInfo<>();res.setList(javaList);res.setTotal(total);return res;}/*** 獲取工作流執行順序* @param processInstanceId* @return*/public List<GanttTaskVO> viewGantt(Long processInstanceId) {HttpHeaders hs = new HttpHeaders();hs.add("Content-Type", CONTENT_TYPE);hs.add("X-Requested-With", X_REQUESTED_WITH);hs.add(TOKEN, dolphinschedulerConfig.getDsdToken());HttpEntity<String> entity = new HttpEntity<String>(hs);StringBuilder urlBuilder = new StringBuilder(dolphinschedulerConfig.getDsdUrl()).append("/projects/").append(dolphinschedulerConfig.getProjectCode()).append("/process-instances/").append(processInstanceId).append("/view-gantt");String url = urlBuilder.toString();JSONObject resultJSON = doGetForObject(url, entity);JSONObject data = (JSONObject) resultJSON.get(DATA);JSONArray jsonArray = (JSONArray) data.get("tasks");List<GanttTaskVO> javaList = jsonArray.toJavaList(GanttTaskVO.class);return javaList;}
}

本文來自互聯網用戶投稿,該文觀點僅代表作者本人,不代表本站立場。本站僅提供信息存儲空間服務,不擁有所有權,不承擔相關法律責任。
如若轉載,請注明出處:http://www.pswp.cn/web/39861.shtml
繁體地址,請注明出處:http://hk.pswp.cn/web/39861.shtml
英文地址,請注明出處:http://en.pswp.cn/web/39861.shtml

如若內容造成侵權/違法違規/事實不符,請聯系多彩編程網進行投訴反饋email:809451989@qq.com,一經查實,立即刪除!

相關文章

項目管理實用表格與應用【項目文件資料分享】

項目管理基礎知識 項目管理可分為五大過程組&#xff08;啟動、規劃、執行、監控、收尾&#xff09;十大知識領域&#xff0c;其中包含49個子過程 項目十大知識領域分為&#xff1a;項目整合管理、項目范圍管理、項目進度管理、項目成本管理、項目質量管理、項目資源管理、項目…

標量場與向量場

標量場與向量場 flyfish 場 是一個函數&#xff0c;它把空間中的每一點關聯到一個數值或一個數學對象&#xff08;如向量、張量等&#xff09;。在物理學中&#xff0c;場可以描述許多物理現象&#xff0c;例如溫度分布、電場、磁場、壓力場等。 標量場 標量場 是一個函數&…

【BUUCTF-PWN】9-ciscn_2019_n_8

不屬于棧溢出&#xff0c;應該是比較簡單的pwn&#xff0c;看懂代碼邏輯使用pwntools 32位&#xff0c;開啟了Stack、NX、PIE保護 執行效果&#xff1a; main函數 使用通義千問詢問的代碼解讀&#xff1a; 即當var數組的第十四個元素是17就可以 這里可以用兩種payload…

Python使用總結之應用程序有哪些配置方式?配置方式對比

Python使用總結之應用程序有哪些配置方式&#xff1f;配置方式對比 在Python程序中&#xff0c;管理配置信息的方法有很多&#xff0c;常見的方式包括使用INI文件、JSON文件、YAML文件、環境變量、以及直接在代碼中定義配置。每種方式都有其獨特的優勢和適用場景。 1. INI文件 …

天環公益原創開發進度網站源碼帶后臺免費分享

天環公益計劃首發原創開發進度網站源碼帶后臺免費分享 后臺地址是&#xff1a;admin.php 后臺沒有賬號密碼 這個沒有數據庫 有能力的可以自己改 天環公益原創開發進度網站 帶后臺

ARM架構服務器/虛擬機編譯部署Tendis(國產化替換Redis)

文章目錄 一、概述 二、安裝相關組件 三、下載最新的Tendis源碼 四、編譯源碼 五、啟動Tendis 六、使用Docker鏡像部署Tendis 七、常見報錯 八、參考鏈接 一、概述 國產化項目要求盡可能使用國產組件,尤其是已存在的項目,需要替換已有組件,比如使用Tendis替換Redis。…

微軟中國全面撤店!我們到現場看了看

ChatGPT狂飆160天&#xff0c;世界已經不是之前的樣子。 更多資源歡迎關注 7月1日&#xff0c;微軟官方發言人向媒體表示&#xff1a; “微軟不斷評估其零售策略以滿足我們的客戶不斷變化的需求&#xff0c;微軟已決定對中國大陸市場的渠道進行整合。客戶仍可通過零售合作伙伴…

校園失物招領系統帶萬字文檔java項目失物招領管理系統java課程設計java畢業設計springboot vue

文章目錄 校園失物招領系統一、項目演示二、項目介紹三、萬字字項目文檔四、部分功能截圖五、部分代碼展示六、底部獲取項目源碼帶萬字文檔&#xff08;9.9&#xffe5;帶走&#xff09; 校園失物招領系統 一、項目演示 校園失物招領系統 二、項目介紹 語言: Java 數據庫&…

JAVA導出數據庫字典到Excel

文章目錄 1、查詢某張表字段信息2、TableVo接收sql查詢得到的數據3、excel導出4、導出案例 1、查詢某張表字段信息 select column_name as columnName, -- 字段名 COLUMN_DEFAULT as colDefault, -- 默認值 column_key as columnKey, -- PRI-主鍵&#xff0c;UNI-唯一鍵&…

【Tools】 Postman 接口測試工具詳解

那年夏天我和你躲在 這一大片寧靜的海 直到后來我們都還在 對這個世界充滿期待 今年冬天你已經不在 我的心空出了一塊 很高興遇見你 讓我終究明白 回憶比真實精彩 &#x1f3b5; 王心凌《那年夏天寧靜的海》 在現代軟件開發中&#xff0c;API&#xff08;…

【Python實戰因果推斷】21_傾向分1

目錄 The Impact of Management Training Adjusting with Regression 之前學習了如何使用線性回歸調整混雜因素。此外&#xff0c;還向您介紹了通過正交化去偏差的概念&#xff0c;這是目前最有用的偏差調整技術之一。不過&#xff0c;您還需要學習另一種技術--傾向加權。這種…

Ionic 卡片:設計和使用指南

Ionic 卡片&#xff1a;設計和使用指南 Ionic 是一個強大的開源框架&#xff0c;用于構建跨平臺的移動應用程序。它結合了 Angular、React 和 Vue 的強大功能&#xff0c;允許開發者使用 Web 技術創建高性能的移動應用。Ionic 卡片是框架中的一個核心組件&#xff0c;用于展示…

js使用插件完成xml轉json

插件&#xff1a;xml2json.min.js 插件文件下載&#xff08;不能上傳附件&#xff09;&#xff1a;https://download.csdn.net/download/zhu_zhu_xia/89513965 html代碼&#xff1a; <!DOCTYPE html> <html lang"en"><head><meta charset&qu…

我認為一般信息管理應用中使用存儲過程高效

總看有些人反對使用存儲過程&#xff0c;原因無非是以下幾點 1.不利于更換數據庫&#xff0c;就是沒有移植性 2.不利用調試和擴展 就依據我們大大小小項目&#xff0c;風風雨雨走過近20年&#xff0c;每個系統的業務邏輯處理幾乎都是用存儲過程實現的&#xff0c;沒發現多不…

p標簽文本段落中因編輯器換行引起的空格問題完美解決方案

目錄 1.修改前的代碼&#xff1a;2.修改后的代碼3.總結 在HTML文檔中&#xff0c;如何要在&#xff08;p標簽&#xff09;內寫一段很長的文本段落&#xff0c;并且沒有 換行。由于IDE或者編輯器界面大小有限或需要在vue中邏輯處理動態顯示文本&#xff0c;一行寫完太長&#x…

Eslint prettier airbnb規范 配置

1.安裝vscode的Eslint和prettier 插件 eslint&#xff1a;代碼質量檢查工具 https://eslint.nodejs.cn/docs/latest/use/getting-started prettier&#xff1a;代碼風格格式化工具 https://www.prettier.cn/docs/index.html /* eslint-config-airbnb-base airbnb 規范 esl…

高德地圖軌跡回放并提示具體信息

先上效果圖 到達某地點后顯示提示語&#xff1a;比如&#xff1a;12&#xff1a;56分駛入康莊大道、左轉駛入xx大道等 <!doctype html> <html> <head><meta charset"utf-8"><meta http-equiv"X-UA-Compatible" content"…

【前端CSS3】CSS顯示模式(黑馬程序員)

文章目錄 一、前言&#x1f680;&#x1f680;&#x1f680;二、CSS元素顯示模式&#xff1a;??????2.1 什么是元素顯示模式2.2 塊元素2.3 行內元素2.4 行塊元素2.5 元素顯示模式的轉換 三、總結&#x1f680;&#x1f680;&#x1f680; 一、前言&#x1f680;&#x1f…

巴圖自動化Modbus協議轉Profinet協議網關模塊連智能儀表與PLC通訊

一、現場要求:PLC作為控制器&#xff0c;儀表設備作為執行設備。執行設備可以實時響應PLC傳送的指令&#xff0c;并將數據反饋給PLC&#xff0c;從而實現PLC對儀表設備的控制和監控&#xff0c;實現對生產過程的精確控制。 二、解決方案:通過巴圖自動化Modbus協議轉Profinet協議…

前端面試題4(瀏覽器對http請求處理過程)

瀏覽器對http請求處理過程 當我們在瀏覽器中輸入URL并按下回車鍵時&#xff0c;瀏覽器會執行一系列步驟來處理HTTP請求并與服務器通信。下面是瀏覽器處理過程 1. 解析URL 瀏覽器首先解析輸入的URL&#xff0c;提取出協議&#xff08;通常是http://或https://&#xff09;、主…