Commit eced1602 by liyuan

Merge branch 'developer' into 'master'

调度基本功能完善版本

See merge request !2
parents f7a2ba2e 8783187d
*.class
*.iml
.idea
target
module_*.xml
byit-myth-job.xml
\ No newline at end of file
# byit-math-job
>分布式任务调度
\ No newline at end of file
>分布式任务调度
## 一、目录结构
### 1.byit-myth-admin
>介绍:
### 2.byit-myth-core
>介绍:
#### 2.1 myth-admin-core
>介绍:
#### 2.2 myth-executor-core
>介绍:
### 3.byit-myth-executor
>介绍:
### 4.byit-myth-register
>介绍:
### 5.byit-myth-rpc
>介绍:
\ No newline at end of file
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.byit</groupId>
<artifactId>byit-mybatis-plugin</artifactId>
<version>1.0.0-SNAPSHOT</version>
<packaging>jar</packaging>
<name>byit-mybatis-plugin</name>
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
<java.version>1.8</java.version>
<maven.compiler.source>1.8</maven.compiler.source>
<maven.compiler.target>1.8</maven.compiler.target>
</properties>
<dependencies>
<dependency>
<groupId>org.mybatis.generator</groupId>
<artifactId>mybatis-generator-core</artifactId>
<version>1.3.5</version>
</dependency>
<dependency>
<groupId>ch.qos.logback</groupId>
<artifactId>logback-classic</artifactId>
<version>1.2.3</version>
<scope>compile</scope>
</dependency>
</dependencies>
<distributionManagement>
<repository>
<id>releases</id>
<url>http://10.0.120.2/repository/maven-releases/</url>
</repository>
<snapshotRepository>
<id>snapshots</id>
<url>http://10.0.120.2/repository/maven-snapshots/</url>
</snapshotRepository>
</distributionManagement>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
<version>1.5.19.RELEASE</version>
</plugin>
</plugins>
</build>
</project>
package com.byit.plugin;
import java.util.List;
import org.mybatis.generator.api.*;
import org.mybatis.generator.api.dom.java.Field;
import org.mybatis.generator.api.dom.java.FullyQualifiedJavaType;
import org.mybatis.generator.api.dom.java.JavaVisibility;
import org.mybatis.generator.api.dom.java.Method;
import org.mybatis.generator.api.dom.java.Parameter;
import org.mybatis.generator.api.dom.java.PrimitiveTypeWrapper;
import org.mybatis.generator.api.dom.java.TopLevelClass;
import org.mybatis.generator.api.dom.xml.Attribute;
import org.mybatis.generator.api.dom.xml.TextElement;
import org.mybatis.generator.api.dom.xml.XmlElement;
public class AddLimitOffsetPlugin extends PluginAdapter {
@Override
public boolean validate(List<String> warnings) {
return true;
}
@Override
public boolean modelExampleClassGenerated(TopLevelClass topLevelClass, IntrospectedTable introspectedTable) {
PrimitiveTypeWrapper integerWrapper = FullyQualifiedJavaType.getIntInstance().getPrimitiveTypeWrapper();
Field limit = new Field();
limit.setName("limit");
limit.setVisibility(JavaVisibility.PRIVATE);
limit.setType(integerWrapper);
topLevelClass.addField(limit);
Method limitSet = new Method();
limitSet.setVisibility(JavaVisibility.PUBLIC);
limitSet.setName("setLimit");
limitSet.addParameter(new Parameter(integerWrapper, "limit"));
limitSet.addBodyLine("this.limit = limit;");
topLevelClass.addMethod(limitSet);
Method limitGet = new Method();
limitGet.setVisibility(JavaVisibility.PUBLIC);
limitGet.setReturnType(integerWrapper);
limitGet.setName("getLimit");
limitGet.addBodyLine("return limit;");
topLevelClass.addMethod(limitGet);
Field offset = new Field();
offset.setName("offset");
offset.setVisibility(JavaVisibility.PRIVATE);
offset.setType(integerWrapper);
topLevelClass.addField(offset);
Method offsetSet = new Method();
offsetSet.setVisibility(JavaVisibility.PUBLIC);
offsetSet.setName("setOffset");
offsetSet.addParameter(new Parameter(integerWrapper, "offset"));
offsetSet.addBodyLine("this.offset = offset;");
topLevelClass.addMethod(offsetSet);
Method offsetGet = new Method();
offsetGet.setVisibility(JavaVisibility.PUBLIC);
offsetGet.setReturnType(integerWrapper);
offsetGet.setName("getOffset");
offsetGet.addBodyLine("return offset;");
topLevelClass.addMethod(offsetGet);
return true;
}
@Override
public boolean sqlMapSelectByExampleWithoutBLOBsElementGenerated(XmlElement element,
IntrospectedTable introspectedTable) {
@SuppressWarnings("unused")
FullyQualifiedTable table = introspectedTable.getFullyQualifiedTable();
// XmlElement lastElement =
// (XmlElement)element.getElements().get(element.getElements().size());
XmlElement isNotNullElement = new XmlElement("if");
isNotNullElement.addAttribute(new Attribute("test", "limit != null"));
isNotNullElement.addElement(new TextElement("limit ${limit}"));
element.getElements().add(isNotNullElement);
isNotNullElement = new XmlElement("if");
isNotNullElement.addAttribute(new Attribute("test", "offset != null"));
isNotNullElement.addElement(new TextElement("offset ${offset}"));
element.getElements().add(isNotNullElement);
return true;
}
@Override
public boolean modelGetterMethodGenerated(Method method, TopLevelClass topLevelClass, IntrospectedColumn introspectedColumn, IntrospectedTable introspectedTable, ModelClassType modelClassType) {
return false;
}
@Override
public boolean modelSetterMethodGenerated(Method method, TopLevelClass topLevelClass, IntrospectedColumn introspectedColumn, IntrospectedTable introspectedTable, ModelClassType modelClassType) {
return false;
}
}
package com.byit.plugin;
import org.mybatis.generator.api.PluginAdapter;
import org.mybatis.generator.internal.util.StringUtility;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.List;
/**
* @program: ddmp-parent
* @description:
* @author: guoqingming
* @create: 2019-03-07 22:55
**/
public class BasePlugin extends PluginAdapter {
protected static final Logger logger = LoggerFactory.getLogger(BasePlugin.class);
@Override
public boolean validate(List<String> warnings) {
// 插件使用前提是targetRuntime为MyBatis3
if (StringUtility.stringHasValue(getContext().getTargetRuntime()) && "MyBatis3".equalsIgnoreCase(getContext().getTargetRuntime()) == false) {
warnings.add("assembly mybatis:插件" + this.getClass().getTypeName() + "要求运行targetRuntime必须为MyBatis3!");
return false;
}
return true;
}
}
package com.byit.plugin;
import java.sql.Types;
import java.util.Properties;
import org.mybatis.generator.api.IntrospectedColumn;
import org.mybatis.generator.api.IntrospectedTable;
import org.mybatis.generator.api.dom.java.FullyQualifiedJavaType;
import org.mybatis.generator.api.dom.java.Method;
import org.mybatis.generator.api.dom.java.TopLevelClass;
import org.mybatis.generator.internal.types.JavaTypeResolverDefaultImpl;
public class DefaultJavaTypeResolverDefault extends JavaTypeResolverDefaultImpl {
public DefaultJavaTypeResolverDefault() {
super();
super.typeMap.put(Types.SMALLINT, new JdbcTypeInformation("SMALLINT", //$NON-NLS-1$
new FullyQualifiedJavaType(Integer.class.getName())));
super.typeMap.put(Types.TINYINT, new JdbcTypeInformation("TINYINT", //$NON-NLS-1$
new FullyQualifiedJavaType(Integer.class.getName())));
typeMap.put(Types.BIT, new JdbcTypeInformation("BIT", //$NON-NLS-1$
new FullyQualifiedJavaType(Integer.class.getName())));
}
@Override
public void addConfigurationProperties(Properties properties) {
super.addConfigurationProperties(properties);
}
}
package com.byit.plugin;
import java.util.List;
import org.mybatis.generator.api.GeneratedXmlFile;
import org.mybatis.generator.api.IntrospectedColumn;
import org.mybatis.generator.api.IntrospectedTable;
import org.mybatis.generator.api.PluginAdapter;
import org.mybatis.generator.api.dom.java.Interface;
import org.mybatis.generator.api.dom.java.Method;
import org.mybatis.generator.api.dom.java.TopLevelClass;
import org.mybatis.generator.config.Context;
import org.mybatis.generator.config.TableConfiguration;
public class DefaultNamePlugin extends PluginAdapter {
@Override
public boolean validate(List<String> warnings) {
return true;
}
private void defaultRename(IntrospectedTable introspectedTable) {
introspectedTable.setSelectByPrimaryKeyStatementId("getById");
introspectedTable.setDeleteByPrimaryKeyStatementId("deleteById");
introspectedTable.setUpdateByPrimaryKeyStatementId("updateById");
introspectedTable.setUpdateByPrimaryKeySelectiveStatementId("updateByIdSelective");
}
/**
* int deleteByPrimaryKey(Long id);
*
* int insert(UserRegister record);
*
* int insertSelective(UserRegister record);
*
* UserRegister selectByPrimaryKey(Long id);
*
* int updateByPrimaryKeySelective(UserRegister record);
*
* int updateByPrimaryKey(UserRegister record);
*/
@Override
public boolean clientGenerated(Interface interfaze, TopLevelClass topLevelClass,
IntrospectedTable introspectedTable) {
defaultRename(introspectedTable);
for (Method method : interfaze.getMethods()) {
String methodName = method.getName();
if (methodName.contains("PrimaryKey")) {
method.setName(methodName.replace("PrimaryKey", "Id"));
}
methodName = method.getName();
if (methodName.contains("select")) {
method.setName(methodName.replace("select", "get"));
}
}
return true;
}
@Override
public boolean sqlMapGenerated(GeneratedXmlFile sqlMap, IntrospectedTable introspectedTable) {
defaultRename(introspectedTable);
return true;
}
@Override
public void setContext(Context context) {
List<TableConfiguration> list = context.getTableConfigurations();
for (TableConfiguration tableConfiguration : list) {
tableConfiguration.setCountByExampleStatementEnabled(false);
tableConfiguration.setDeleteByExampleStatementEnabled(false);
tableConfiguration.setSelectByExampleStatementEnabled(false);
tableConfiguration.setUpdateByExampleStatementEnabled(false);
}
super.setContext(context);
}
@Override
public boolean modelGetterMethodGenerated(Method method, TopLevelClass topLevelClass, IntrospectedColumn introspectedColumn, IntrospectedTable introspectedTable, ModelClassType modelClassType) {
return false;
}
@Override
public boolean modelSetterMethodGenerated(Method method, TopLevelClass topLevelClass, IntrospectedColumn introspectedColumn, IntrospectedTable introspectedTable, ModelClassType modelClassType) {
return false;
}
}
package com.byit.plugin;
import java.util.List;
import org.mybatis.generator.api.IntrospectedColumn;
import org.mybatis.generator.api.IntrospectedTable;
import org.mybatis.generator.api.PluginAdapter;
import org.mybatis.generator.api.dom.java.FullyQualifiedJavaType;
import org.mybatis.generator.api.dom.java.Interface;
import org.mybatis.generator.api.dom.java.Method;
import org.mybatis.generator.api.dom.java.Parameter;
import org.mybatis.generator.api.dom.java.TopLevelClass;
import org.mybatis.generator.api.dom.xml.Attribute;
import org.mybatis.generator.api.dom.xml.Document;
import org.mybatis.generator.api.dom.xml.TextElement;
import org.mybatis.generator.api.dom.xml.XmlElement;
public class DeleteLogicByIdsPlugin extends PluginAdapter {
/**
* {@inheritDoc}
*/
@Override
public boolean validate(List<String> warnings) {
return true;
}
/**
* {@inheritDoc}
*/
@Override
public boolean clientSelectByExampleWithBLOBsMethodGenerated(Method method,
Interface interfaze, IntrospectedTable introspectedTable) {
interfaze.addMethod(generateDeleteLogicByIds(method,
introspectedTable));
return true;
}
/**
* {@inheritDoc}
*/
@Override
public boolean clientSelectByExampleWithoutBLOBsMethodGenerated(
Method method, Interface interfaze,
IntrospectedTable introspectedTable) {
interfaze.addMethod(generateDeleteLogicByIds(method,
introspectedTable));
return true;
}
/**
* {@inheritDoc}
*/
@Override
public boolean clientSelectByExampleWithBLOBsMethodGenerated(Method method,
TopLevelClass topLevelClass, IntrospectedTable introspectedTable) {
topLevelClass.addMethod(generateDeleteLogicByIds(method,
introspectedTable));
return true;
}
/**
* {@inheritDoc}
*/
@Override
public boolean clientSelectByExampleWithoutBLOBsMethodGenerated(
Method method, TopLevelClass topLevelClass,
IntrospectedTable introspectedTable) {
topLevelClass.addMethod(generateDeleteLogicByIds(method,
introspectedTable));
return true;
}
@Override
public boolean sqlMapDocumentGenerated(Document document, IntrospectedTable introspectedTable) {
String tableName = introspectedTable.getAliasedFullyQualifiedTableNameAtRuntime();//数据库表名
XmlElement parentElement = document.getRootElement();
// 产生分页语句前半部分
XmlElement deleteLogicByIdsElement = new XmlElement("update");
deleteLogicByIdsElement.addAttribute(new Attribute("id", "deleteLogicByIds"));
deleteLogicByIdsElement.addElement(
new TextElement(
"update " + tableName + " set deleteFlag = #{deleteFlag,jdbcType=INTEGER} where id in "
+ " <foreach item=\"item\" index=\"index\" collection=\"ids\" open=\"(\" separator=\",\" close=\")\">#{item}</foreach> "
));
parentElement.addElement(deleteLogicByIdsElement);
return super.sqlMapDocumentGenerated(document, introspectedTable);
}
private Method generateDeleteLogicByIds(Method method, IntrospectedTable introspectedTable) {
Method m = new Method("deleteLogicByIds");
m.setVisibility(method.getVisibility());
m.setReturnType(FullyQualifiedJavaType.getIntInstance());
m.addParameter(new Parameter(FullyQualifiedJavaType.getIntInstance(), "deleteFlag", "@Param(\"deleteFlag\")"));
m.addParameter(new Parameter(new FullyQualifiedJavaType("Integer[]"), "ids", "@Param(\"ids\")"));
context.getCommentGenerator().addGeneralMethodComment(m,
introspectedTable);
return m;
}
@Override
public boolean modelGetterMethodGenerated(Method method, TopLevelClass topLevelClass, IntrospectedColumn introspectedColumn, IntrospectedTable introspectedTable, ModelClassType modelClassType) {
return false;
}
@Override
public boolean modelSetterMethodGenerated(Method method, TopLevelClass topLevelClass, IntrospectedColumn introspectedColumn, IntrospectedTable introspectedTable, ModelClassType modelClassType) {
return false;
}
}
package com.byit.plugin;
import com.byit.plugin.utils.FormatTools;
import com.byit.plugin.utils.JavaElementGeneratorTools;
import com.byit.plugin.utils.XmlElementGeneratorTools;
import org.mybatis.generator.api.IntrospectedColumn;
import org.mybatis.generator.api.IntrospectedTable;
import org.mybatis.generator.api.dom.java.*;
import org.mybatis.generator.api.dom.xml.*;
import org.mybatis.generator.codegen.ibatis2.sqlmap.elements.AbstractXmlElementGenerator;
import org.mybatis.generator.codegen.mybatis3.ListUtilities;
import org.mybatis.generator.config.GeneratedKey;
import java.util.List;
/**
* @program: ddmp-parent
* @description:
* @author: guoqingming
* @create: 2019-03-07 22:53
**/
public class InsertIgnorePlugin extends BasePlugin{
public InsertIgnorePlugin() {
}
public static final String METHOD_INSERT_IGNORE = "insertIgnore"; // 方法名
@Override
public boolean validate(List<String> warnings) {
// 该插件只支持MYSQL
if ("com.mysql.jdbc.Driver".equalsIgnoreCase(this.getContext().getJdbcConnectionConfiguration().getDriverClass()) == false
&& "com.mysql.cj.jdbc.Driver".equalsIgnoreCase(this.getContext().getJdbcConnectionConfiguration().getDriverClass()) == false) {
warnings.add("assembly mybatis:插件" + this.getClass().getTypeName() + "只支持MySQL数据库!");
return false;
}
return super.validate(warnings);
}
@Override
public boolean clientGenerated(Interface interfaze, TopLevelClass topLevelClass, IntrospectedTable introspectedTable) {
Method mUpsert = JavaElementGeneratorTools.generateMethod(
METHOD_INSERT_IGNORE,
JavaVisibility.DEFAULT,
FullyQualifiedJavaType.getIntInstance(),
new Parameter(JavaElementGeneratorTools.getModelTypeWithoutBLOBs(introspectedTable), "record")
);
// interface 增加方法
FormatTools.addMethodWithBestPosition(interfaze, mUpsert);
return super.clientGenerated(interfaze, topLevelClass, introspectedTable);
}
@Override
public boolean sqlMapDocumentGenerated(Document document, IntrospectedTable introspectedTable) {
XmlElement insertIgnoreElement = new XmlElement("insert");
//添加ID
insertIgnoreElement.addAttribute(new Attribute("id", METHOD_INSERT_IGNORE));
// 参数类型
FullyQualifiedJavaType parameterType = JavaElementGeneratorTools.getModelTypeWithoutBLOBs(introspectedTable);
insertIgnoreElement.addAttribute(new Attribute("parameterType", //$NON-NLS-1$
parameterType.getFullyQualifiedName()));
// insertIgnoreElement.addAttribute(new Attribute("parameterType", "map"));
GeneratedKey gk = introspectedTable.getGeneratedKey();
if (gk != null) {
IntrospectedColumn introspectedColumn = introspectedTable
.getColumn(gk.getColumn());
// if the column is null, then it's a configuration error. The
// warning has already been reported
if (introspectedColumn != null) {
if (gk.isJdbcStandard()) {
insertIgnoreElement.addAttribute(new Attribute(
"useGeneratedKeys", "true")); //$NON-NLS-1$ //$NON-NLS-2$
insertIgnoreElement.addAttribute(new Attribute(
"keyProperty", introspectedColumn.getJavaProperty())); //$NON-NLS-1$
insertIgnoreElement.addAttribute(new Attribute(
"keyColumn", introspectedColumn.getActualColumnName())); //$NON-NLS-1$
} else {
insertIgnoreElement.addElement(new XMLElementGenerator().getSelectKeyPublic(introspectedColumn, gk));
}
}
}
//insert
insertIgnoreElement.addElement(new TextElement("insert ignore into " + introspectedTable.getFullyQualifiedTableNameAtRuntime()));
for (Element element : XmlElementGeneratorTools.generateKeys(ListUtilities.removeIdentityAndGeneratedAlwaysColumns(introspectedTable.getAllColumns()), true)) {
insertIgnoreElement.addElement(element);
}
insertIgnoreElement.addElement(new TextElement("values"));
for (Element element : XmlElementGeneratorTools.generateValues(ListUtilities.removeIdentityAndGeneratedAlwaysColumns(introspectedTable.getAllColumns()), "")) {
insertIgnoreElement.addElement(element);
}
document.getRootElement().addElement(insertIgnoreElement);
return super.sqlMapDocumentGenerated(document, introspectedTable);
}
private static class XMLElementGenerator extends AbstractXmlElementGenerator {
public XmlElement getSelectKeyPublic(IntrospectedColumn introspectedColumn,
GeneratedKey generatedKey) {
return this.getSelectKey(introspectedColumn, generatedKey);
}
@Override
public void addElements(XmlElement parentElement) {
}
}
}
package com.byit.plugin;
import java.util.ArrayList;
import java.util.List;
import org.mybatis.generator.api.IntrospectedColumn;
import org.mybatis.generator.api.IntrospectedTable;
import org.mybatis.generator.api.PluginAdapter;
import org.mybatis.generator.api.dom.java.FullyQualifiedJavaType;
import org.mybatis.generator.api.dom.java.Method;
import org.mybatis.generator.api.dom.java.TopLevelClass;
public class LombokAnnotationPlugin extends PluginAdapter {
@Override
public boolean validate(List<String> list) {
return false;
}
@Override
public boolean modelBaseRecordClassGenerated(TopLevelClass topLevelClass, IntrospectedTable introspectedTable) {
topLevelClass.addAnnotation("@Data");
topLevelClass.addImportedType(new FullyQualifiedJavaType("lombok.Data"));
List<Method> methods = topLevelClass.getMethods();
List<Method> remove = new ArrayList<Method>();
for (Method method : methods) {
if (method.getBodyLines().size() < 2) {
remove.add(method);
}
}
methods.removeAll(remove);
return true;
}
@Override
public boolean modelGetterMethodGenerated(Method method, TopLevelClass topLevelClass, IntrospectedColumn introspectedColumn, IntrospectedTable introspectedTable, ModelClassType modelClassType) {
return false;
}
@Override
public boolean modelSetterMethodGenerated(Method method, TopLevelClass topLevelClass, IntrospectedColumn introspectedColumn, IntrospectedTable introspectedTable, ModelClassType modelClassType) {
return false;
}
}
\ No newline at end of file
package com.byit.plugin;
import java.util.List;
import org.mybatis.generator.api.CommentGenerator;
import org.mybatis.generator.api.IntrospectedColumn;
import org.mybatis.generator.api.IntrospectedTable;
import org.mybatis.generator.api.PluginAdapter;
import org.mybatis.generator.api.dom.java.Field;
import org.mybatis.generator.api.dom.java.FullyQualifiedJavaType;
import org.mybatis.generator.api.dom.java.JavaVisibility;
import org.mybatis.generator.api.dom.java.Method;
import org.mybatis.generator.api.dom.java.Parameter;
import org.mybatis.generator.api.dom.java.TopLevelClass;
import org.mybatis.generator.api.dom.xml.Attribute;
import org.mybatis.generator.api.dom.xml.TextElement;
import org.mybatis.generator.api.dom.xml.XmlElement;
public class MySQLPaginationPlugin extends PluginAdapter
{
@Override
public boolean modelExampleClassGenerated(TopLevelClass topLevelClass, IntrospectedTable introspectedTable)
{
// add field, getter, setter for limit clause
addPage(topLevelClass, introspectedTable, "page");
return super.modelExampleClassGenerated(topLevelClass, introspectedTable);
}
@Override
public boolean sqlMapSelectByExampleWithoutBLOBsElementGenerated(XmlElement element, IntrospectedTable introspectedTable)
{
XmlElement page = new XmlElement("if");
page.addAttribute(new Attribute("test", "page != null"));
page.addElement(new TextElement("limit #{page.begin} , #{page.length}"));
element.addElement(page);
return super.sqlMapUpdateByExampleWithoutBLOBsElementGenerated(element, introspectedTable);
}
/**
* @param topLevelClass
* @param introspectedTable
* @param name
*/
private void addPage(TopLevelClass topLevelClass, IntrospectedTable introspectedTable, String name)
{
topLevelClass.addImportedType(new FullyQualifiedJavaType("net.javaw.mybatis.generator.Page"));
CommentGenerator commentGenerator = context.getCommentGenerator();
Field field = new Field();
field.setVisibility(JavaVisibility.PROTECTED);
field.setType(new FullyQualifiedJavaType("net.javaw.mybatis.generator.Page"));
field.setName(name);
commentGenerator.addFieldComment(field, introspectedTable);
topLevelClass.addField(field);
char c = name.charAt(0);
String camel = Character.toUpperCase(c) + name.substring(1);
Method method = new Method();
method.setVisibility(JavaVisibility.PUBLIC);
method.setName("set" + camel);
method.addParameter(new Parameter(new FullyQualifiedJavaType("net.javaw.mybatis.generator.Page"), name));
method.addBodyLine("this." + name + "=" + name + ";");
commentGenerator.addGeneralMethodComment(method, introspectedTable);
topLevelClass.addMethod(method);
method = new Method();
method.setVisibility(JavaVisibility.PUBLIC);
method.setReturnType(new FullyQualifiedJavaType("net.javaw.mybatis.generator.Page"));
method.setName("get" + camel);
method.addBodyLine("return " + name + ";");
commentGenerator.addGeneralMethodComment(method, introspectedTable);
topLevelClass.addMethod(method);
}
/**
* This plugin is always valid - no properties are required
*/
@Override
public boolean validate(List<String> warnings)
{
return true;
}
@Override
public boolean modelGetterMethodGenerated(Method method, TopLevelClass topLevelClass, IntrospectedColumn introspectedColumn, IntrospectedTable introspectedTable, ModelClassType modelClassType) {
return false;
}
@Override
public boolean modelSetterMethodGenerated(Method method, TopLevelClass topLevelClass, IntrospectedColumn introspectedColumn, IntrospectedTable introspectedTable, ModelClassType modelClassType) {
return false;
}
}
package com.byit.plugin.utils;
import org.mybatis.generator.api.CommentGenerator;
import org.mybatis.generator.api.IntrospectedTable;
import org.mybatis.generator.api.dom.java.*;
import org.mybatis.generator.api.dom.xml.Attribute;
import org.mybatis.generator.api.dom.xml.Element;
import org.mybatis.generator.api.dom.xml.TextElement;
import org.mybatis.generator.api.dom.xml.XmlElement;
import java.util.Iterator;
import java.util.List;
import java.util.Set;
import java.util.TreeSet;
/**
* @program: ddmp-parent
* @description:
* @author: guoqingming
* @create: 2019-03-07 23:02
**/
public class FormatTools {
/**
* 在最佳位置添加方法
* @param innerClass
* @param method
*/
public static void addMethodWithBestPosition(InnerClass innerClass, Method method) {
addMethodWithBestPosition(method, innerClass.getMethods());
}
/**
* 在最佳位置添加方法
* @param interfacz
* @param method
*/
public static void addMethodWithBestPosition(Interface interfacz, Method method) {
// import
Set<FullyQualifiedJavaType> importTypes = new TreeSet<>();
// 返回
if (method.getReturnType() != null) {
importTypes.add(method.getReturnType());
importTypes.addAll(method.getReturnType().getTypeArguments());
}
// 参数 比较特殊的是ModelColumn生成的Column
for (Parameter parameter : method.getParameters()) {
boolean flag = true;
for (String annotation : parameter.getAnnotations()) {
if (annotation.startsWith("@Param")) {
importTypes.add(new FullyQualifiedJavaType("org.apache.ibatis.annotations.Param"));
if (annotation.matches(".*selective.*")) {
flag = false;
}
}
}
if (flag) {
importTypes.add(parameter.getType());
importTypes.addAll(parameter.getType().getTypeArguments());
}
}
interfacz.addImportedTypes(importTypes);
addMethodWithBestPosition(method, interfacz.getMethods());
}
/**
* 在最佳位置添加方法
* @param innerEnum
* @param method
*/
public static void addMethodWithBestPosition(InnerEnum innerEnum, Method method) {
addMethodWithBestPosition(method, innerEnum.getMethods());
}
/**
* 在最佳位置添加方法
* @param topLevelClass
* @param method
*/
public static void addMethodWithBestPosition(TopLevelClass topLevelClass, Method method) {
addMethodWithBestPosition(method, topLevelClass.getMethods());
}
/**
* 在最佳位置添加节点
* @param rootElement
* @param element
*/
public static void addElementWithBestPosition(XmlElement rootElement, XmlElement element) {
// sql 元素都放在sql后面
if (element.getName().equals("sql")) {
int index = 0;
for (Element ele : rootElement.getElements()) {
if (ele instanceof XmlElement && ((XmlElement) ele).getName().equals("sql")) {
index++;
}
}
rootElement.addElement(index, element);
} else {
// 根据id 排序
String id = getIdFromElement(element);
if (id == null) {
rootElement.addElement(element);
} else {
List<Element> elements = rootElement.getElements();
int index = -1;
for (int i = 0; i < elements.size(); i++) {
Element ele = elements.get(i);
if (ele instanceof XmlElement) {
String eleId = getIdFromElement((XmlElement) ele);
if (eleId != null) {
if (eleId.startsWith(id)) {
if (index == -1) {
index = i;
}
} else if (id.startsWith(eleId)) {
index = i + 1;
}
}
}
}
if (index == -1 || index >= elements.size()) {
rootElement.addElement(element);
} else {
elements.add(index, element);
}
}
}
}
/**
* 找出节点ID值
* @param element
* @return
*/
private static String getIdFromElement(XmlElement element) {
for (Attribute attribute : element.getAttributes()) {
if (attribute.getName().equals("id")) {
return attribute.getValue();
}
}
return null;
}
/**
* 获取最佳添加位置
* @param method
* @param methods
* @return
*/
private static void addMethodWithBestPosition(Method method, List<Method> methods) {
int index = -1;
for (int i = 0; i < methods.size(); i++) {
Method m = methods.get(i);
if (m.getName().equals(method.getName())) {
if (m.getParameters().size() <= method.getParameters().size()) {
index = i + 1;
} else {
index = i;
}
} else if (m.getName().startsWith(method.getName())) {
if (index == -1) {
index = i;
}
} else if (method.getName().startsWith(m.getName())) {
index = i + 1;
}
}
if (index == -1 || index >= methods.size()) {
methods.add(methods.size(), method);
} else {
methods.add(index, method);
}
}
/**
* 替换已有方法注释
* @param commentGenerator
* @param method
* @param introspectedTable
*/
public static void replaceGeneralMethodComment(CommentGenerator commentGenerator, Method method, IntrospectedTable introspectedTable) {
method.getJavaDocLines().clear();
commentGenerator.addGeneralMethodComment(method, introspectedTable);
}
/**
* 替换已有注释
* @param commentGenerator
* @param element
*/
public static void replaceComment(CommentGenerator commentGenerator, XmlElement element) {
Iterator<Element> elementIterator = element.getElements().iterator();
boolean flag = false;
while (elementIterator.hasNext()) {
Element ele = elementIterator.next();
if (ele instanceof TextElement && ((TextElement) ele).getContent().matches("<!--")) {
flag = true;
}
if (flag) {
elementIterator.remove();
}
if (ele instanceof TextElement && ((TextElement) ele).getContent().matches("-->")) {
flag = false;
}
}
XmlElement tmpEle = new XmlElement("tmp");
commentGenerator.addComment(tmpEle);
for (int i = tmpEle.getElements().size() - 1; i >= 0; i--) {
element.addElement(0, tmpEle.getElements().get(i));
}
}
}
package com.byit.plugin.utils;
import org.mybatis.generator.api.IntrospectedTable;
import org.mybatis.generator.api.dom.java.*;
import static org.mybatis.generator.internal.util.messages.Messages.getString;
public class JavaElementGeneratorTools {
/**
* 生成静态常量
* @param fieldName 常量名称
* @param javaType 类型
* @param initString 初始化字段
* @return
*/
public static Field generateStaticFinalField(String fieldName, FullyQualifiedJavaType javaType, String initString) {
Field field = new Field(fieldName, javaType);
field.setVisibility(JavaVisibility.PUBLIC);
field.setStatic(true);
field.setFinal(true);
if (initString != null) {
field.setInitializationString(initString);
}
return field;
}
/**
* 生成属性
* @param fieldName 常量名称
* @param visibility 可见性
* @param javaType 类型
* @param initString 初始化字段
* @return
*/
public static Field generateField(String fieldName, JavaVisibility visibility, FullyQualifiedJavaType javaType, String initString) {
Field field = new Field(fieldName, javaType);
field.setVisibility(visibility);
if (initString != null) {
field.setInitializationString(initString);
}
return field;
}
/**
* 生成方法
* @param methodName 方法名
* @param visibility 可见性
* @param returnType 返回值类型
* @param parameters 参数列表
* @return
*/
public static Method generateMethod(String methodName, JavaVisibility visibility, FullyQualifiedJavaType returnType, Parameter... parameters) {
Method method = new Method(methodName);
method.setVisibility(visibility);
method.setReturnType(returnType);
if (parameters != null) {
for (Parameter parameter : parameters) {
method.addParameter(parameter);
}
}
return method;
}
/**
* 生成方法实现体
* @param method 方法
* @param bodyLines 方法实现行
* @return
*/
public static Method generateMethodBody(Method method, String... bodyLines) {
if (bodyLines != null) {
for (String bodyLine : bodyLines) {
method.addBodyLine(bodyLine);
}
}
return method;
}
/**
* 生成Filed的Set方法
* @param field field
* @return
*/
public static Method generateSetterMethod(Field field) {
Method method = generateMethod(
"set" + field.getName().substring(0, 1).toUpperCase() + field.getName().substring(1),
JavaVisibility.PUBLIC,
null,
new Parameter(field.getType(), field.getName())
);
return generateMethodBody(method, "this." + field.getName() + " = " + field.getName() + ";");
}
/**
* 生成Filed的Get方法
* @param field field
* @return
*/
public static Method generateGetterMethod(Field field) {
Method method = generateMethod(
"get" + field.getName().substring(0, 1).toUpperCase() + field.getName().substring(1),
JavaVisibility.PUBLIC,
field.getType()
);
return generateMethodBody(method, "return this." + field.getName() + ";");
}
/**
* 获取Model没有BLOBs类时的类型
* @param introspectedTable
* @return
*/
public static FullyQualifiedJavaType getModelTypeWithoutBLOBs(IntrospectedTable introspectedTable) {
FullyQualifiedJavaType type;
if (introspectedTable.getRules().generateBaseRecordClass()) {
type = new FullyQualifiedJavaType(introspectedTable.getBaseRecordType());
} else if (introspectedTable.getRules().generatePrimaryKeyClass()) {
type = new FullyQualifiedJavaType(introspectedTable.getPrimaryKeyType());
} else {
throw new RuntimeException(getString("RuntimeError.12"));
}
return type;
}
/**
* 获取Model有BLOBs类时的类型
* @param introspectedTable
* @return
*/
public static FullyQualifiedJavaType getModelTypeWithBLOBs(IntrospectedTable introspectedTable) {
FullyQualifiedJavaType type;
if (introspectedTable.getRules().generateRecordWithBLOBsClass()) {
type = new FullyQualifiedJavaType(introspectedTable.getRecordWithBLOBsType());
} else {
// the blob fields must be rolled up into the base class
type = new FullyQualifiedJavaType(introspectedTable.getBaseRecordType());
}
return type;
}
}
......@@ -12,4 +12,134 @@
<artifactId>byit-myth-admin</artifactId>
<dependencies>
<dependency>
<groupId>myth-job</groupId>
<artifactId>myth-admin-core</artifactId>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
</dependency>
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>fastjson</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-thymeleaf</artifactId>
</dependency>
<!-- https://mvnrepository.com/artifact/net.sourceforge.nekohtml/nekohtml -->
<dependency>
<groupId>net.sourceforge.nekohtml</groupId>
<artifactId>nekohtml</artifactId>
<version>1.9.22</version>
</dependency>
<dependency>
<groupId>io.springfox</groupId>
<artifactId>springfox-swagger2</artifactId>
</dependency>
<dependency>
<groupId>io.springfox</groupId>
<artifactId>springfox-swagger-ui</artifactId>
</dependency>
<dependency>
<groupId>com.google.guava</groupId>
<artifactId>guava</artifactId>
<version>20.0</version>
</dependency>
<dependency>
<groupId>myth-job</groupId>
<artifactId>myth-executor-api</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-devtools</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>myth-job</groupId>
<artifactId>myth-web-core</artifactId>
<version>1.0-SNAPSHOT</version>
</dependency>
</dependencies>
<build>
<finalName>${project.artifactId}</finalName>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
<configuration>
<!-- 指定SpringBoot程序的main函数入口类 -->
<mainClass>com.byit.AdminApplication</mainClass>
</configuration>
<executions>
<execution>
<goals>
<goal>repackage</goal>
</goals>
</execution>
</executions>
</plugin>
<plugin>
<artifactId>maven-compiler-plugin</artifactId>
<configuration>
<source>1.8</source>
<target>1.8</target>
<encoding>UTF-8</encoding>
<compilerArguments>
<!-- 打包本地jar包 -->
<extdirs>${project.basedir}/lib</extdirs>
</compilerArguments>
</configuration>
</plugin>
<plugin>
<groupId>org.mybatis.generator</groupId>
<artifactId>mybatis-generator-maven-plugin</artifactId>
<version>1.3.5</version>
<dependencies>
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<version>5.1.44</version>
</dependency>
<dependency>
<groupId>org.mybatis.spring.boot</groupId>
<artifactId>mybatis-spring-boot-starter</artifactId>
<version>1.2.0</version>
</dependency>
<dependency>
<groupId>com.byit</groupId>
<artifactId>byit-mybatis-plugin</artifactId>
<version>1.0.0-SNAPSHOT</version>
</dependency>
</dependencies>
<configuration>
<configurationFile>${basedir}/src/main/resources/Generator-config.xml</configurationFile>
<verbose>true</verbose>
<overwrite>true</overwrite>
<skip>false</skip>
</configuration>
</plugin>
</plugins>
</build>
</project>
\ No newline at end of file
package com.byit;
import com.byit.annotations.EnablePluginClient;
import com.byit.rpc.remoting.provider.annotation.RpcService;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.annotation.Bean;
import org.springframework.transaction.annotation.EnableTransactionManagement;
import org.springframework.web.cors.CorsConfiguration;
import org.springframework.web.cors.UrlBasedCorsConfigurationSource;
import org.springframework.web.filter.CorsFilter;
/**
* @program: byit-myth-job->AdminApplication
* @description: 管理启动类
* @author: huangfu
* @date: 2019/11/20 11:44
**/
@SpringBootApplication
@RpcService(http_type = true)
@EnablePluginClient
public class AdminApplication {
public static void main(String[] args) {
SpringApplication.run(AdminApplication.class,args);
}
}
package com.byit.api;
import com.byit.conf.MythJobAutoConfigure;
import com.byit.dto.executor.JobRunResultDto;
import com.byit.packet.response.PluginRpcResponsePacket;
import com.byit.thread.JavaTaskCallbackThread;
import com.byit.thread.LogCallbackThread;
import io.swagger.annotations.Api;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
/**
* @author huangfu
*/
@Api(tags = "回调api")
@RestController
@RequestMapping("api/callback")
public class ApiCallbackController {
@PostMapping(value = "callbackRes")
public void callbackRes(@RequestBody PluginRpcResponsePacket pluginRpcResponsePacket){
MythJobAutoConfigure.LOG_CALLBACK.execute(new JavaTaskCallbackThread(pluginRpcResponsePacket));
}
}
package com.byit.api;
import com.byit.dto.plugin.CollectData;
import com.byit.dto.web.ResponseResult;
import com.byit.model.RunRecording;
import com.byit.model.vo.RunRecordingVo;
import com.byit.service.ApiFlowService;
import com.byit.service.FlowService;
import io.swagger.annotations.Api;
import io.swagger.annotations.ApiOperation;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import javax.annotation.Resource;
import java.text.ParseException;
import java.util.Date;
import java.util.List;
import java.util.Map;
/**
* @description: 工作流操作的API接口
* @author: gml
* @create: 2019-12-30 14:20
*/
@Api(tags = "工作流api")
@RestController
@RequestMapping("api/flow")
public class ApiFlowController {
@Resource
private ApiFlowService apiFlowService;
@Resource
private FlowService flowService;
@PostMapping("publish")
@ApiOperation("发布工作流,并开始调度")
public String publishFlow(String param) throws Exception {
apiFlowService.publishFlow(param);
return "SUCCESS";
}
@PostMapping("deleteFlow")
@ApiOperation("删除工作流,若当前工作流被依赖则删除失败,只允许删除不被依赖的工作流,当前工作流依赖其他工作流不影响")
public ResponseResult deleteFlow(String param){
apiFlowService.deleteFlow(param);
return ResponseResult.ok("SUCCESS");
}
@PostMapping("start")
@ApiOperation("开始工作流的调度,将工作流启用调度")
public ResponseResult start(String param) throws ParseException {
apiFlowService.start(param);
return ResponseResult.ok("SUCCESS");
}
@PostMapping("exist")
@ApiOperation("判断工作流是否存在")
public ResponseResult exist(String param) throws ParseException {
Boolean result = apiFlowService.exist(param);
return ResponseResult.ok(result);
}
@PostMapping("repealSchedule")
@ApiOperation("撤销工作流调度,只可以撤销总工作流,若为内嵌工作流不允许撤销")
public ResponseResult repealSchedule(String param) throws ParseException {
apiFlowService.repealSchedule(param);
return ResponseResult.ok("SUCCESS");
}
@PostMapping("killJob")
@ApiOperation("杀死节点")
public ResponseResult killJob(String param) throws InterruptedException {
Boolean result = flowService.killJob(param);
return ResponseResult.ok(result);
}
@PostMapping("killFlow")
@ApiOperation("杀死工作流")
public ResponseResult killFlow(String param) throws InterruptedException {
Boolean result = flowService.killFlow(param);
return ResponseResult.ok(result);
}
@PostMapping("stopSchedule")
@ApiOperation("暂停工作流所有调度")
public ResponseResult stopSchedule(String param){
return ResponseResult.ok(apiFlowService.stopSchedule(param));
}
@PostMapping("stopScheduleByRunId")
@ApiOperation("暂停指定的工作流调度")
public ResponseResult stopScheduleByRunId(String runId){
return ResponseResult.ok(apiFlowService.stopScheduleByRunId(runId));
}
@PostMapping("reStartSchedule")
@ApiOperation("重新开始某次调度")
public ResponseResult reStartSchedule(String runIds){
apiFlowService.reStartSchedule(runIds);
return ResponseResult.ok("SUCCESS");
}
@PostMapping("reRunJob")
@ApiOperation("重跑节点")
public ResponseResult reRunJob(String param){
apiFlowService.reRunJob(param);
return ResponseResult.ok("SUCCESS");
}
@PostMapping("reRunFlow")
@ApiOperation("重跑工作流")
public ResponseResult reRunFlow(String param){
apiFlowService.reRunFlow(param);
return ResponseResult.ok("SUCCESS");
}
@PostMapping("madeSuccess")
@ApiOperation("手动置为成功")
public ResponseResult madeSuccess(String param){
apiFlowService.madeSuccess(param);
return ResponseResult.ok("SUCCESS");
}
/**
*
* @param param startTime
* endTime
* flowName
* workspaceName
* @return
*/
@PostMapping("loadScheduleResult")
@ApiOperation("加载运行记录")
public ResponseResult loadScheduleResult(String param){
List<RunRecording> result = apiFlowService.loadScheduleResult(param);
return ResponseResult.ok(result);
}
/**
*
* @param param
* @return
*/
@PostMapping("loadScheduleLog")
@ApiOperation("加载运行日志")
public ResponseResult loadScheduleLog(String param){
RunRecordingVo result = apiFlowService.loadScheduleLog(param);
return ResponseResult.ok(result);
}
/**
* 补批节点
* @param param
* @return
*/
@PostMapping("/repairJob")
@ApiOperation("补批")
public ResponseResult repairJob(String param){
apiFlowService.repairJob(param);
return ResponseResult.ok("SUCCESS");
}
/**
* 补批工作流
* @param param
* @return
*/
@PostMapping("/repairFlow")
@ApiOperation("补批工作流")
public ResponseResult repairFlow(String param){
apiFlowService.repairFlow(param);
return ResponseResult.ok("SUCCESS");
}
/**
* 获取整体运行的统计数据
* @param param
* @return
*/
@PostMapping("/loadStatisticData")
@ApiOperation("获取运行的统计数据")
public ResponseResult loadStatisticData(String param){
CollectData collectData = apiFlowService.loadStatisticData(param);
return ResponseResult.ok(collectData);
}
/**
* 获取节点运行的统计数据
* @param param
* @return
*/
@PostMapping("/loadNodeStatisticData")
@ApiOperation("获取运行的统计数据")
public ResponseResult loadNodeStatisticData(String param){
CollectData collectData = apiFlowService.loadNodeStatisticData(param);
return ResponseResult.ok(collectData);
}
@PostMapping("/loadCurrentStatus")
@ApiOperation("获取当前的工作流运行状态")
public ResponseResult loadCurrentStatus(String param){
Map<String, RunRecording> statusMap = apiFlowService.loadCurrentStatus(param);
return ResponseResult.ok(statusMap);
}
}
package com.byit.api;
import com.byit.dto.web.ResponseResult;
import com.byit.model.JobTaskRunLog;
import com.byit.model.JobTaskRunLogWithBLOBs;
import com.byit.service.ApiNodeService;
import io.swagger.annotations.Api;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import javax.annotation.Resource;
import java.util.List;
import java.util.Map;
@Api(tags = "任务节点api")
@RestController
@RequestMapping("api/node")
@Slf4j
public class ApiNodeController {
@Resource
private ApiNodeService apiNodeService;
@Resource
private StringRedisTemplate stringRedisTemplate;
@PostMapping("runNode")
public String runNode(String param){
log.debug("----------=--runNode方法接收到参数为:{}--------------",param);
String monitorKey = apiNodeService.runNode(param);
return monitorKey;
}
@PostMapping("killNode")
public String killNode(Integer logId){
log.debug("----------=--killNode方法接收到参数为:{}--------------",logId);
apiNodeService.killNode(logId);
return "SUCCESS";
}
@PostMapping("runHistory")
public ResponseResult runHistory(String nodeId){
List<JobTaskRunLogWithBLOBs> jobTaskRunLogList = apiNodeService.runHistory(nodeId);
return ResponseResult.ok(jobTaskRunLogList);
}
@PostMapping("autoAddJavaTask")
public String autoAddJavaTask(String param) throws Exception {
apiNodeService.autoAddJavaTask(param);
return "SUCCESS";
}
@PostMapping("addJavaTask")
public String addJavaTask(String param)throws Exception {
apiNodeService.addJavaTask(param);
return "SUCCESS";
}
@PostMapping("updateJavaTask")
public String updateJavaTask(String param)throws Exception {
apiNodeService.updateJavaTask(param);
return "SUCCESS";
}
@PostMapping("deleteJavaTask")
public String deleteJavaTask(String jobName){
apiNodeService.deleteJavaTask(jobName);
return "SUCCESS";
}
@PostMapping("existJavaTask")
public Boolean existJavaTask(String jobName){
Boolean result = apiNodeService.existJavaTask(jobName);
return result;
}
@PostMapping("loadLogByJobName")
public ResponseResult loadLogByJobName(String param){
List<JobTaskRunLog> jobTaskRunLogList = apiNodeService.loadLogByJobName(param);
return ResponseResult.ok(jobTaskRunLogList);
}
@PostMapping("loadLogByTaskName")
public ResponseResult loadLogByTaskName(String param){
List<JobTaskRunLog> jobTaskRunLogList = apiNodeService.loadLogByTaskName(param);
return ResponseResult.ok(jobTaskRunLogList);
}
@PostMapping("runJavaTask")
public ResponseResult runJavaTask(String jobName){
apiNodeService.runJavaTask(jobName);
return ResponseResult.ok("SUCCESS");
}
@PostMapping("runTask")
public ResponseResult runTask(String param){
apiNodeService.runTask(param);
return ResponseResult.ok("SUCCESS");
}
@PostMapping("loadCurrentStatusByJobName")
public ResponseResult loadCurrentStatusByJobName(String jobNames){
Map<String, JobTaskRunLog> jobTaskRunLogMap = apiNodeService.loadCurrentStatusByJobName(jobNames);
return ResponseResult.ok(jobTaskRunLogMap);
}
}
package com.byit.api;
import com.byit.dto.web.ResponseResult;
import com.byit.service.ApiWorkspaceService;
import io.swagger.annotations.Api;
import io.swagger.annotations.ApiOperation;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import javax.annotation.Resource;
/**
* @description: 工作空间操作的API接口
* @author: gml
* @create: 2019-12-31 11:19
*/
@Api(tags = "工作空间api")
@RestController
@RequestMapping("api/workspace")
public class ApiWorkspaceController {
@Resource
private ApiWorkspaceService workspaceService;
@PostMapping("add")
@ApiOperation("新建工作空间")
public ResponseResult add(String workspaceName){
workspaceService.add(workspaceName);
return ResponseResult.ok("SUCCESS");
}
@PostMapping("exist")
@ApiOperation("是否存在")
public Boolean exist(String workspaceName){
return workspaceService.exist(workspaceName);
}
}
package com.byit.config;
import com.byit.rpc.registry.impl.RegistryServiceRegistry;
import com.byit.rpc.remoting.invoker.impl.RpcSpringInvokerFactory;
import com.byit.rpc.remoting.provider.impl.RpcSpringProviderFactory;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.HashMap;
/**
* @author huangfu
**/
@Configuration
@Slf4j
public class AdminRegisterConfig {
@Value("${myth-rpc.registry.address}")
private String address;
@Value("${myth-rpc.registry.biz}")
private String biz;
@Value("${myth-rpc.registry.env}")
private String env;
// @Value("${myth-rpc.registry.port}")
// private int port;
/**
* 初始化 rpc客户端 属于消费者
* @return
*/
@Bean
public RpcSpringInvokerFactory invokerFactory(){
RpcSpringInvokerFactory invokerFactory = new RpcSpringInvokerFactory();
invokerFactory.setServiceRegistryClass(RegistryServiceRegistry.class);
invokerFactory.setServiceRegistryParam(new HashMap<String,String>(){{
put(RegistryServiceRegistry.REGISTRY_ADDRESS,address);
put(RegistryServiceRegistry.BIZ, biz);
put(RegistryServiceRegistry.ENV,env);
}});
log.info(">>>>>>>>>> byit-myth-admin client invoker config 初始化成功");
return invokerFactory;
}
/**
* 设置对外提供服务的接口 这个其实是配合网关调用的方式 属于生产者
* @return
*/
@Bean
public RpcSpringProviderFactory rpcSpringProviderFactory() {
RpcSpringProviderFactory providerFactory = new RpcSpringProviderFactory();
//providerFactory.setPort(port);
providerFactory.setServiceRegistryClass(RegistryServiceRegistry.class);
providerFactory.setServiceRegistryParam(new HashMap<String, String>() {{
put(RegistryServiceRegistry.REGISTRY_ADDRESS, address);
put(RegistryServiceRegistry.BIZ, biz);
put(RegistryServiceRegistry.ENV, env);
}});
log.info(">>>>>>>>>>>>>>> byit-myth-admin 向注册中心注册config初始化完成");
return providerFactory;
}
}
package com.byit.config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import springfox.documentation.builders.ApiInfoBuilder;
import springfox.documentation.builders.PathSelectors;
import springfox.documentation.builders.RequestHandlerSelectors;
import springfox.documentation.service.ApiInfo;
import springfox.documentation.service.Contact;
import springfox.documentation.spi.DocumentationType;
import springfox.documentation.spring.web.plugins.Docket;
import springfox.documentation.swagger2.annotations.EnableSwagger2;
@Configuration
@EnableSwagger2
public class SwaggerConfig{
@Bean
public Docket api() {
return new Docket(DocumentationType.SWAGGER_2)
.select()
.apis(RequestHandlerSelectors.basePackage("com.byit"))
.paths(PathSelectors.any())
.build()
.apiInfo(apiInfo());
}
//构建 api文档的详细信息函数,注意这里的注解引用的是哪个
private ApiInfo apiInfo() {
return new ApiInfoBuilder()
//页面标题
.title("Myth-Job 项目集成Swagger接口文档")
//创建人
.contact(new Contact("gml", "", ""))
//版本号
.version("1.0")
//描述
.description("接口描述")
.build();
}
}
\ No newline at end of file
package com.byit.config;
import org.springframework.context.annotation.Configuration;
import org.springframework.web.servlet.config.annotation.ResourceHandlerRegistry;
import org.springframework.web.servlet.config.annotation.WebMvcConfigurerAdapter;
/**
* @description:
* @author: gml
* @create: 2020-01-06 13:53
*/
@Configuration
public class WebMvcConfig extends WebMvcConfigurerAdapter {
@Override
public void addResourceHandlers(ResourceHandlerRegistry registry) {
registry.addResourceHandler("swagger-ui.html")
.addResourceLocations("classpath:/META-INF/resources/");
registry.addResourceHandler("/webjars/**")
.addResourceLocations("classpath:/META-INF/resources/webjars/");
}
}
package com.byit.controller;
import com.byit.exception.DataValidationException;
import com.byit.model.vo.EmailAlarmVo;
import com.byit.service.EmailAlarmService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.UUID;
/**
* @program: byit-myth-job->EmailController
* @description: TODO
* @author: huangfu
* @date: 2019/12/30 11:46
**/
@RestController
public class EmailController {
private final EmailAlarmService emailAlarmService;
@Autowired
public EmailController(EmailAlarmService emailAlarmService) {
this.emailAlarmService = emailAlarmService;
}
@RequestMapping("send")
public String send(){
EmailAlarmVo emailAlarmVo = new EmailAlarmVo();
emailAlarmVo.setFlowId(1);
emailAlarmVo.setFlowName("测试工作流");
emailAlarmVo.setRunId(UUID.randomUUID( ).toString());
emailAlarmVo.setVersionName("V1");
emailAlarmVo.setAlarmEmail("huangfukexing@byitgroup.com");
emailAlarmService.sendEmail(emailAlarmVo);
return "success";
}
}
package com.byit.controller;
import com.byit.model.vo.FlowVo;
import com.byit.model.vo.NodeVo;
import com.byit.service.FlowService;
import io.swagger.annotations.Api;
import io.swagger.annotations.ApiOperation;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import javax.annotation.Resource;
import java.util.List;
/**
* @description: 工作流
* @author: gml
* @create: 2019-12-23 15:25
*/
@Api(tags = "工作流")
@RestController
@RequestMapping("flow")
public class FlowController {
@Resource
private FlowService flowService;
@PostMapping("save")
@ApiOperation("保存工作流信息")
public String saveFlow(@RequestBody FlowVo flowVo){
flowService.saveJobFlow(flowVo);
return "SUCCESS";
}
@PostMapping("create")
@ApiOperation("创建工作流")
public String createFlow(@RequestBody FlowVo flowVo){
flowService.createFlow(flowVo);
return "SUCCESS";
}
@PostMapping("delete")
@ApiOperation("删除工作流")
public String deleteFlow(@RequestBody Integer flowId){
flowService.deleteFlow(flowId);
return "SUCCESS";
}
@PostMapping("get")
@ApiOperation("查询工作流")
public FlowVo getFlow(@RequestBody Integer flowId){
FlowVo flow = flowService.getFlowId(flowId);
return flow;
}
@PostMapping("start")
@ApiOperation("启动工作流的调度")
public Boolean startFlow(@RequestBody FlowVo flowVo){
Boolean result = flowService.startFlow(flowVo);
return result;
}
@PostMapping("find")
public List<NodeVo> findFlow(@RequestBody Integer workSpaceId){
return null;
}
}
package com.byit.controller;
import com.byit.conf.MythJobAutoConfigure;
import com.byit.dto.executor.DispatchResponseDto;
import com.byit.dto.executor.JobRunResultDto;
import com.byit.dto.executor.PluginBeanJobInfo;
import com.byit.model.JobTask;
import com.byit.service.JobTaskService;
import com.byit.service.RunScriptService;
import com.byit.thread.LogCallbackThread;
import com.byit.util.SourceObj2TargetObjUtil;
import com.byit.utils.ValidationUtil;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;
import java.util.List;
/**
* @program: byit-myth-job->JobHandel
* @description: 测试添加任务
* @author: huangfu
* @date: 2019/11/20 11:57
**/
@RestController
@RequestMapping("job")
public class JobController {
@Autowired
private RunScriptService runScriptService;
private final JobTaskService jobTaskService;
@Autowired
public JobController(JobTaskService jobTaskService) {
this.jobTaskService = jobTaskService;
}
@PostMapping(value = "addJob")
public String addJob(@RequestBody PluginBeanJobInfo pluginBeanJobInfo){
JobTask jobTask = SourceObj2TargetObjUtil.pluginBeanJobInfo2JobTask(pluginBeanJobInfo);
jobTaskService.addMythJobTask(jobTask);
return "SUCCESS";
}
@PostMapping(value = "callbackRes")
public void callbackRes(@RequestBody JobRunResultDto jobRunResultDto){
MythJobAutoConfigure.LOG_CALLBACK.execute(new LogCallbackThread(jobRunResultDto));
}
@GetMapping(value = "getJobTask")
public List<JobTask> getJobInfo(){
ValidationUtil.dataNotNull(null,"不能为null");
return jobTaskService.findJobTaskByTriggerNextTimeLessThanEqual(11);
}
@GetMapping("test")
public DispatchResponseDto test(){
DispatchResponseDto dispatchResponseDto = runScriptService.runScript(null);
System.out.println("---------------------");
return dispatchResponseDto;
}
}
package com.byit.service;
import com.byit.dto.plugin.CollectData;
import com.byit.dto.plugin.PluginFlow;
import com.byit.model.RunRecording;
import com.byit.model.vo.RunRecordingVo;
import java.text.ParseException;
import java.util.List;
import java.util.Map;
/**
* @description: 工作流的api请求业务处理接口
* @author: gml
* @create: 2019-12-30 14:34
*/
public interface ApiFlowService {
void deleteFlow(String param);
void start(String param) throws ParseException;
void repealSchedule(String param) throws ParseException;
void killSchedule(String param);
String stopSchedule(String param);
void reStartSchedule(String runIds);
void publishFlow(String param) throws Exception;
/**
* 判断工作流石佛存在
* @param workspaceId
* @param flow
* @return
*/
void validateFlow(Integer workspaceId, PluginFlow flow, boolean isInner);
/**
* 重跑任务
* @param param
*/
void reRunJob(String param);
/**
* 手动置为成功
* @param param
*/
void madeSuccess(String param);
/**
* 重跑工作流
* @param param
*/
void reRunFlow(String param);
/**
* 加载运行记录
* @param param
* @return
*/
List<RunRecording> loadScheduleResult(String param);
/**
* 补批
* @param param
*/
void repairJob(String param);
/**
* 补批工作流
* @param param
*/
void repairFlow(String param);
/**
* 加载运行日志
* @param param
* @return
*/
RunRecordingVo loadScheduleLog(String param);
/**
* 统计汇总数据
* @param param
* @return
*/
CollectData loadStatisticData(String param);
/**
* 获取节点的汇总数据
* @param param
* @return
*/
CollectData loadNodeStatisticData(String param);
/**
* 获取工作流的当前状态
* @param param
* @return
*/
Map<String, RunRecording> loadCurrentStatus(String param);
/**
* 功能描述 判断工作流是否存在
* @author gml
* @date 2020-05-27 10:12
* @param param
* @return java.lang.Boolean
*/
Boolean exist(String param);
Boolean stopScheduleByRunId(String runId);
}
package com.byit.service;
import com.byit.model.JobTaskRunLog;
import com.byit.model.JobTaskRunLogWithBLOBs;
import java.text.ParseException;
import java.util.List;
import java.util.Map;
/**
* 任务节点service
*/
public interface ApiNodeService {
/**
* 单独运行节点
* @param param
* @return
*/
String runNode(String param);
/**
* 根据nodeid查询运行历史
* @param nodeId
* @return
*/
List<JobTaskRunLogWithBLOBs> runHistory(String nodeId);
/**
* 功能描述 自动发布任务的接口
* @author gml
* @date 2020-04-14 10:55
* @param param
* @return void
*/
void autoAddJavaTask(String param) throws Exception;
/**
* 功能描述 发布任务接口
* @author gml
* @date 2020-04-14 11:27
* @param param
* @return void
*/
void addJavaTask(String param) throws Exception;
/**
* 立即运行任务不需要验证是否存在
* @param param
* @throws Exception
*/
void runTask(String param) ;
/**
* 功能描述 更新任务配置信息
* @author gml
* @date 2020-04-14 11:28
* @param param
* @return void
*/
void updateJavaTask(String param) throws Exception;
/**
* 功能描述 删除任务
* @author gml
* @date 2020-04-14 11:28
* @param jobName
* @return void
*/
void deleteJavaTask(String jobName);
/**
* 功能描述 是否存在任务
* @author gml
* @date 2020-04-14 11:54
* @param jobName
* @return java.lang.Boolean
*/
Boolean existJavaTask(String jobName);
/**
* 功能描述 根据jobName获取日志
* @author gml
* @date 2020-04-14 16:10
* @param param
* @return java.util.List<com.byit.model.JobTaskRunLog>
*/
List<JobTaskRunLog> loadLogByJobName(String param);
List<JobTaskRunLog> loadLogByTaskName(String param);
void runJavaTask(String jobName);
/**
* 获取当前任务的状态
* @param jobNames
* @return
*/
Map<String, JobTaskRunLog> loadCurrentStatusByJobName(String jobNames);
void killNode(Integer logId);
}
package com.byit.service;
/**
* @description: 工作空间操作API的业务逻辑处理接口
* @author: gml
* @create: 2019-12-31 11:21
*/
public interface ApiWorkspaceService {
void add(String workspaceName);
/**
* 判断工作空间是否存在
* @param workspaceName
* @return
*/
Boolean exist(String workspaceName);
}
package com.byit.service;
import com.byit.executor.api.ScriptExecutorService;
import com.byit.dto.executor.DispatchResponseDto;
import com.byit.rpc.remoting.invoker.annotation.RpcReference;
import org.springframework.stereotype.Service;
@Service
public class TestServiceImpl {
@RpcReference
private ScriptExecutorService scriptExecutorService;
public DispatchResponseDto test(){
return scriptExecutorService.runScript(null);
}
}
package com.byit.service.impl;
import com.byit.mapper.WorkspaceMapper;
import com.byit.model.Workspace;
import com.byit.service.ApiWorkspaceService;
import com.byit.utils.ValidationUtil;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import javax.annotation.Resource;
import java.util.Date;
/**
* @description: 工作空间操作API的业务逻辑处理实现类
* @author: gml
* @create: 2019-12-31 11:22
*/
@Service
@Transactional(rollbackFor = Exception.class)
public class ApiWorkspaceServiceImpl implements ApiWorkspaceService {
@Resource
private WorkspaceMapper workspaceMapper;
@Override
public void add(String workspaceName) {
ValidationUtil.dataNotBank(workspaceName, "工作空间名称为空!");
Workspace workspace = workspaceMapper.getByName(workspaceName);
ValidationUtil.isTrueValidation(null != workspace, "工作空间已存在!");
workspace = new Workspace();
workspace.setWorkspaceName(workspaceName);
workspace.setAddTime(new Date());
workspaceMapper.insertSelective(workspace);
}
@Override
public Boolean exist(String workspaceName) {
ValidationUtil.dataNotBank(workspaceName, "工作空间名称不允许为空!");
Workspace workspace = workspaceMapper.getByName(workspaceName);
if (null == workspace){
return false;
}
return true;
}
}
package com.byit.view;
import com.byit.dto.FlowConditionDto;
import com.byit.model.vo.FlowImportantAllVo;
import com.byit.model.vo.FlowVersionDetailedVo;
import com.byit.model.vo.FlowVersionVo;
import com.byit.service.FlowVersionService;
import com.byit.model.vo.FlowVersionVoDep;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.List;
/**
* 工作流视图界面
* @author huangfu
*/
@RestController
@RequestMapping("view/flow/version")
public class FlowVersionViewController {
private final FlowVersionService flowVersionService;
public FlowVersionViewController(FlowVersionService flowVersionService) {
this.flowVersionService = flowVersionService;
}
@PostMapping("getAllFlowVersion")
public List<FlowVersionVoDep> getAllFlowVersion(){
return flowVersionService.findAllFlow();
}
@PostMapping("findAllVersion")
public List<FlowVersionVo> findAllVersion(@RequestBody FlowConditionDto flowConditionDto){
return flowVersionService.findAllVersion(flowConditionDto);
}
@RequestMapping("findAllByFlowId")
public List<FlowVersionDetailedVo> findAllByFlowId(Integer flowId) {
return flowVersionService.findAllByFlowId(flowId);
}
}
package com.byit.view;
import com.byit.dto.FlowConditionDto;
import com.byit.model.Flow;
import com.byit.model.vo.FlowImportantAllVo;
import com.byit.model.vo.FlowViewVo;
import com.byit.service.FlowService;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.List;
/**
* 当前运行工作流
* @author huangfu
*/
@RestController
@RequestMapping("view/flow/")
public class FlowViewController {
private final FlowService flowService;
public FlowViewController(FlowService flowService) {
this.flowService = flowService;
}
@PostMapping("findAllThisVersionFlow")
public List<Flow> findAllThisVersionFlow(@RequestBody FlowConditionDto flowConditionDto){
return flowService.findAllFlow(flowConditionDto);
}
@PostMapping("findAllFlowViewVo")
public List<FlowViewVo> findAllFlowViewVo(@RequestBody FlowConditionDto flowConditionDto){
return flowService.findAllFlowViewVo(flowConditionDto);
}
/**
* 查询工作流摘要
* @return
*/
@PostMapping("findAllFlowImportant")
public List<FlowImportantAllVo> findAllFlowImportant(@RequestBody FlowConditionDto flowConditionDto) {
return flowService.findAllFlowImportant(flowConditionDto);
}
}
package com.byit.view;
import com.byit.dto.plugin.JavaTask;
import com.byit.model.vo.JavaTaskLogImportantVo;
import com.byit.service.JavaTaskService;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.List;
/**
* @author huangfu
*/
@RestController
@RequestMapping("view/java")
public class JavaTaskViewController {
private final JavaTaskService javaTaskService;
public JavaTaskViewController(JavaTaskService javaTaskService) {
this.javaTaskService = javaTaskService;
}
@RequestMapping("findAll")
public List<JavaTask> findAll(String jobName){
return javaTaskService.findAll(jobName);
}
@RequestMapping("findAllJavaTaskLogImportantVo")
public List<JavaTaskLogImportantVo> findAllJavaTaskLogImportantVo(String jobName) {
return javaTaskService.findAllJavaTaskLogImportantVo(jobName);
}
}
package com.byit.view;
import com.byit.model.NodeVersion;
import com.byit.model.vo.FlowNodeVersionVo;
import com.byit.service.NodeVersionService;
import org.apache.ibatis.annotations.Param;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.List;
/**
* @author huangfu
*/
@RestController
@RequestMapping("view/node")
public class NodeVersionViewController {
private final NodeVersionService nodeVersionService;
public NodeVersionViewController(NodeVersionService nodeVersionService) {
this.nodeVersionService = nodeVersionService;
}
@RequestMapping("/getFlowNodeVersionVo")
public List<FlowNodeVersionVo> getFlowNodeVersionVo(){
return nodeVersionService.findAll();
}
@RequestMapping("/findAllNodeVersionByFlowId")
public List<NodeVersion> findAllNodeVersionByFlowId(Integer id){
return nodeVersionService.findAllByFlowId(id);
}
}
package com.byit.view;
import com.byit.model.Node;
import com.byit.service.FlowService;
import com.byit.service.NodeService;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.List;
/**
* @author huangfu
*/
@RestController
@RequestMapping("view/node")
public class NodeViewController {
private final NodeService nodeService;
public NodeViewController(NodeService nodeService) {
this.nodeService = nodeService;
}
@RequestMapping("findNodes")
public List<Node> findNodes(Integer flowId){
return nodeService.findNodeByFlowIdAndVersionName(flowId);
}
}
package com.byit.view;
import com.byit.model.JobTaskRunLog;
import com.byit.model.vo.RunLogVo;
import com.byit.service.JobTaskRunLogService;
import org.apache.ibatis.annotations.Param;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.List;
/**
* @author huangfu
*/
@RestController
@RequestMapping("view/log")
public class RunLogViewController {
private final JobTaskRunLogService jobTaskRunLogService;
public RunLogViewController(JobTaskRunLogService jobTaskRunLogService) {
this.jobTaskRunLogService = jobTaskRunLogService;
}
@RequestMapping("findAllRunLog")
public List<RunLogVo> findAllRunLogIdByFlowIdAndRunId(@Param("flowId") Integer flowId, @Param("runId") String runId){
return jobTaskRunLogService.findAllRunLogIdByFlowIdAndRunId(flowId,runId);
}
@RequestMapping("findNodeLogByHandlerName")
public List<RunLogVo> findNodeLogByHandlerName(String handlerName){
return jobTaskRunLogService.findNodeLogByHandlerName(handlerName);
}
@RequestMapping("findImmediatelyNode")
public List<JobTaskRunLog> findImmediatelyNode(String nodeName) {
return jobTaskRunLogService.findImmediatelyNode(nodeName);
}
}
package com.byit.view;
import com.byit.dto.FlowConditionDto;
import com.byit.model.RunRecording;
import com.byit.model.vo.RunRecordingViewVo;
import com.byit.service.RunRecordingService;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.List;
/**
* 运行记录视图查询
* @author huangfu
*/
@RestController
@RequestMapping("view/runRecording")
public class RunRecordingViewController {
private final RunRecordingService runRecordingService;
public RunRecordingViewController(RunRecordingService runRecordingService) {
this.runRecordingService = runRecordingService;
}
@RequestMapping("findAll")
public List<RunRecording> findAll(@RequestBody FlowConditionDto flowConditionDto){
return runRecordingService.findAll(flowConditionDto);
}
@RequestMapping("findAllRunIng")
public List<RunRecording> findAllRunIng(@RequestBody FlowConditionDto flowConditionDto) {
return runRecordingService.findAllRunIng(flowConditionDto);
}
@RequestMapping("findAllRunRecordingViewVo")
public List<RunRecordingViewVo> findAllRunRecordingViewVo(@RequestBody FlowConditionDto flowConditionDto){
return runRecordingService.findAllRunRecordingViewVo(flowConditionDto);
}
@RequestMapping("findAllError")
public List<RunRecordingViewVo> findAllError(@RequestBody FlowConditionDto flowConditionDto){
return runRecordingService.findAllError(flowConditionDto);
}
}
package com.byit.view;
import com.byit.model.Workspace;
import com.byit.service.WorkspaceService;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.List;
/**
* @author huangfu
*/
@RestController
@RequestMapping("view/workspace")
public class WorkspaceViewController {
private final WorkspaceService workspaceService;
public WorkspaceViewController(WorkspaceService workspaceService) {
this.workspaceService = workspaceService;
}
@RequestMapping("findAllWorkspace")
public List<Workspace> findAllWorkspace(){
return workspaceService.findAll();
}
}
<?xml version="1.0" encoding="UTF-8" ?>
<!DOCTYPE generatorConfiguration PUBLIC
"-//mybatis.org//DTD MyBatis Generator Configuration 1.0//EN"
"http://mybatis.org/dtd/mybatis-generator-config_1_0.dtd" >
<generatorConfiguration>
<context id="context" targetRuntime="MyBatis3">
<!-- 识别关键字 -->
<property name="autoDelimitKeywords" value="true"/>
<property name="beginningDelimiter" value="`"/>
<property name="endingDelimiter" value="`"/>
<!-- 配置相关插件-->
<plugin type="org.mybatis.generator.plugins.SerializablePlugin"/>
<plugin type="com.byit.plugin.DefaultNamePlugin"/>
<plugin type="com.byit.plugin.LombokAnnotationPlugin"/>
<!-- 统一格式-->
<commentGenerator type="com.byit.plugin.NormalCommentGenerator">
<property name="suppressAllComments" value="false"/>
<property name="addRemarkComments" value="true"/>
<property name="dateFormat" value="yyyy-MM-dd"/>
</commentGenerator>
<!--<jdbcConnection driverClass="com.mysql.jdbc.Driver"
connectionURL="jdbc:mysql://10.1.2.104:3306/byitact?tinyInt1isBit=false&amp;
useUnicode=true&amp;characterEncoding=UTF-8&amp;zeroDateTimeBehavior=convertToNull"
userId="dev_oper" password="0755oper"/>-->
<jdbcConnection driverClass="com.mysql.jdbc.Driver"
connectionURL="jdbc:mysql://10.0.10.118:3306/myth-job?tinyInt1isBit=false&amp;
useUnicode=true&amp;characterEncoding=UTF-8&amp;zeroDateTimeBehavior=convertToNull"
userId="root" password="123456"/>
<!-- 处理TINYINT(1) 等 -->
<javaTypeResolver type="com.byit.plugin.DefaultJavaTypeResolverDefault">
<property name="forceBigDecimals" value="true"/>
</javaTypeResolver>
<!-- domaim=po持久化对象-->
<javaModelGenerator targetPackage="com.byit.model" targetProject="src/main/java">
<property name="enableSubPackages" value="false"/>
<property name="trimStrings" value="true"/>
</javaModelGenerator>
<!-- xml生成地址 -->
<sqlMapGenerator targetPackage="mapper" targetProject="src/main/resources">
<property name="enableSubPackages" value="false"/>
<property name="rootInterface" value="idata.dmp.dc.mapper"/>
</sqlMapGenerator>
<!-- mapper接口地址 -->
<javaClientGenerator targetPackage="com.byit.mapper" targetProject="src/main/java"
type="XMLMAPPER">
<property name="enableSubPackages" value="false"/>
</javaClientGenerator>
<table tableName="waiting_record" domainObjectName="WaitingRecord" />
<table tableName="waiting_task" domainObjectName="WaitingTask" />
<!--<table tableName="t_publish_result" domainObjectName="PublishResult" />
<table tableName="t_publish_approve" domainObjectName="PublishApprove" />-->
</context>
</generatorConfiguration>
███╗ ███╗██╗ ██╗████████╗██╗ ██╗ ██╗ ██████╗ ██████╗
████╗ ████║╚██╗ ██╔╝╚══██╔══╝██║ ██║ ██║██╔═══██╗██╔══██╗
██╔████╔██║ ╚████╔╝ ██║ ███████║ █████╗ ██║██║ ██║██████╔╝
██║╚██╔╝██║ ╚██╔╝ ██║ ██╔══██║ ╚════╝ ██ ██║██║ ██║██╔══██╗
██║ ╚═╝ ██║ ██║ ██║ ██║ ██║ ╚█████╔╝╚██████╔╝██████╔╝
╚═╝ ╚═╝ ╚═╝ ╚═╝ ╚═╝ ╚═╝ ╚════╝ ╚═════╝ ╚═════╝
spring:
datasource:
driver-class-name: com.mysql.jdbc.Driver
url: jdbc:mysql://10.0.120.30:3307/myth-job?Unicode=true&characterEncoding=UTF-8&useSSL=true
username: root
password: root
redis:
database: 0
host: 10.0.120.208
password:
port: 6379
timeout: 3000
pool:
max-active: 8
max-idle: 8
max-wait: -1
min-idle: 0
mail:
host: smtp.byitgroup.com
username: ddmp@byitgroup.com
password: Byit12345
default-encoding: UTF-8
protocol: smtp
properties:
from: ddmp@byitgroup.com
thymeleaf:
prefix: classpath:/templates/
suffix: .html
cache: false
enabled: true
check-template: true
check-template-location: true
encoding: utf-8
mode: HTML5
mybatis:
mapper-locations: /mapper/*.xml
myth-rpc:
registry:
address: http://10.0.120.208:8080/myth-register
env: dev
biz: byit-myth-job
logging:
path: /data/mythjob
file: myth_log_file
authentication:
user:
header-name: token
expire: 43200 # 外部token有效期为12小时
pub-key: client/pub.key # 解密
file:
system:
ip: 10.0.120.2
port: 88
myth-job:
filestystem: FASTDFS
snapshoot-date: 60 #快照的保存时间 单位天
fdfs:
so-timeout: 1500
http-port: 88
connect-timeout: 600
pool:
jmx-enabled: false
tracker-list:
- 10.0.120.2:22122
myth:
plugin:
env: ${myth-rpc.registry.env}
biz: ${myth-rpc.registry.biz}
register:
url: ${myth-rpc.registry.address}
######################哨兵模式#####################
#redisson:
# master-name: myMaster
# sentinel-addresses:
# - 127.0.0.1:3306
# - 127.7.7.7:6379
# timeout: 60000
######################哨兵模式#####################
######################单机环境#####################
redisson:
address: redis://10.0.120.208:6379
database: 3
timeout: 60000
######################单机环境#####################
######################redis分布式锁#####################
lock:
type: redis
######################redis分布式锁#####################
######################DB行锁#####################
#lock:
# type: db
######################DB行锁#####################
\ No newline at end of file
spring:
datasource:
driver-class-name: com.mysql.jdbc.Driver
url: jdbc:mysql://10.0.10.118:3306/myth-job?Unicode=true&characterEncoding=UTF-8&useSSL=true
username: root
password: 123456
mail:
host: smtp.byitgroup.com
username: ddmp@byitgroup.com
password: Byit12345
default-encoding: UTF-8
protocol: smtp
properties:
from: ddmp@byitgroup.com
redis:
database: 0
host: 10.0.120.208
port: 6379
password:
mybatis:
mapper-locations: /mapper/*.xml
myth-rpc:
registry:
address: http://localhost:8080/myth-register
env: huangfu
biz: byit-myth-job
logging:
path: /data/mythjob
file: myth_log_file
authentication:
user:
header-name: token
expire: 43200 # 外部token有效期为12小时
pub-key: client/pub.key # 解密
file:
system:
ip: 10.0.120.2
port: 88
myth-job:
filestystem: FASTDFS
snapshoot-date: 60 #快照的保存时间 单位天
fdfs:
so-timeout: 1500
connect-timeout: 600
http-port: 88
pool:
jmx-enabled: false
tracker-list:
- 10.0.120.2:22122
myth:
plugin:
env: ${myth-rpc.registry.env}
biz: ${myth-rpc.registry.biz}
register:
url: ${myth-rpc.registry.address}
######################哨兵模式#####################
#redisson:
# master-name: myMaster
# sentinel-addresses:
# - 127.0.0.1:3306
# - 127.7.7.7:6379
# timeout: 60000
######################哨兵模式#####################
######################单机环境#####################
redisson:
address: redis://10.0.120.208:6379
database: 3
timeout: 60000
######################单机环境#####################
######################redis分布式锁#####################
lock:
type: redis
######################redis分布式锁#####################
######################DB行锁#####################
#lock:
# type: db
######################DB行锁#####################
\ No newline at end of file
spring:
datasource:
driver-class-name: com.mysql.jdbc.Driver
url: jdbc:mysql://10.0.120.30:3307/myth-job?Unicode=true&characterEncoding=UTF-8&useSSL=true
username: root
password: root
redis:
database: 0
host: 10.0.120.208
port: 6379
password:
timeout: 3000
pool:
max-active: 8
max-idle: 8
max-wait: -1
min-idle: 0
mail:
host: smtp.byitgroup.com
username: ddmp@byitgroup.com
password: Byit12345
default-encoding: UTF-8
protocol: smtp
properties:
from: ddmp@byitgroup.com
mybatis:
mapper-locations: /mapper/*.xml
myth-rpc:
registry:
address: http://10.0.120.208:8080/myth-register
env: pro
biz: byit-myth-job
logging:
path: /data/mythjob
file: myth_log_file
authentication:
user:
header-name: token
expire: 43200 # 外部token有效期为12小时
pub-key: client/pub.key # 解密
file:
system:
ip: 10.0.120.2
port: 88
myth-job:
filestystem: FASTDFS
snapshoot-date: 60 #快照的保存时间 单位天
fdfs:
so-timeout: 1500
connect-timeout: 600
http-port: 8888
pool:
jmx-enabled: false
tracker-list:
- 10.0.120.216:22122
- 10.0.120.217:22122
- 10.0.120.218:22122
myth:
plugin:
env: ${myth-rpc.registry.env}
biz: ${myth-rpc.registry.biz}
register:
url: ${myth-rpc.registry.address}
######################哨兵模式#####################
#redisson:
# master-name: myMaster
# sentinel-addresses:
# - 127.0.0.1:3306
# - 127.7.7.7:6379
# timeout: 60000
######################哨兵模式#####################
######################单机环境#####################
redisson:
address: redis://10.0.120.208:6379
database: 3
timeout: 60000
######################单机环境#####################
######################redis分布式锁#####################
lock:
type: redis
######################redis分布式锁#####################
######################DB行锁#####################
#lock:
# type: db
######################DB行锁#####################
\ No newline at end of file
spring:
datasource:
driver-class-name: com.mysql.jdbc.Driver
url: jdbc:mysql://${CM_IP}/${CM_DB}?Unicode=true&characterEncoding=UTF-8&useSSL=true
username: ${CM_USER}
password: ${CM_PWD}
redis:
database: 0
host: ${Redis_IP}
port: ${Redis_port}
password:
timeout: 3000
pool:
max-active: 8
max-idle: 8
max-wait: -1
min-idle: 0
mail:
host: smtp.byitgroup.com
username: ddmp@byitgroup.com
password: Byit12345
default-encoding: UTF-8
protocol: smtp
properties:
from: ddmp@byitgroup.com
mybatis:
mapper-locations: /mapper/*.xml
myth-rpc:
registry:
address: http://${Eureka_IP}/myth-register
env: test
biz: byit-myth-job
logging:
path: /data/mythjob
file: myth_log_file
authentication:
user:
header-name: token
expire: 43200 # 外部token有效期为12小时
pub-key: client/pub.key # 解密
file:
system:
ip: ${FILE_IP}
port: ${FILW_PORT}
myth-job:
filestystem: FASTDFS
snapshoot-date: 60 #快照的保存时间 单位天
fdfs:
so-timeout: 1500
connect-timeout: 600
http-port: ${FDFS_PORT}
pool:
jmx-enabled: false
tracker-list:
- ${FDFS_TK1}
- ${FDFS_TK2}
- ${FDFS_TK3}
myth:
plugin:
env: ${myth-rpc.registry.env}
biz: ${myth-rpc.registry.biz}
register:
url: ${myth-rpc.registry.address}
######################哨兵模式#####################
#redisson:
# master-name: myMaster
# sentinel-addresses:
# - 127.0.0.1:3306
# - 127.7.7.7:6379
# timeout: 60000
######################哨兵模式#####################
######################单机环境#####################
redisson:
address: redis://${Redis_IP}:${Redis_port}
database: 3
timeout: 60000
######################单机环境#####################
######################redis分布式锁#####################
lock:
type: redis
######################redis分布式锁#####################
######################DB行锁#####################
#lock:
# type: db
######################DB行锁#####################
\ No newline at end of file
server:
port: 8998
context-path: /myth-job-admin
spring:
profiles:
active: local
application:
name: myth-job-admin
<?xml version="1.0" encoding="UTF-8"?>
<configuration>
<appender name="console" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{50} - %msg%n</pattern>
</encoder>
</appender>
<appender name="file-debug" class="ch.qos.logback.core.rolling.RollingFileAppender">
<rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
<fileNamePattern>${LOG_PATH}/${LOG_FILE}_debug.%d{yyyy-MM-dd}.log</fileNamePattern>
</rollingPolicy>
<encoder>
<pattern>%d{HH:mm:ss.SSS} %contextName [%thread] %-5level %logger{36} - %msg%n</pattern>
</encoder>
<filter class="ch.qos.logback.classic.filter.LevelFilter">
<level>DEBUG</level>
<onMatch>ACCEPT</onMatch>
<onMismatch>DENY</onMismatch>
</filter>
</appender>
<appender name="file-info" class="ch.qos.logback.core.rolling.RollingFileAppender">
<rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
<fileNamePattern>${LOG_PATH}/${LOG_FILE}_info.%d{yyyy-MM-dd}.log</fileNamePattern>
</rollingPolicy>
<encoder>
<pattern>%d{HH:mm:ss.SSS} %contextName [%thread] %-5level %logger{36} - %msg%n</pattern>
</encoder>
<filter class="ch.qos.logback.classic.filter.LevelFilter">
<level>INFO</level>
<onMatch>ACCEPT</onMatch>
<onMismatch>DENY</onMismatch>
</filter>
</appender>
<appender name="file-error" class="ch.qos.logback.core.rolling.RollingFileAppender">
<rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
<fileNamePattern>${LOG_PATH}/${LOG_FILE}_error.%d{yyyy-MM-dd}.log</fileNamePattern>
</rollingPolicy>
<encoder>
<pattern>%d{HH:mm:ss.SSS} %contextName [%thread] %-5level %logger{36} - %msg%n</pattern>
</encoder>
<filter class="ch.qos.logback.classic.filter.LevelFilter">
<level>ERROR</level>
<onMatch>ACCEPT</onMatch>
<onMismatch>DENY</onMismatch>
</filter>
</appender>
<appender name="file-myth-job" class="ch.qos.logback.core.rolling.RollingFileAppender">
<rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
<fileNamePattern>${LOG_PATH}/${LOG_FILE}_myth_job.%d{yyyy-MM-dd}.log</fileNamePattern>
</rollingPolicy>
<encoder>
<pattern>%d{HH:mm:ss.SSS} %contextName [%thread] %-5level %logger{36} - %msg%n</pattern>
</encoder>
</appender>
<logger name="com.byit.thread.helper" level="debug" additivity="false">
<appender-ref ref="console"/>
</logger>
<!--<logger name="com.byit.factory.DaemonScanThreadRunHelperRedisLock" level="debug" additivity="false">
<appender-ref ref="console"/>
<appender-ref ref="file-debug"/>
<appender-ref ref="file-info"/>
</logger>-->
<logger name="com.byit.selector" level="info" additivity="false">
<appender-ref ref="console"/>
</logger>
<logger name="com.byit.task" level="info" additivity="false">
<appender-ref ref="console"/>
</logger>
<logger name="org.springframework.web" level="debug">
<appender-ref ref="console" />
<appender-ref ref="file-debug" />
</logger>
<logger name="com.ibatis" level="debug">
<appender-ref ref="console" />
<appender-ref ref="file-debug" />
</logger>
<logger name="com.ibatis.common.jdbc.SimpleDataSource" level="debug" >
<appender-ref ref="console" />
<appender-ref ref="file-debug" />
</logger>
<logger name="com.ibatis.common.jdbc.ScriptRunner" level="debug">
<appender-ref ref="console" />
<appender-ref ref="file-debug" />
</logger>
<logger name="com.ibatis.sqlmap.engine.impl.SqlMapClientDelegate" level="debug" >
<appender-ref ref="console" />
<appender-ref ref="file-debug" />
</logger>
<logger name="java.sql.Connection" level="debug" >
<appender-ref ref="console" />
<appender-ref ref="file-debug" />
</logger>
<logger name="java.sql.Statement" level="debug" >
<appender-ref ref="console" />
<appender-ref ref="file-debug" />
</logger>
<logger name="java.sql.PreparedStatement" level="debug" >
<appender-ref ref="console" />
<appender-ref ref="file-debug" />
</logger>
<root level="info">
<appender-ref ref="console" />
<appender-ref ref="file-debug" />
<appender-ref ref="file-info" />
<appender-ref ref="file-error" />
</root>
</configuration>
\ No newline at end of file
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<parent>
<artifactId>byit-myth-core</artifactId>
<groupId>myth-job</groupId>
<version>1.0-SNAPSHOT</version>
</parent>
<modelVersion>4.0.0</modelVersion>
<groupId>myth-job</groupId>
<artifactId>myth-admin-core</artifactId>
<dependencies>
<dependency>
<groupId>io.netty</groupId>
<artifactId>netty-all</artifactId>
</dependency>
<dependency>
<groupId>myth-job</groupId>
<artifactId>myth-executor-api</artifactId>
</dependency>
<dependency>
<groupId>myth-job</groupId>
<artifactId>myth-core-common</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-mail</artifactId>
</dependency>
<dependency>
<groupId>org.mybatis.spring.boot</groupId>
<artifactId>mybatis-spring-boot-starter</artifactId>
</dependency>
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-jdbc</artifactId>
</dependency>
<dependency>
<groupId>io.springfox</groupId>
<artifactId>springfox-swagger2</artifactId>
</dependency>
<dependency>
<groupId>io.springfox</groupId>
<artifactId>springfox-swagger-ui</artifactId>
</dependency>
<dependency>
<groupId>myth-job</groupId>
<artifactId>byit-myth-rpc</artifactId>
</dependency>
<dependency>
<groupId>myth-job</groupId>
<artifactId>myth-dto-core</artifactId>
<version>1.0-SNAPSHOT</version>
<scope>compile</scope>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-redis</artifactId>
<version>1.3.2.RELEASE</version>
</dependency>
<dependency>
<groupId>redis.clients</groupId>
<artifactId>jedis</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-thymeleaf</artifactId>
</dependency>
<dependency>
<groupId>myth-job</groupId>
<artifactId>plugin-spring-boot-starter</artifactId>
<version>1.0-SNAPSHOT</version>
</dependency>
<dependency>
<groupId>org.redisson</groupId>
<artifactId>redisson</artifactId>
<version>3.8.2</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.mybatis.generator</groupId>
<artifactId>mybatis-generator-maven-plugin</artifactId>
<version>1.3.5</version>
<dependencies>
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<version>5.1.44</version>
</dependency>
<dependency>
<groupId>org.mybatis.spring.boot</groupId>
<artifactId>mybatis-spring-boot-starter</artifactId>
<version>1.2.0</version>
</dependency>
<dependency>
<groupId>com.byit</groupId>
<artifactId>byit-mybatis-plugin</artifactId>
<version>1.0.0-SNAPSHOT</version>
</dependency>
</dependencies>
<configuration>
<configurationFile>${basedir}/src/main/resources/Generator-config.xml</configurationFile>
<verbose>true</verbose>
<overwrite>true</overwrite>
<skip>false</skip>
</configuration>
</plugin>
</plugins>
</build>
</project>
\ No newline at end of file
package com.byit.conf;
import org.mybatis.spring.annotation.MapperScan;
import org.springframework.context.annotation.*;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
/**
* @program: byit-myth-job->AppConf
* @description: 稳健配置
* @author: huangfu
* @date: 2019/12/9 14:51
**/
@Configuration
@MapperScan("com.byit.mapper")
public class MythJobAutoConfigure {
/**
* 日志的回调线程池,主要是将执行结果写到库里面
* LinkedBlockingQueue 不指定容量就变成了无界队列
* 判断核心线程数是否已满,核心线程数大小和corePoolSize参数有关,未满则创建线程执行任务
* 若核心线程池已满,判断队列是否满,队列是否满和workQueue参数有关,若未满则加入队列中
* 若队列已满,判断线程池是否已满,线程池是否已满和maximumPoolSize参数有关,若未满创建线程执行任务
* 若线程池已满,则采用拒绝策略处理无法执执行的任务,拒绝策略和handler参数有关
*/
public static final ThreadPoolExecutor LOG_CALLBACK = new ThreadPoolExecutor(
10,
100,
60L,
TimeUnit.SECONDS,
new LinkedBlockingQueue<Runnable>(100),
r ->new Thread(r, "MythJob Thread of Job Run Log Callback Warehouse-" + r.hashCode()));
/**
* 高级任务的线程池
*/
public static final ThreadPoolExecutor ADVANCED_JOB_THREAD_POOL = new ThreadPoolExecutor(
10,
100,
60L,
TimeUnit.SECONDS,
new LinkedBlockingQueue<Runnable>(200),
r ->new Thread(r, "MythJob Thread of Job advanced run pool-" + r.hashCode()));
/**
* 低级任务的线程池
*/
public static final ThreadPoolExecutor LOW_LEVEL_JOB_THREAD_POOL = new ThreadPoolExecutor(
20,
200,
60L,
TimeUnit.SECONDS,
new LinkedBlockingQueue<Runnable>(300),
r ->new Thread(r, "MythJob Thread of Job low level run pool-" + r.hashCode()));
}
package com.byit.conf;
import com.byit.factory.DaemonScanThreadRunHelper;
import com.byit.factory.DaemonScanThreadRunHelperDbLock;
import com.byit.thread.BaseThreadRunHelper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationContext;
import org.springframework.stereotype.Component;
import java.util.Map;
/**
* @program: byit-myth-job->MythJobScheduler
* @description: 扫描线程的生命周期
* @author: huangfu
* @date: 2019/12/9 16:53
**/
@Component
@Slf4j
public class MythJobScheduler implements InitializingBean, DisposableBean {
private final ApplicationContext applicationContext;
private final DaemonScanThreadRunHelper daemonScanThreadRunHelper;
@Autowired
public MythJobScheduler(ApplicationContext applicationContext, DaemonScanThreadRunHelper daemonScanThreadRunHelper) {
this.applicationContext = applicationContext;
this.daemonScanThreadRunHelper = daemonScanThreadRunHelper;
}
/**
* 销毁方法
*/
@Override
public void destroy() {
daemonScanThreadRunHelper.logoutThreadGroup();
}
/**
* 初始化方法
*/
@Override
public void afterPropertiesSet() {
helperAutoAdd();
//启用扫描线程
daemonScanThreadRunHelper.runDaemonThreads();
}
/**
* 自动注册
*/
public void helperAutoAdd(){
Map<String, BaseThreadRunHelper> beansOfType = applicationContext.getBeansOfType(BaseThreadRunHelper.class);
beansOfType.forEach((key, value) -> {
log.info("-------------自动加载帮助类,{}-------------",key);
daemonScanThreadRunHelper.addThread(value);
});
}
}
package com.byit.conf;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.redis.connection.RedisConnectionFactory;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.data.redis.serializer.StringRedisSerializer;
@Configuration
public class RedisConfig {
@Bean
StringRedisTemplate stringRedisTemplate(RedisConnectionFactory connectionFactory){
StringRedisTemplate stringRedisTemplate = new StringRedisTemplate();
stringRedisTemplate.setConnectionFactory(connectionFactory);
stringRedisTemplate.setDefaultSerializer(new StringRedisSerializer());
return stringRedisTemplate;
}
}
package com.byit.conf;
import com.byit.conf.properties.RedissonProperties;
import com.byit.util.lock.DistributedLocker;
import com.byit.util.lock.RedissLockUtil;
import com.byit.util.lock.RedissonDistributedLocker;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.redisson.Redisson;
import org.redisson.api.RedissonClient;
import org.redisson.config.Config;
import org.redisson.config.SentinelServersConfig;
import org.redisson.config.SingleServerConfig;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
* 分布式锁配置对象
* @author huangfu
*/
@Configuration
@ConditionalOnClass(Config.class)
@EnableConfigurationProperties(RedissonProperties.class)
@Slf4j
public class RedissonAutoConfiguration {
@Autowired
private RedissonProperties redssionProperties;
/**
* 哨兵模式自动装配
* @return
*/
@Bean
@ConditionalOnProperty(name="redisson.master-name")
public RedissonClient redissonSentinel() {
Config config = new Config();
SentinelServersConfig serverConfig = config.useSentinelServers()
.addSentinelAddress(redssionProperties.getSentinelAddresses())
.setMasterName(redssionProperties.getMasterName())
.setTimeout(redssionProperties.getTimeout())
.setMasterConnectionPoolSize(redssionProperties.getMasterConnectionPoolSize())
.setSlaveConnectionPoolSize(redssionProperties.getSlaveConnectionPoolSize())
.setDatabase(redssionProperties.getDatabase());
if(StringUtils.isNotBlank(redssionProperties.getPassword())) {
serverConfig.setPassword(redssionProperties.getPassword());
}
return Redisson.create(config);
}
/**
* 单机模式自动装配
* @return
*/
@Bean
@ConditionalOnProperty(name="redisson.address")
public RedissonClient redissonSingle() {
Config config = new Config();
SingleServerConfig serverConfig = config.useSingleServer()
.setAddress(redssionProperties.getAddress())
.setTimeout(redssionProperties.getTimeout())
.setConnectionPoolSize(redssionProperties.getConnectionPoolSize())
.setConnectionMinimumIdleSize(redssionProperties.getConnectionMinimumIdleSize())
.setDatabase(redssionProperties.getDatabase());
if(StringUtils.isNotBlank(redssionProperties.getPassword())) {
serverConfig.setPassword(redssionProperties.getPassword());
}
return Redisson.create(config);
}
/**
* 装配locker类,并将实例注入到RedissLockUtil中
* @return
*/
@Bean
public DistributedLocker distributedLocker(RedissonClient redissonSingle) {
RedissonDistributedLocker locker = new RedissonDistributedLocker();
locker.setRedissonClient(redissonSingle);
log.info("----------redis分布式锁{},被设置加载---------",locker);
RedissLockUtil.setLocker(locker);
DistributedLocker redissLock = RedissLockUtil.getRedissLock();
if(redissLock == null) {
throw new RuntimeException("redis分布式锁设置异常,请联系调度中心开发团队!");
}
return locker;
}
}
package com.byit.conf;
import cn.hutool.http.HttpRequest;
import com.alibaba.fastjson.JSON;
import com.byit.dto.web.ResponseResult;
import com.byit.packet.response.PluginRpcResponsePacket;
import com.byit.param.ResultCallback;
import io.netty.channel.ChannelHandlerContext;
import static com.alibaba.fastjson.serializer.SerializerFeature.WriteClassName;
/**
* @author huangfu
*/
public class RpcResultHttpCallback implements ResultCallback {
private String gatewayIpAndPort;
public RpcResultHttpCallback(String gatewayIpAndPort) {
this.gatewayIpAndPort = "http://"+gatewayIpAndPort+"/myth-job-admin/api/callback/callbackRes";
}
@Override
public void resultCallback(ChannelHandlerContext ctx, PluginRpcResponsePacket pluginRpcResponsePacket) {
Object result = pluginRpcResponsePacket.getResult();
String responseStr = JSON.toJSONString(pluginRpcResponsePacket, WriteClassName);
HttpRequest post = HttpRequest.post(gatewayIpAndPort);
post.header("token", "Token");
post.body("param="+responseStr).execute();
System.out.println("-------消息回复成功-----------");
}
}
package com.byit.conf.properties;
import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;
import java.util.List;
/**
* redisson分布式锁配置类
* @author huangfu
*/
@Data
@ConfigurationProperties(prefix = "redisson")
public class RedissonProperties {
private int timeout = 3000;
private String address;
private String password;
private int connectionPoolSize = 64;
private int connectionMinimumIdleSize=10;
private int slaveConnectionPoolSize = 250;
private int masterConnectionPoolSize = 250;
private String[] sentinelAddresses;
private String masterName;
private int database = 2;
}
package com.byit.dto;
import lombok.Builder;
import lombok.Data;
import java.util.List;
@Data
public class EmailMappingDto {
private String nodeName;
private String startDate;
private String endDate;
private String consuming;
private String runStatus;
private List<String> logPaths;
}
package com.byit.dto;
import lombok.Data;
import lombok.ToString;
import java.io.Serializable;
/**
* @author huangfu
*/
@Data
@ToString
public class FlowConditionDto implements Serializable {
private Integer workspaceId;
private String flowName;
}
package com.byit.dto;
import lombok.Data;
@Data
public class StatisticsConditionDto {
private long startTime;
private long endTime;
private Integer workspaceId;
private String flowName;
}
package com.byit.enums;
import lombok.AllArgsConstructor;
/**
* @program: byit-myth-job->AdminEnums
* @description: 调度中心的异常枚举
* @author: huangfu
* @date: 2019/12/25 15:57
**/
@AllArgsConstructor
public enum AdminEnums implements IEnum {
SEND_EMAIL_FAILURE("100500","邮件发送失败"),
SEND_EMAIL_SUCCESS("100200","邮件发送成功")
;
String code;
String msg;
@Override
public String getCode() {
return this.code;
}
@Override
public String getMsg() {
return this.msg;
}
}
package com.byit.enums;
import lombok.Getter;
@Getter
public enum DagCheckEnum {
LOOP(0, "存在环路"),
FREE(1, "存在游离节点"),
PASS(2, "校验通过"),
STARTNAME_WRONG(3, "开始节点必须是"),
ENDNAME_WRONG(4, "结束节点必须是"),
;
private Integer code;
private String msg;
private DagCheckEnum(Integer code, String msg){
this.code = code;
this.msg = msg;
}
public Integer getCode(){
return this.code;
}
public String getMsg(){
return this.msg;
}
public static String getNameByCode(Integer code) {
for (DagCheckEnum type : DagCheckEnum.values()) {
if (type.getCode().equals(code)) {
return type.getMsg();
}
}
return null;
}
}
package com.byit.enums;
/**
* @author Administrator
*/
public enum EmailEnum {
IS_ALARM_YES("1","已经告警"),
IS_ALARM_NO("0","没有告警"),
EMAIL_NOT_ALARML("0","不告警"),
EMAIL_END_ALARML("1","完成时告警"),
EMAIL_ERROR_ALARML("2","错误时告警"),
EMAIL_SUCCESS_ALARML("3","成功时告警"),
;
private String code;
private String msg;
EmailEnum(String code, String msg) {
this.code = code;
this.msg = msg;
}
public String getCode() {
return code;
}
}
package com.byit.enums;
/**
* @Description 执行状态
* @Author guo_m
* @Date 2020-03-31
*/
public enum ExecuteStatusEnum {
RUNING("0","运行中"),
SUCCESS("1", "成功"),
FAIL("2", "失败"),
REPAIR_SUCCESS("3", "补批成功"),
REPAIR_FAIL("4", "补批失败"),
KILL("5", "杀死"),
PARENT_NODE_FAIL("6","上级节点执行失败")
;
private String code;
private String msg;
private ExecuteStatusEnum(String code, String msg){
this.code = code;
this.msg = msg;
}
public String getCode(){
return this.code;
}
public String getMsg(){
return this.msg;
}
}
package com.byit.enums;
/**
* @description: 工作流属性的枚举类
* @author: gml
* @create: 2019-12-26 15:17
*/
public enum FlowPropertyEnum {
IS_INNER("0", "是内嵌工作流"),
ISNOT_INNER("1", "不是内嵌工作流"),
IS_START("0", "启动"),
NO_START("1", "未启动"),
MANUAL_MODE("2", "手动执行"),
SCHEDULE_MODE("1", "周期执行"),
NO_SCHEDULE("2", "不跟随工作流调度"),
SCHEDULE("1", "跟随工作流调度"),
NO_ALARML("0", "不告警"),
FINISH_ALARML("1", "完成时告警"),
SUCCESS_ALARML("2", "成功时告警"),
FAIL_ALARML("3", "失败时告警"),
IS_CURRENTVERSION("0", "版本表是当前版本的工作流"),
ISNOT_CURRENTVERSION("1", "版本表不是当前版本的工作流"),
FLOW_RUN_ING("2","工作流运行中"),
SCAN("1","标识扫描"),
NOT_SCAN("2","不允许扫描")
;
private String code;
private String name;
private FlowPropertyEnum(String code, String name){
this.code = code;
this.name = name;
}
public String getCode(){
return this.code;
}
public String getName(){
return this.name;
}
}
package com.byit.enums;
/**
* @author huangfu
*/
public enum JobTriggerStatusEnums {
STOP("0","暂停"),
START("1","开始")
;
public String getCode() {
return code;
}
public void setCode(String code) {
this.code = code;
}
public String getName() {
return name;
}
public void setName(String name) {
this.name = name;
}
JobTriggerStatusEnums(String code, String name) {
this.code = code;
this.name = name;
}
private String code;
private String name;
}
package com.byit.enums;
/**
* 节点名称 枚举
* @author huangfu
*/
public enum NodeNameEnum {
START_NODE("start","开始节点"),
END_NODE("end","结束节点"),
;
private String nodeName;
private String details;
public String getNodeName() {
return nodeName;
}
public String getDetails() {
return details;
}
NodeNameEnum(String nodeName, String details) {
this.nodeName = nodeName;
this.details = details;
}
}
package com.byit.enums;
/**
* @description: 节点属性枚举
* @author: gml
* @create: 2019-12-26 15:36
*/
public enum NodePropertyEnum {
IS_VIRTUAL("0", "是虚节点"),
ISNOT_VIRTUAL("1", "不是虚节点"),
ON_FORK("0", "在工作流调度中"),
OFF_FORK("1", "不在工作流调度中"),
ADVANCED_NODE("2","高级节点"),
LOW_LEVEL_NODE("1","低级节点"),
STRONG_NODE("0","强引用"),
WEAK_NODE("1","弱引用")
;
private String code;
private String name;
private NodePropertyEnum(String code, String name){
this.code = code;
this.name = name;
}
public String getCode(){
return this.code;
}
public String getName(){
return this.name;
}
}
package com.byit.enums;
/**
* 节点运行状态枚举
* @author huangfu
*/
public enum NodeRunStatusPropertyEnum {
RUN_ING("0","运行中"),
RUN_SUCCESS("1","成功"),
RUN_FAILURE("2","失败"),
RE_RUN_SUCCESS("3","补批成功"),
RE_RUN_FAILURE("4","补批失败"),
KILL("5","kill"),
PARENT_NODE_FAILED("6","上级节点执行失败"),
NODE_RELY_ERROR("7","节点依赖错误"),
;
private String code;
private String msg;
public String getCode() {
return code;
}
public void setCode(String code) {
this.code = code;
}
public String getMsg() {
return msg;
}
public void setMsg(String msg) {
this.msg = msg;
}
NodeRunStatusPropertyEnum(String code, String msg) {
this.code = code;
this.msg = msg;
}
}
package com.byit.enums;
/**
* @description: 节点类型
* @author: gml
* @create: 2019-12-25 12:01
*/
public enum NodeTypeEnum {
SHELL("SHELL","SCRIPT"),
JAVA("JAVA","JAVA"),
PYTHON("PYTHON","SCRIPT"),
SQL("SQL","SCRIPT"),
SCRIPT("SCRIPT","SCRIPT"),
;
private String code;
private String type;
NodeTypeEnum(String code, String type) {
this.code = code;
this.type = type;
}
public String getType() {
return type;
}
public String getCode(){
return this.code;
}
public NodeTypeEnum getTypeByCode(String code){
for (NodeTypeEnum typeEnum : NodeTypeEnum.values()) {
if (typeEnum.getCode().equals(code)){
return typeEnum;
}
}
return null;
}
}
package com.byit.enums;
/**
* @description: 运行记录的枚举类
* @author huangfu
*/
public enum RunRecordingEnum {
RUN_FLOW_SUCCESS("1","工作流运行成功")
,RUN_FLOW_FAILURE("2","工作流运行失败")
,RUN_FLOW_RE_SUCCESS("3","补批成功")
,RUN_FLOW_RE_FAILURE("4","补批失败")
,RUN_FLOW_KILL("5","工作流进程被杀死")
,FLOW_STATUS_NOT_RUN("1","工作流未开始")
,FLOW_STATUS_RUN_ING("2","工作流运行中")
,FLOW_STATUS_IS_STOP("3","工作流被暂停")
,FLOW_STATUS_IS_END("4","工作流已经完结")
,FAIL_FAST_NO("0","不快速失败")
,FAIL_FAST_YES("1","快速失败")
;
private String code;
private String message;
RunRecordingEnum(String code, String message) {
this.code = code;
this.message = message;
}
public String getCode() {
return code;
}
public String getMessage() {
return message;
}}
package com.byit.enums;
/**
* @Description
* @Author guo_m
* @Date 2020-03-31
*/
public enum ScheduleStatusEnum {
UN_START("1", "未运行"),
STARTING("2", "运行中"),
FINISH("4", "成功"),
STOP("3", "暂停"),
ERROR("5", "失败"),
KILL("6", "杀死"),
;
private String code;
private String msg;
ScheduleStatusEnum(String code, String msg){
this.code = code;
this.msg = msg;
}
public String getCode(){
return this.code;
}
public String getMsg(){
return this.msg;
}
}
package com.byit.event;
import org.springframework.context.ApplicationEvent;
/**
* 完结工作流的事件
* @author huangfu
*/
public class EndFlowEvent extends ApplicationEvent {
private Integer flowId;
/**
* 创建工作流完结的事件
*
* @param source the object on which the event initially occurred (never {@code null})
* @param flowId 工作流id
*/
public EndFlowEvent(Object source,Integer flowId) {
super(source);
this.flowId = flowId;
}
public Integer getFlowId() {
return flowId;
}
}
package com.byit.event;
import org.springframework.context.ApplicationEvent;
/**
* 创建工作流处理完毕事件
* 作用:主要是工作流处理完毕后需要将对应的工作流改为不扫描
* @author huangfu
*/
public class FlowScanEndEvent extends ApplicationEvent {
private Integer flowId;
/**
* 创建工作流添加进实例表完毕后的事件
*
* @param source the object on which the event initially occurred (never {@code null})
* @param flowId 工作流id
*/
public FlowScanEndEvent(Object source,Integer flowId) {
super(source);
this.flowId = flowId;
}
public Integer getFlowId() {
return flowId;
}
}
\ No newline at end of file
package com.byit.exceptions;
/**
* 上级节点执行失败的异常类
* @author huangfu
*/
public class SuperiorNodeRunException extends RuntimeException {
public SuperiorNodeRunException(String message) {
super(message);
}
}
package com.byit.factory;
import com.byit.thread.BaseThreadRunHelper;
import com.byit.util.lock.RedissLockUtil;
import lombok.extern.slf4j.Slf4j;
import java.util.HashSet;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
/**
* @author huangfu
*/
@Slf4j
public abstract class BaseDaemonScanThreadRunHelper implements DaemonScanThreadRunHelper {
/**
* 控制一组守护线程是否运行
*/
public volatile static boolean THREAD_GROUP_STOP = false;
/**
* 初始睡眠时间 单位毫秒
*/
public static final Long INIT_SLEEP_DATE = 10000L;
/**
* 关闭等待时间 单位毫秒
*/
public static final Long CLOSE_WAIT_TIME = 1000L;
/**
* 循环间隔 单位毫秒
*/
public static final Long CYCLE_INTERVAL = 1000L;
/**
* 线程运行必须原料
*/
public static final Map<String,BaseThreadRunHelper> THREAD_RUN_OPT = new ConcurrentHashMap<>(8);
/**
* 线程存储 用于停止线程
*/
public static final Map<String,Thread> THREADS_MAP = new ConcurrentHashMap<>(8);
/**
* 用于磁存储所有的锁
*/
public static final Set<String> LOCK_NAMES = new HashSet<>(8);
/**
* 线程添加 将线程扫描器添加进线程管理池
* @param threadRunHelper 扫描器
*/
@Override
public void addThread(BaseThreadRunHelper threadRunHelper) {
String lockName = threadRunHelper.getLockName();
THREAD_RUN_OPT.put(lockName,threadRunHelper);
LOCK_NAMES.add(lockName);
}
/**
* 销毁器,将线程管理池里面的扫描器全部注销掉
*/
@Override
public void logoutThreadGroup() {
THREAD_GROUP_STOP = true;
Set<Map.Entry<String, Thread>> threadExamples = THREADS_MAP.entrySet();
threadExamples.forEach(threadExample ->{
log.warn("------------开始注销线程{}------------",threadExample);
String threadName = threadExample.getKey();
Thread thread = threadExample.getValue();
dateAligned(CLOSE_WAIT_TIME,threadName);
//判断线程是否处于终止状态
if(thread.getState() != Thread.State.TERMINATED){
thread.interrupt();
try {
thread.join();
} catch (InterruptedException e) {
e.printStackTrace( );
}
}
log.warn("--------{}线程被注销---------",threadName);
});
LOCK_NAMES.forEach(lockName ->{
log.warn("--------{}锁消除---------",lockName);
RedissLockUtil.unlock(lockName);
});
}
/**
* 对齐时钟。整秒运行
* @param waitTime 休眠时间
* @param threadName 线程名称
*/
public void dateAligned(long waitTime,String threadName){
try {
log.debug("-----------线程{}开始休眠,休眠时间{}---------",threadName,waitTime);
TimeUnit.MILLISECONDS.sleep(waitTime - System.currentTimeMillis() % 1000);
} catch (InterruptedException e) {
log.warn("----------------【{}线程被中断】-----------------------",threadName);
}
}
}
package com.byit.factory;
import com.byit.thread.BaseThreadRunHelper;
/**
* 线程构建帮助器
* @author huangfu
*/
public interface DaemonScanThreadRunHelper {
/**
* 添加线程
* @param threadRunHelper
*/
void addThread(BaseThreadRunHelper threadRunHelper);
/**
* 运行线程
*/
void runDaemonThreads();
/**
* 注销线程组
*/
void logoutThreadGroup();
}
package com.byit.factory;
import com.byit.thread.BaseThreadRunHelper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.context.annotation.Primary;
import org.springframework.stereotype.Component;
import javax.sql.DataSource;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
/**
* 扫描线程运行帮助 使用数据库锁
* @author huangfu
*/
@Slf4j
@Component
@ConditionalOnExpression("'${lock.type}'.equals('db')")
public class DaemonScanThreadRunHelperDbLock extends BaseDaemonScanThreadRunHelper {
/**
* 运行线程
*/
@Override
public void runDaemonThreads(){
Set<Map.Entry<String, BaseThreadRunHelper>> entries = THREAD_RUN_OPT.entrySet();
entries.forEach(threadExamples ->{
String examplesKey = threadExamples.getKey();
BaseThreadRunHelper examplesValue = threadExamples.getValue();
buildThread(examplesKey,examplesValue);
});
}
/**
* 线程构建
* @param lockName 行锁名称
* @param examplesValue 线程运行资源
*/
private void buildThread(String lockName, BaseThreadRunHelper examplesValue){
log.info("-----------开始构建扫描线程,线程锁为{}----------------",lockName);
DataSource dataSource = examplesValue.getDataSource();
if(dataSource != null){
String threadClassName = examplesValue.getClass().getSimpleName();
String threadName = "myth-job#【"+threadClassName+"】";
Thread exampleThread = new Thread(() ->{
//初始化睡眠
dateAligned(INIT_SLEEP_DATE,threadName);
log.info("---------------{}线程启动成功----------------",threadName);
while (!THREAD_GROUP_STOP){
dateAligned(CYCLE_INTERVAL,threadName);
//定义睡眠变量
Long sleepTime = 0L;
Connection conn = null;
Boolean connAutoCommit = null;
PreparedStatement preparedStatement = null;
try {
/**
* 添加行锁
*/
conn = dataSource.getConnection();
connAutoCommit = conn.getAutoCommit();
conn.setAutoCommit(false);
preparedStatement = conn.prepareStatement("SELECT * FROM job_lock WHERE LOCK_NAME = '"+lockName+"' FOR UPDATE ");
preparedStatement.execute();
//调用业务操作
sleepTime = examplesValue.start();
}catch (Exception e){
if(!THREAD_GROUP_STOP){
e.printStackTrace();
}
}finally {
//提交行锁
if (conn != null) {
try {
conn.commit();
} catch (Exception e) {
if (!THREAD_GROUP_STOP) {
log.error("--------------------【提交行锁出错】---------------------");
}
}
}
//恢复自动提交
if (conn != null) {
try {
conn.setAutoCommit(connAutoCommit);
} catch (Exception e) {
if (!THREAD_GROUP_STOP) {
log.error("--------------------【恢复自动提交出错】---------------------");
}
}
}
//关闭数据库执行器
if (preparedStatement != null) {
try {
preparedStatement.close();
} catch (SQLException e) {
if (!THREAD_GROUP_STOP) {
log.error("--------------------【关闭执行器出错】---------------------");
}
}
}
//关闭数据库连接
if (conn != null) {
try {
conn.close();
} catch (SQLException e) {
if (!THREAD_GROUP_STOP) {
log.error("--------------------【关闭执行器出错】---------------------");
}
}
}
}
sleepTime = sleepTime==null?examplesValue.UNIVERSAL_WAIT_TIME:sleepTime;
dateAligned(sleepTime,threadName);
}
});
exampleThread.setName(threadName);
exampleThread.setDaemon(true);
exampleThread.start();
THREADS_MAP.put(threadName,exampleThread);
}else{
System.out.println("--------------出错了-----------");
}
}
}
package com.byit.factory;
import com.byit.thread.BaseThreadRunHelper;
import com.byit.util.lock.RedissLockUtil;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Component;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.TimeUnit;
/**
* 使用redis分布式锁
* @author huangfu
*/
@Component
@Slf4j
@ConditionalOnExpression("'${lock.type}'.equals('redis')")
public class DaemonScanThreadRunHelperRedisLock extends BaseDaemonScanThreadRunHelper {
@Override
public void runDaemonThreads() {
Set<Map.Entry<String, BaseThreadRunHelper>> entries = THREAD_RUN_OPT.entrySet();
entries.forEach(threadExamples ->{
String examplesKey = threadExamples.getKey();
BaseThreadRunHelper examplesValue = threadExamples.getValue();
buildThread(examplesKey,examplesValue);
});
}
private void buildThread(String lockName, BaseThreadRunHelper examplesValue){
log.info("-----------开始构建扫描线程,线程锁为{}----------------",lockName);
String threadClassName = examplesValue.getClass().getSimpleName();
String threadName = "myth-job#【"+threadClassName+"】";
Thread exampleThread = new Thread(() ->{
//初始化睡眠
dateAligned(INIT_SLEEP_DATE,threadName);
log.info("---------------{}线程启动成功----------------",threadName);
while (!THREAD_GROUP_STOP){
//定义睡眠变量
Long sleepTime = 0L;
try{
//加锁 60秒后超时
if (RedissLockUtil.trlock(lockName, TimeUnit.SECONDS, 59)) {
log.debug("------------{},加锁成功,锁名称为{}----------",threadName,lockName);
dateAligned(CYCLE_INTERVAL,threadName);
//调用业务操作
sleepTime = examplesValue.start();
RedissLockUtil.unlock(lockName);
}
}catch (Exception e) {
RedissLockUtil.unlock(lockName);
if(!THREAD_GROUP_STOP){
log.error("---------{}---------",e.getMessage());
}
}
sleepTime = sleepTime==null || sleepTime<=0?examplesValue.UNIVERSAL_WAIT_TIME:sleepTime;
dateAligned(sleepTime,threadName);
}
});
exampleThread.setName(threadName);
exampleThread.setDaemon(true);
exampleThread.start();
THREADS_MAP.put(threadName,exampleThread);
}
}
package com.byit.flowservice;
import com.byit.model.JobTask;
/**
* 做节点验证的数据
* 具体注释看下方方法级的注释
* @author huangfu
*/
public interface NodeVerification {
/**
* 上级节点的状态
* 这个方法主要就是检测当前节点的上级节点的状态
* 当上级节点处于未执行 执行中 失败重试中 时返回的状态为 false 即当前节点不可运行
* 当上级节点全部执行成功时 返回false
* 当上级节点失败重试完毕切结果依旧是失败的情况下 会抛出一个异常信息,即上级节点执行失败
* @param thisJobTask 当前节点
* @return 成功 失败 运行中 未重试完毕
*/
boolean superiorNodeStatus(JobTask thisJobTask);
}
package com.byit.flowservice.impl;
import cn.hutool.core.collection.CollectionUtil;
import com.byit.enums.NodePropertyEnum;
import com.byit.enums.NodeRunStatusPropertyEnum;
import com.byit.enums.ScheduleTypeEnum;
import com.byit.exceptions.SuperiorNodeRunException;
import com.byit.flowservice.NodeVerification;
import com.byit.job.exceptions.BusinessException;
import com.byit.model.JobTask;
import com.byit.model.JobTaskRunLog;
import com.byit.service.JobTaskRunLogService;
import com.byit.util.MythStringUtil;
import org.apache.commons.lang3.StringUtils;
import org.springframework.stereotype.Component;
import java.util.List;
import java.util.stream.Collectors;
/**
* 具体的详细注释参见{@link NodeVerification}
* @author huangfu
*/
@Component
public class NodeVerificationImpl implements NodeVerification {
public static final String NODEDEPEND_SPLIT = ",";
/**
* 日志节点操作
*/
private final JobTaskRunLogService jobTaskRunLogService;
public NodeVerificationImpl(JobTaskRunLogService jobTaskRunLogService) {
this.jobTaskRunLogService = jobTaskRunLogService;
}
/**
* 具体这个方法的概述参见{@link NodeVerification#superiorNodeStatus(com.byit.model.JobTask)}
*
* 这个方法判断上级是否执行完毕的逻辑
* 1.查询该节点的依赖节点
* 2.然后返回所有依赖节点的日志信息,判断上级节点是否和查询出来的数目相同,相同就证明上级节点已经全部完成了
* 3.判断上级节点是否全部成功
* 成功:true
* 失败:
* 判断失败的节点是否已经重试完毕
* 完毕:
* 判断该节点是否是弱引用 弱引用直接返回true
* 抛出异常
* 重试中:false
* @param thisJobTask 当前节点
* @return 返回的是该节点是否可以执行
*/
@Override
public boolean superiorNodeStatus(JobTask thisJobTask) {
//获取该节点的运行标识
//TODO 这里有个问题,就是第一个节点是寻找他的重跑主工作流里面依赖的节点 但是除了第一个 都是找本节点的功能项
String runId = thisJobTask.getRunId();
// if(ScheduleTypeEnum.REPEAT.getCode().equals(thisJobTask.getScheduleType())) {
// runId = thisJobTask.getReRunId();
// }
//查询该节点的依赖节点
String nodeDepend = thisJobTask.getNodeDepend();
if(StringUtils.isBlank(nodeDepend)){
//对于这个操作,外部捕获到这个异常后应该将该节点写入日志,并设置异常信息
throw new BusinessException(NodeRunStatusPropertyEnum.NODE_RELY_ERROR.getMsg());
}
String[] split = nodeDepend.split(NODEDEPEND_SPLIT);
//根据依赖节点查询对应的日志信息
List<Integer> nodeDependIntegers = MythStringUtil.stringArrayConvertIntegerArray(split);
//这里返回的是上级节点的日志执行情况 把运行中的数据给过滤掉了
List<JobTaskRunLog> jobTaskRunLogList = jobTaskRunLogService.findJobTaskRunLogNotEndNodeByRunCodeCount(nodeDependIntegers, runId);
//判断上级节点是否和查询出来的数目相同
if (nodeDependIntegers.size() == jobTaskRunLogList.size()){
//如果腹肌节点全部完成 那么判断父级节点是否全部成功
//过滤失败的节点
List<JobTaskRunLog> errorJobLog = jobTaskRunLogList.stream().
filter(jobTaskRunLog -> (
NodeRunStatusPropertyEnum.RUN_FAILURE.getCode().equals(jobTaskRunLog.getRunCode())
|| NodeRunStatusPropertyEnum.RE_RUN_FAILURE.getCode().equals(jobTaskRunLog.getRunCode())
|| NodeRunStatusPropertyEnum.PARENT_NODE_FAILED.getCode().equals(jobTaskRunLog.getRunCode())))
.collect(Collectors.toList());
//当上级节点的失败个数为0时 返回true
if(CollectionUtil.isEmpty(errorJobLog)){
return true;
}else {
//查看失败节点的重试次数是不是为0
for (JobTaskRunLog jobTaskRunLog : errorJobLog) {
if (jobTaskRunLog.getFailedRemainingCount() != null) {
if (jobTaskRunLog.getFailedRemainingCount() > 0){
//上级节点的失败原因还不能是被杀死和上级节点执行失败的,只有这样他才有重试的资格
if (!(NodeRunStatusPropertyEnum.KILL.getCode().equals(jobTaskRunLog.getRunCode()) ||
NodeRunStatusPropertyEnum.PARENT_NODE_FAILED.getCode().equals(jobTaskRunLog.getRunCode()))){
return false;
}
}
}
}
//这里还有一层判断,就是当上级节点执行失败了,而且重试完了,不能立即断定该节点就一定要执行快速失败,因为如果该节点是弱引用
//那么该节点依旧能够执行,这一段逻辑之所以不在首行进行判断是因为,无论他是不是弱引用,他都要等待上级节点执行完毕后才能执行
//判断该节点是否是弱引用
if (NodePropertyEnum.WEAK_NODE.getCode().equals(thisJobTask.getSuperSuccessRun())) {
return true;
}
throw new SuperiorNodeRunException("上级节点执行失败");
}
}
return false;
}
}
package com.byit.job;
import com.byit.task.JavaTaskJobTask;
import io.netty.util.HashedWheelTimer;
import io.netty.util.TimerTask;
import lombok.extern.slf4j.Slf4j;
import java.util.concurrent.TimeUnit;
/**
* @program: byit-myth-job->WorkRoulette
* @description: 工作轮盘,所有的添加任务 删除任务 暂停任务都在此列
* @author: huangfu
* @date: 2019/11/15 14:59
**/
@Slf4j
public class WorkRoulette {
/**
* HASHED_WHEEL_TIMER:工作轮盘
* ThreadFactory:创建work线程
* tickDuration:每个刻度的时间
* ticsPerWheel:轮盘一圈大小
* maxPendingTimeouts:最大等待处理超时
*/
private static final HashedWheelTimer HASHED_WHEEL_TIMER = new HashedWheelTimer(r -> new Thread(r, "HASHED_WHEEL_TIMER" + r.hashCode()), 1, TimeUnit.SECONDS, 8, true, 0);
public static final int INT = 1000;
public static void addJob(TimerTask timerTask,long triggerNextTime) {
log.info("-----任务{}毫秒后执行-------",triggerNextTime-System.currentTimeMillis());
//设置这个的根部原因是因为保证任务的抛出在事务提交动作完成之后抛出
long time = triggerNextTime-System.currentTimeMillis();
if(time < INT){
time = INT;
}
HASHED_WHEEL_TIMER.newTimeout(timerTask, TimeUnit.MILLISECONDS.toNanos(time), TimeUnit.NANOSECONDS);
}
}
package com.byit.listener;
import com.byit.enums.FlowPropertyEnum;
import com.byit.event.EndFlowEvent;
import com.byit.event.FlowScanEndEvent;
import com.byit.model.Flow;
import com.byit.service.FlowService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Component;
import java.util.Date;
/**
* 工作流的事件监听操作
* @author huangfu
*/
@Component
@Slf4j
public class FlowEventListener {
private final FlowService flowService;
public FlowEventListener(FlowService flowService) {
this.flowService = flowService;
}
/**
* 工作添加进实例表后的事件监听
* @param flowScanEndEvent 事假信息
*/
@EventListener
public void flowScanEndEventListener(FlowScanEndEvent flowScanEndEvent){
log.info("-----------监听到事件{},工作流添加进实例完成事件-------",flowScanEndEvent);
Flow flow = Flow.builder()
.flowId(flowScanEndEvent.getFlowId())
.scanMark(FlowPropertyEnum.NOT_SCAN.getCode())
.build();
flowService.updateByIdSelective(flow);
}
/**
* 工作流完成事件监听
* @param endFlowEvent 事件信息
*/
@EventListener
public void flowEndEventListener(EndFlowEvent endFlowEvent){
log.info("-----------监听到事件{},工作流实例完成事件-------",endFlowEvent);
Flow flow = Flow.builder()
.flowId(endFlowEvent.getFlowId())
.scanMark(FlowPropertyEnum.SCAN.getCode())
.build();
flowService.updateByIdSelective(flow);
}
}
package com.byit.mapper;
import com.byit.model.EmailAlarm;
import com.byit.model.vo.EmailAlarmVo;
import org.apache.ibatis.annotations.Param;
import org.springframework.stereotype.Repository;
import java.util.List;
/**
* @author huangfu
*/
@Repository
public interface EmailAlarmMapper {
/**
* 查询没有告警的邮箱
* @return
*/
List<EmailAlarm> findEmailAlarmByAlarmResult();
/**
* 根据id查询邮件
* @param id
* @return
*/
EmailAlarm findEmailAlarmById(Integer id);
/**
* 删除一个邮件
* @param id
* @return
*/
int deleteById(Integer id);
/**
* 保存邮件
* @param record
* @return
*/
int saveEmailAlarm(EmailAlarm record);
/**
* 批量保存邮件
* @param emailAlarms
* @return
*/
int saveEmailAlarms(@Param("emailAlarms") List<EmailAlarm> emailAlarms);
/**
* 修改邮件发送情况
* @param record
* @return
*/
int updateEmailAlarm(EmailAlarm record);
}
\ No newline at end of file
package com.byit.mapper;
import com.byit.dto.FlowConditionDto;
import com.byit.dto.StatisticsConditionDto;
import com.byit.model.Flow;
import org.apache.ibatis.annotations.Param;
import java.util.List;
public interface FlowMapper {
/**
* 查询当天将要运行的工作流
* @param statisticsConditionDto
* @return 所有的流信息
*/
List<Flow> findAllByThisDayFlow(StatisticsConditionDto statisticsConditionDto);
/**
* 查询全部的任务流数据
* @return
*/
List<Flow> findAllFlow();
/**
* 根据ID查询
* @param id
* @return
*/
Flow findFlowById(Integer id);
/**
* 查询半个小时内即将要执行的工作流
* @param triggerNextTime
* @return
*/
List<Flow> findHalfAnHourFlow(@Param("triggerNextTime") Long triggerNextTime);
int deleteById(Integer flowId);
int insertSelective(Flow record);
Flow getById(Integer flowId);
int updateByIdSelective(Flow record);
/**
* 根据工作空间和工作流名称查找是否存在工作流
* @param workspaceId
* @param flowName
* @return
*/
Flow getByWorkSpaceAndName(@Param("workspaceId") Integer workspaceId, @Param("flowName") String flowName);
/**
* 根据workspace查找所属的工作流
* @param workspaceId
* @return
*/
List<Flow> findByWorkspace(Integer workspaceId);
/**
* 获取已经启动的工作流
* @return
*/
List<Flow> findONStartUp();
/**
* 获取全部的工作流
* @return
*/
List<Flow> findAll();
List<Flow> findAllByCondition(FlowConditionDto flowConditionDto);
}
\ No newline at end of file
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