Skip to content
Projects
Groups
Snippets
Help
This project
Loading...
Sign in / Register
Toggle navigation
B
byit-myth-job
Overview
Overview
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
liyuan
byit-myth-job
Commits
6e3322fc
Commit
6e3322fc
authored
Dec 31, 2019
by
guominglei
Browse files
Options
Browse Files
Download
Plain Diff
Merge remote-tracking branch 'origin/developer' into developer
parents
0f88f263
b489a913
Show whitespace changes
Inline
Side-by-side
Showing
20 changed files
with
455 additions
and
65 deletions
+455
-65
EmailController.java
...in/src/main/java/com/byit/controller/EmailController.java
+38
-0
application.yml
byit-myth-admin/src/main/resources/application.yml
+9
-4
logback-spring.xml
byit-myth-admin/src/main/resources/logback-spring.xml
+3
-3
pom.xml
byit-myth-core/myth-admin-core/pom.xml
+0
-11
AdminEnums.java
...h-admin-core/src/main/java/com/byit/enums/AdminEnums.java
+2
-1
EmailAlarmMapper.java
...-core/src/main/java/com/byit/mapper/EmailAlarmMapper.java
+31
-4
JobTaskRunLog.java
...dmin-core/src/main/java/com/byit/model/JobTaskRunLog.java
+7
-1
RunRecording.java
...admin-core/src/main/java/com/byit/model/RunRecording.java
+6
-0
EmailAlarmVo.java
...in-core/src/main/java/com/byit/model/vo/EmailAlarmVo.java
+108
-0
EmailAlarmService.java
...ore/src/main/java/com/byit/service/EmailAlarmService.java
+45
-0
EmailAlarmServiceImpl.java
...ain/java/com/byit/service/impl/EmailAlarmServiceImpl.java
+103
-0
JobScheduleHelper.java
...core/src/main/java/com/byit/thread/JobScheduleHelper.java
+4
-2
LogScanHelper.java
...min-core/src/main/java/com/byit/thread/LogScanHelper.java
+46
-17
EmailAlarmMapper.xml
...admin-core/src/main/resources/mapper/EmailAlarmMapper.xml
+3
-3
JobTaskRunLogMapper.xml
...in-core/src/main/resources/mapper/JobTaskRunLogMapper.xml
+16
-3
RunRecordingMapper.xml
...min-core/src/main/resources/mapper/RunRecordingMapper.xml
+14
-2
JobResultEnum.java
...ommon/src/main/java/com/byit/job/enums/JobResultEnum.java
+12
-6
ReturnResult.java
...re-common/src/main/java/com/byit/job/vo/ReturnResult.java
+5
-5
RunJobThread.java
...ector-plugin/src/main/java/com/byit/rpc/RunJobThread.java
+1
-1
Mains.java
...nt/byit-demo-client/src/main/java/com/byit/job/Mains.java
+2
-2
No files found.
byit-myth-admin/src/main/java/com/byit/controller/EmailController.java
0 → 100644
View file @
6e3322fc
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"
;
}
}
byit-myth-admin/src/main/resources/application.yml
View file @
6e3322fc
...
...
@@ -4,10 +4,15 @@ spring:
url
:
jdbc:mysql://10.0.10.118:3306/myth-job?Unicode=true&characterEncoding=UTF-8&useSSL=true
username
:
root
password
:
123456
# jpa:
# hibernate:
# ddl-auto: none
# show-sql: false
mail
:
host
:
smtp.163.com
#我自己的SMTP服务器地址
username
:
huangfusuper@163.com
#登录用户名
password
:
huangfu0110
#授权密码
default-encoding
:
UTF-8
protocol
:
smtp
#协议
properties
:
from
:
huangfusuper@163.com
#真实邮箱
mybatis
:
mapper-locations
:
/mapper/*.xml
...
...
byit-myth-admin/src/main/resources/logback-spring.xml
View file @
6e3322fc
...
...
@@ -2,9 +2,6 @@
<configuration>
<appender
name=
"console"
class=
"ch.qos.logback.core.ConsoleAppender"
>
<!-- <filter class="ch.qos.logback.classic.filter.ThresholdFilter">
<level>ERROR</level>
</filter>-->
<encoder>
<pattern>
%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{50} - %msg%n
</pattern>
</encoder>
...
...
@@ -64,6 +61,9 @@
<logger
name=
"com.byit.thread"
level=
"debug"
additivity=
"false"
>
<appender-ref
ref=
"console"
/>
</logger>
<logger
name=
"com.byit.selector"
level=
"debug"
additivity=
"false"
>
<appender-ref
ref=
"console"
/>
</logger>
<logger
name=
"org.springframework.web"
level=
"debug"
/>
...
...
byit-myth-core/myth-admin-core/pom.xml
View file @
6e3322fc
...
...
@@ -45,17 +45,6 @@
<artifactId>
mysql-connector-java
</artifactId>
</dependency>
<!--自动装载配置-->
<dependency>
<groupId>
org.springframework.boot
</groupId>
<artifactId>
spring-boot-autoconfigure
</artifactId>
</dependency>
<dependency>
<groupId>
org.springframework.boot
</groupId>
<artifactId>
spring-boot-configuration-processor
</artifactId>
</dependency>
<dependency>
<groupId>
org.springframework.boot
</groupId>
<artifactId>
spring-boot-starter-jdbc
</artifactId>
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/enums/AdminEnums.java
View file @
6e3322fc
...
...
@@ -11,7 +11,8 @@ import lombok.AllArgsConstructor;
**/
@AllArgsConstructor
public
enum
AdminEnums
implements
IEnum
{
TEST
(
"2132131"
,
"滚犊子"
)
SEND_EMAIL_FAILURE
(
"100500"
,
"邮件发送失败"
),
SEND_EMAIL_SUCCESS
(
"100200"
,
"邮件发送成功"
)
;
String
code
;
String
msg
;
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/mapper/EmailAlarmMapper.java
View file @
6e3322fc
package
com
.
byit
.
mapper
;
import
com.byit.model.EmailAlarm
;
import
com.byit.model.vo.EmailAlarmVo
;
import
org.springframework.stereotype.Repository
;
/**
* @author huangfu
*/
@Repository
public
interface
EmailAlarmMapper
{
int
deleteById
(
Integer
id
);
/**
* 根据id查询邮件
* @param id
* @return
*/
EmailAlarm
findEmailAlarmById
(
Integer
id
);
int
insertSelective
(
EmailAlarm
record
);
/**
* 删除一个邮件
* @param id
* @return
*/
int
deleteById
(
Integer
id
);
EmailAlarm
getById
(
Integer
id
);
/**
* 保存邮件
* @param record
* @return
*/
int
saveEmailAlarm
(
EmailAlarm
record
);
int
updateByIdSelective
(
EmailAlarm
record
);
/**
* 修改邮件发送情况
* @param record
* @return
*/
int
updateEmailAlarm
(
EmailAlarm
record
);
}
\ No newline at end of file
byit-myth-core/myth-admin-core/src/main/java/com/byit/model/JobTaskRunLog.java
View file @
6e3322fc
...
...
@@ -76,7 +76,7 @@ public class JobTaskRunLog implements Serializable {
/**
* 运行结果
*/
@ApiModelProperty
(
"
运行结果
"
)
@ApiModelProperty
(
"
1 成功 2 失败 3 补批成功 4 补批失败
"
)
private
String
runCode
;
/**
...
...
@@ -154,6 +154,11 @@ public class JobTaskRunLog implements Serializable {
@ApiModelProperty
(
"运行标识"
)
private
String
runId
;
/**
* 补批运行标识
*/
@ApiModelProperty
(
"补批运行标识"
)
private
String
reRunId
;
/**
*/
private
static
final
long
serialVersionUID
=
1L
;
}
\ No newline at end of file
byit-myth-core/myth-admin-core/src/main/java/com/byit/model/RunRecording.java
View file @
6e3322fc
...
...
@@ -126,6 +126,12 @@ public class RunRecording implements Serializable {
*/
@ApiModelProperty
(
"是否是内嵌工作流 0 否, 1 是"
)
private
String
isInner
;
/**
* 是否执行快速失败 0 否, 1 是
*/
@ApiModelProperty
(
"是否执行快速失败 0 否, 1 是"
)
private
String
failFast
;
/**
*/
private
static
final
long
serialVersionUID
=
1L
;
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/model/vo/EmailAlarmVo.java
0 → 100644
View file @
6e3322fc
package
com
.
byit
.
model
.
vo
;
import
com.byit.annotation.annotationselector.EmailValidation
;
import
com.byit.annotation.annotationselector.NotBank
;
import
com.byit.annotation.annotationselector.NotNull
;
import
com.fasterxml.jackson.annotation.JsonIgnore
;
import
io.swagger.annotations.ApiModel
;
import
io.swagger.annotations.ApiModelProperty
;
import
java.io.File
;
import
java.io.Serializable
;
import
java.util.Date
;
import
java.util.List
;
import
lombok.AllArgsConstructor
;
import
lombok.Builder
;
import
lombok.Data
;
import
lombok.NoArgsConstructor
;
/**
* 邮箱告警
* @author gen
*/
@ApiModel
@Data
@Builder
@AllArgsConstructor
@NoArgsConstructor
public
class
EmailAlarmVo
implements
Serializable
{
/**
* 主键
*/
@ApiModelProperty
(
"主键"
)
private
Integer
id
;
/**
* 工作流Id
*/
@ApiModelProperty
(
"工作流Id"
)
@NotNull
(
errorMessage
=
"工作流id不能为空"
)
private
Integer
flowId
;
/**
* 工作流名称
*/
@ApiModelProperty
(
"工作流名称"
)
@NotBank
(
errorMessage
=
"工作流名称不能为空"
)
private
String
flowName
;
/**
* 运行标识
*/
@ApiModelProperty
(
"运行标识"
)
@NotBank
(
errorMessage
=
"运行标识不能为空"
)
private
String
runId
;
/**
* 版本名称
*/
@ApiModelProperty
(
"版本名称"
)
@NotBank
(
errorMessage
=
"版本名称不能为空"
)
private
String
versionName
;
/**
* 告警邮箱
*/
@ApiModelProperty
(
"告警邮箱"
)
@NotBank
(
errorMessage
=
"告警邮箱不能为空,多个联系人可用,分割"
)
@EmailValidation
private
String
alarmEmail
;
/**
* 当前工作流版本的告警的时机(0 不告警, 1 完成时告警, 2 失败时告警, 3 成功时告警)
*/
@ApiModelProperty
(
"当前工作流版本的告警的时机(0 不告警, 1 完成时告警, 2 失败时告警, 3 成功时告警)"
)
private
String
alarmAction
;
/**
* 告警结果(0, 未告警 1,告警成功 2,告警失败)
*/
@ApiModelProperty
(
"告警结果(0, 未告警 1,告警成功 2,告警失败)"
)
private
String
alarmResult
;
/**
* 发送时间
*/
@ApiModelProperty
(
"发送时间"
)
private
Date
sendTime
;
/**
* 告警内容
*/
@ApiModelProperty
(
"告警内容"
)
private
String
alarmContent
;
@ApiModelProperty
(
"告警主题"
)
private
String
alarmTitle
;
@ApiModelProperty
(
"附件"
)
@JsonIgnore
private
List
<
File
>
fileList
;
/**
*/
private
static
final
long
serialVersionUID
=
1L
;
}
\ No newline at end of file
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/EmailAlarmService.java
0 → 100644
View file @
6e3322fc
package
com
.
byit
.
service
;
import
com.byit.model.EmailAlarm
;
import
com.byit.model.vo.EmailAlarmVo
;
/**
* @program: byit-myth-job->EmailAlarmService
* @description: 邮件发送的业务接口
* @author: huangfu
* @date: 2019/12/30 10:59
**/
public
interface
EmailAlarmService
{
/**
* 根据id查询邮件
* @param id
* @return
*/
EmailAlarm
findEmailAlarmById
(
Integer
id
);
/**
* 删除一个邮件
* @param id
* @return
*/
int
deleteById
(
Integer
id
);
/**
* 保存邮件
* @param record
* @return
*/
int
saveEmailAlarm
(
EmailAlarm
record
);
/**
* 发送邮件
* @param record
*/
void
sendEmail
(
EmailAlarmVo
record
);
/**
* 修改邮件发送情况
* @param record
* @return
*/
int
updateEmailAlarm
(
EmailAlarm
record
);
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/impl/EmailAlarmServiceImpl.java
0 → 100644
View file @
6e3322fc
package
com
.
byit
.
service
.
impl
;
import
cn.hutool.core.collection.CollectionUtil
;
import
com.byit.annotation.DataValidation
;
import
com.byit.annotation.ParamValidation
;
import
com.byit.enums.AdminEnums
;
import
com.byit.job.exceptions.BusinessException
;
import
com.byit.mapper.EmailAlarmMapper
;
import
com.byit.model.EmailAlarm
;
import
com.byit.model.vo.EmailAlarmVo
;
import
com.byit.service.EmailAlarmService
;
import
lombok.extern.slf4j.Slf4j
;
import
org.springframework.beans.BeanUtils
;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.mail.javamail.JavaMailSenderImpl
;
import
org.springframework.mail.javamail.MimeMessageHelper
;
import
org.springframework.stereotype.Service
;
import
javax.mail.MessagingException
;
import
java.util.Date
;
/**
* @program: byit-myth-job->EmailAlarmServiceImpl
* @description: 邮件发送的业务接口
* @author: huangfu
* @date: 2019/12/30 11:00
**/
@Service
@Slf4j
public
class
EmailAlarmServiceImpl
implements
EmailAlarmService
{
private
final
EmailAlarmMapper
emailAlarmMapper
;
private
final
JavaMailSenderImpl
javaMailSender
;
@Autowired
public
EmailAlarmServiceImpl
(
EmailAlarmMapper
emailAlarmMapper
,
JavaMailSenderImpl
javaMailSender
)
{
this
.
emailAlarmMapper
=
emailAlarmMapper
;
this
.
javaMailSender
=
javaMailSender
;
}
@Override
public
EmailAlarm
findEmailAlarmById
(
Integer
id
)
{
return
emailAlarmMapper
.
findEmailAlarmById
(
id
);
}
@Override
public
int
deleteById
(
Integer
id
)
{
return
emailAlarmMapper
.
deleteById
(
id
);
}
@Override
public
int
saveEmailAlarm
(
EmailAlarm
record
)
{
return
emailAlarmMapper
.
saveEmailAlarm
(
record
);
}
@Override
@DataValidation
public
void
sendEmail
(
@ParamValidation
EmailAlarmVo
record
)
{
record
.
setSendTime
(
new
Date
());
record
.
setAlarmTitle
(
"测试邮件"
);
record
.
setAlarmContent
(
"<div style='color:#F00;font-size:1000px'>测试邮件</div>"
);
try
{
//发送一个支持复杂类型的邮件
MimeMessageHelper
messageHelper
=
new
MimeMessageHelper
(
javaMailSender
.
createMimeMessage
(),
true
);
//邮件发信人
messageHelper
.
setFrom
(
javaMailSender
.
getJavaMailProperties
().
getProperty
(
"from"
));
//邮件收信人
messageHelper
.
setTo
(
record
.
getAlarmEmail
().
split
(
","
));
//邮件主题
messageHelper
.
setSubject
(
record
.
getAlarmTitle
());
//邮件内容
messageHelper
.
setText
(
record
.
getAlarmContent
(),
true
);
if
(
CollectionUtil
.
isNotEmpty
(
record
.
getFileList
())){
record
.
getFileList
().
forEach
(
file
->{
try
{
messageHelper
.
addAttachment
(
file
.
getName
(),
file
);
}
catch
(
MessagingException
e
)
{
record
.
setAlarmResult
(
"2"
);
log
.
error
(
"-----------添加附件失败【{}】--------------"
,
e
.
getMessage
());
}
});
}
//正式发送邮件
javaMailSender
.
send
(
messageHelper
.
getMimeMessage
());
record
.
setAlarmResult
(
"1"
);
}
catch
(
Exception
e
){
record
.
setAlarmResult
(
"2"
);
log
.
error
(
"-----------邮件发送失败【{}】--------------"
,
e
.
getMessage
());
throw
new
BusinessException
(
AdminEnums
.
SEND_EMAIL_FAILURE
);
}
EmailAlarm
emailAlarm
=
new
EmailAlarm
(
);
BeanUtils
.
copyProperties
(
record
,
emailAlarm
);
//保存邮件
saveEmailAlarm
(
emailAlarm
);
}
@Override
public
int
updateEmailAlarm
(
EmailAlarm
record
)
{
return
emailAlarmMapper
.
updateEmailAlarm
(
record
);
}
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/thread/JobScheduleHelper.java
View file @
6e3322fc
...
...
@@ -329,8 +329,10 @@ public class JobScheduleHelper{
private
Integer
saveLog
(
JobTaskSchedule
mythJobTaskSchedule
){
JobTaskRunLogWithBLOBs
jobTaskRunLog
=
new
JobTaskRunLogWithBLOBs
();
jobTaskRunLog
.
setRunId
(
mythJobTaskSchedule
.
getRunId
());
jobTaskRunLog
.
setFlowName
(
mythJobTaskSchedule
.
getVersionName
());
//jobTaskRunLog.setRunId(mythJobTaskSchedule.getRunId());
jobTaskRunLog
.
setFlowId
(
mythJobTaskSchedule
.
getFlowId
());
jobTaskRunLog
.
setFlowName
(
"1"
);
jobTaskRunLog
.
setRunId
(
"qwer-tyui-opas-dfgh"
);
jobTaskRunLog
.
setNodeName
(
mythJobTaskSchedule
.
getNodeName
());
jobTaskRunLog
.
setRunParams
(
mythJobTaskSchedule
.
getRunParam
());
jobTaskRunLog
.
setFailedRemainingCount
(
mythJobTaskSchedule
.
getFailedRetryCount
());
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/thread/LogScanHelper.java
View file @
6e3322fc
...
...
@@ -42,8 +42,10 @@ public class LogScanHelper {
}
private
final
String
FAILURE_CODE
=
"100500"
;
private
final
String
SUCCESS_CODE
=
"100200"
;
private
final
String
SUCCESS_CODE
=
"1"
;
private
final
String
FAILURE_CODE
=
"2"
;
private
final
String
RE_SUCCESS_CODE
=
"3"
;
private
final
String
RE_FAILURE_CODE
=
"4"
;
/**
* 扫描虚节点的线程是否停止
*/
...
...
@@ -89,7 +91,7 @@ public class LogScanHelper {
notEndVirtualNodes
.
forEach
(
notEndVirtualNode
->{
log
.
debug
(
"------------【开始查询虚拟节点的执行情况】----------------"
);
//根据运行标识和工作流id查询运行日志
RunRecording
runRecordingByFlowIdAndRunId
=
runRecordingService
.
findRunRecordingByFlowIdAndRunId
(
notEndVirtualNode
.
getFlowId
(
),
notEndVirtualNode
.
getRunId
(
));
RunRecording
runRecordingByFlowIdAndRunId
=
runRecordingService
.
findRunRecordingByFlowIdAndRunId
(
notEndVirtualNode
.
getFlowId
(
),
notEndVirtualNode
.
getRunId
());
//判断当前的工作流是否已经完结
if
(
runRecordingByFlowIdAndRunId
!=
null
&&
"4"
.
equals
(
runRecordingByFlowIdAndRunId
.
getFlowStatus
())){
log
.
debug
(
"------------【查询到有已经完成的虚拟节点修改日志】----------------"
);
...
...
@@ -107,7 +109,7 @@ public class LogScanHelper {
e
.
printStackTrace
();
}
finally
{
//释放资源
freedResource
(
conn
,
preparedStatement
,
connAutoCommit
);
freedResource
(
conn
,
preparedStatement
,
connAutoCommit
,
virtualNodeScanIsStop
);
}
}
...
...
@@ -133,25 +135,47 @@ public class LogScanHelper {
conn
.
setAutoCommit
(
false
);
preparedStatement
=
conn
.
prepareStatement
(
"SELECT * FROM JOB_LOCK WHERE LOCK_NAME = 'log_node_callback_lock' FOR UPDATE "
);
preparedStatement
.
execute
();
System
.
out
.
println
(
"--------------火球锁成功------------------"
);
//查询的是失败的或者是已经结束的节点(完成的)
List
<
JobTaskRunLog
>
jobTaskRunLogEndOrFailureNode
=
jobTaskRunLogService
.
findJobTaskRunLogEndOrFailureNode
(
);
if
(
CollectionUtil
.
isNotEmpty
(
jobTaskRunLogEndOrFailureNode
)){
log
.
debug
(
"-------------【查询到有完结而且未告警的节点】-------------"
);
Map
<
String
,
List
<
JobTaskRunLog
>>
jobLogMap
=
jobTaskRunLogEndOrFailureNode
.
stream
(
).
collect
(
Collectors
.
groupingBy
(
JobTaskRunLog:
:
getRunCode
));
//假设成功的code码是000200 失败的是000500
Map
<
String
,
List
<
JobTaskRunLog
>>
jobLogMap
=
jobTaskRunLogEndOrFailureNode
.
stream
(
)
.
collect
(
Collectors
.
groupingBy
(
JobTaskRunLog:
:
getRunCode
));
//失败的
List
<
JobTaskRunLog
>
failureNodes
=
jobLogMap
.
get
(
FAILURE_CODE
);
List
<
JobTaskRunLog
>
reFailure
=
jobLogMap
.
get
(
RE_FAILURE_CODE
);
if
(
CollectionUtil
.
isNotEmpty
(
reFailure
)){
failureNodes
.
addAll
(
reFailure
);
}
//成功的
List
<
JobTaskRunLog
>
successNodes
=
jobLogMap
.
get
(
SUCCESS_CODE
);
List
<
JobTaskRunLog
>
reSuccessNodes
=
jobLogMap
.
get
(
RE_SUCCESS_CODE
);
if
(
CollectionUtil
.
isNotEmpty
(
reSuccessNodes
)){
successNodes
.
addAll
(
reSuccessNodes
);
}
//遍历成功的节点 修改执行记录表
if
(
CollectionUtil
.
isNotEmpty
(
successNodes
)){
successNodes
.
forEach
(
successNode
->{
String
runId
=
successNode
.
getRunId
(
);
Integer
flowId
=
successNode
.
getFlowId
(
);
runRecordingService
.
updateRunRecordingByFlowIdAndRunId
(
RunRecording
.
builder
().
runId
(
runId
).
flowId
(
flowId
).
flowRunResult
(
"1"
).
flowStatus
(
"4"
).
build
());
successNode
.
setAlertEnd
(
"1"
);
jobTaskRunLogService
.
updateJobTaskRunLog
(
successNode
);
});
}
else
{
dateAligned
(
1000
);
}
/*if(CollectionUtil.isNotEmpty(failureNodes)){
failureNodes.forEach(failureNode ->{
String runId = failureNode.getRunId( );
Integer flowId = failureNode.getFlowId( );
});
}*/
}
else
{
dateAligned
(
20000
);
}
...
...
@@ -159,7 +183,7 @@ public class LogScanHelper {
}
catch
(
Exception
e
){
e
.
printStackTrace
();
}
finally
{
freedResource
(
conn
,
preparedStatement
,
connAutoCommit
);
freedResource
(
conn
,
preparedStatement
,
connAutoCommit
,
notAlarmedNodeScanIsStop
);
}
}
});
...
...
@@ -174,12 +198,12 @@ public class LogScanHelper {
* @param preparedStatement 执行器
* @param connAutoCommit 原来的提交状态
*/
private
void
freedResource
(
Connection
conn
,
PreparedStatement
preparedStatement
,
boolean
connAutoCommit
){
private
void
freedResource
(
Connection
conn
,
PreparedStatement
preparedStatement
,
boolean
connAutoCommit
,
boolean
isStop
){
if
(
conn
!=
null
){
try
{
conn
.
commit
();
}
catch
(
SQLException
e
)
{
if
(!
virtualNodeScanI
sStop
){
if
(!
i
sStop
){
log
.
error
(
"--------------------【提交行锁出错】---------------------"
);
}
}
...
...
@@ -189,28 +213,33 @@ public class LogScanHelper {
assert
conn
!=
null
;
conn
.
setAutoCommit
(
connAutoCommit
);
}
catch
(
SQLException
e
)
{
if
(!
virtualNodeScanI
sStop
){
if
(!
i
sStop
){
log
.
error
(
"--------------------【恢复自动提交出错】---------------------"
);
}
}
if
(
preparedStatement
!=
null
){
try
{
conn
.
close
();
preparedStatement
.
close
();
}
catch
(
SQLException
e
)
{
if
(!
virtualNodeScanIsStop
){
log
.
error
(
"--------------------【关闭链接出错】---------------------"
);
if
(!
isStop
){
log
.
error
(
"--------------------【关闭执行器出错】---------------------"
);
}
}
}
if
(
preparedStatement
!=
null
){
try
{
preparedStatement
.
close
();
conn
.
close
();
}
catch
(
SQLException
e
)
{
if
(!
virtualNodeScanI
sStop
){
log
.
error
(
"--------------------【关闭执行器
出错】---------------------"
);
if
(!
i
sStop
){
log
.
error
(
"--------------------【关闭链接
出错】---------------------"
);
}
}
}
public
void
stopVirtualNodeScanThread
(){
this
.
virtualNodeScanIsStop
=
true
;
}
/**
...
...
byit-myth-core/myth-admin-core/src/main/resources/mapper/EmailAlarmMapper.xml
View file @
6e3322fc
...
...
@@ -26,7 +26,7 @@
<!-- generated @mbg.generated date: 2019-12-25 -->
alarm_content
</sql>
<select
id=
"
get
ById"
parameterType=
"java.lang.Integer"
resultMap=
"ResultMapWithBLOBs"
>
<select
id=
"
findEmailAlarm
ById"
parameterType=
"java.lang.Integer"
resultMap=
"ResultMapWithBLOBs"
>
<!-- generated @mbg.generated date: 2019-12-25 -->
select
<include
refid=
"Base_Column_List"
/>
...
...
@@ -40,7 +40,7 @@
delete from email_alarm
where id = #{id,jdbcType=INTEGER}
</delete>
<insert
id=
"
insertSelective
"
parameterType=
"com.byit.model.EmailAlarm"
>
<insert
id=
"
saveEmailAlarm
"
parameterType=
"com.byit.model.EmailAlarm"
>
<!-- generated @mbg.generated date: 2019-12-25 -->
insert into email_alarm
<trim
prefix=
"("
suffix=
")"
suffixOverrides=
","
>
...
...
@@ -108,7 +108,7 @@
</if>
</trim>
</insert>
<update
id=
"update
ByIdSelective
"
parameterType=
"com.byit.model.EmailAlarm"
>
<update
id=
"update
EmailAlarm
"
parameterType=
"com.byit.model.EmailAlarm"
>
<!-- generated @mbg.generated date: 2019-12-25 -->
update email_alarm
<set>
...
...
byit-myth-core/myth-admin-core/src/main/resources/mapper/JobTaskRunLogMapper.xml
View file @
6e3322fc
...
...
@@ -25,6 +25,7 @@
<result
column=
"alert_end"
jdbcType=
"CHAR"
property=
"alertEnd"
/>
<result
column=
"job_type"
jdbcType=
"CHAR"
property=
"jobType"
/>
<result
column=
"run_id"
jdbcType=
"VARCHAR"
property=
"runId"
/>
<result
column=
"re_run_id"
jdbcType=
"VARCHAR"
property=
"reRunId"
/>
</resultMap>
<resultMap
extends=
"BaseResultMap"
id=
"ResultMapWithBLOBs"
type=
"com.byit.model.JobTaskRunLogWithBLOBs"
>
<result
column=
"run_msg"
jdbcType=
"LONGVARCHAR"
property=
"runMsg"
/>
...
...
@@ -33,7 +34,7 @@
<sql
id=
"Base_Column_List"
>
log_id, failed_remaining_count, version_name, flow_id, flow_name, job_group_id, handler_name,
node_name, is_virtual, run_code, run_params, start_time, run_type, trigger_code,
trigger_time, job_group_ip, map_flow_id, run_command, end_time, node_id,job_type,alert_end,run_id
trigger_time, job_group_ip, map_flow_id, run_command, end_time, node_id,job_type,alert_end,run_id
,re_run_id
</sql>
<sql
id=
"Blob_Column_List"
>
run_msg, trigger_msg
...
...
@@ -42,7 +43,7 @@
<select
id=
"findJobTaskRunLogEndOrFailureNode"
resultMap=
"BaseResultMap"
>
select
<include
refid=
"Base_Column_List"
/>
from job_task_run_log
where (node_name='END' or run_code = '
000500
') and alert_end = '0'
where (node_name='END' or run_code = '
2' or run_code = '4
') and alert_end = '0'
</select>
<select
id=
"findNotEndVirtualNode"
resultMap=
"BaseResultMap"
>
...
...
@@ -142,6 +143,9 @@
<if
test=
"runId != null"
>
run_id,
</if>
<if
test=
"reRunId != null"
>
re_run_id,
</if>
</trim>
<trim
prefix=
"values ("
suffix=
")"
suffixOverrides=
","
>
<if
test=
"logId != null"
>
...
...
@@ -217,7 +221,10 @@
#{triggerMsg,jdbcType=LONGVARCHAR},
</if>
<if
test=
"runId != null"
>
run_id = #{runId,jdbcType=VARCHAR},
#{runId,jdbcType=VARCHAR},
</if>
<if
test=
"reRunId != null"
>
#{reRunId,jdbcType=VARCHAR},
</if>
</trim>
</insert>
...
...
@@ -296,6 +303,9 @@
<if
test=
"runId != null"
>
run_id = #{runId,jdbcType=VARCHAR},
</if>
<if
test=
"reRunId != null"
>
re_run_id = #{reRunId,jdbcType=VARCHAR},
</if>
</set>
where log_id = #{logId,jdbcType=INTEGER}
</update>
...
...
@@ -368,6 +378,9 @@
<if
test=
"runId != null"
>
run_id = #{runId,jdbcType=VARCHAR},
</if>
<if
test=
"reRunId != null"
>
re_run_id = #{reRunId,jdbcType=VARCHAR},
</if>
</set>
where log_id = #{logId,jdbcType=INTEGER}
</update>
...
...
byit-myth-core/myth-admin-core/src/main/resources/mapper/RunRecordingMapper.xml
View file @
6e3322fc
...
...
@@ -20,11 +20,12 @@
<result
column=
"end_time"
jdbcType=
"TIMESTAMP"
property=
"endTime"
/>
<result
column=
"is_alarm"
jdbcType=
"CHAR"
property=
"isAlarm"
/>
<result
column=
"is_inner"
jdbcType=
"CHAR"
property=
"isInner"
/>
<result
column=
"fail_fast"
jdbcType=
"CHAR"
property=
"failFast"
/>
</resultMap>
<sql
id=
"Base_Column_List"
>
recording_id, run_id, alarm_email, dispatch_ip, flow_name, flow_run_result, flow_status,
flow_timeout, flow_version_name, alarml_action, priority, trigger_time, principal,
flow_id, start_time, end_time, is_alarm,is_inner
flow_id, start_time, end_time, is_alarm,is_inner
,fail_fast
</sql>
<select
id=
"findRunRecordingByFlowIdAndRunId"
resultMap=
"BaseResultMap"
>
...
...
@@ -103,6 +104,9 @@
<if
test=
"isInner != null"
>
is_inner,
</if>
<if
test=
"failFast != null"
>
fail_fast,
</if>
</trim>
<trim
prefix=
"values ("
suffix=
")"
suffixOverrides=
","
>
<if
test=
"recordingId != null"
>
...
...
@@ -159,10 +163,12 @@
<if
test=
"isInner != null"
>
#{isInner,jdbcType=CHAR},
</if>
<if
test=
"failFast != null"
>
#{fail_fast,jdbcType=CHAR},
</if>
</trim>
</insert>
<update
id=
"updateRunRecordingById"
parameterType=
"com.byit.model.RunRecording"
>
<!-- generated @mbg.generated date: 2019-12-25 -->
update run_recording
<set>
<if
test=
"runId != null"
>
...
...
@@ -216,6 +222,9 @@
<if
test=
"isInner != null"
>
is_inner = #{isInner,jdbcType=CHAR},
</if>
<if
test=
"failFast != null"
>
fail_fast = #{failFast,jdbcType=CHAR},
</if>
</set>
where recording_id = #{recordingId,jdbcType=INTEGER}
</update>
...
...
@@ -274,6 +283,9 @@
<if
test=
"isInner != null"
>
is_inner = #{isInner,jdbcType=CHAR},
</if>
<if
test=
"failFast != null"
>
fail_fast = #{failFast,jdbcType=CHAR},
</if>
</set>
where flow_id = #{flowId,jdbcType=INTEGER} and run_id = #{runId,jdbcType=VARCHAR}
</update>
...
...
byit-myth-core/myth-core-common/src/main/java/com/byit/job/enums/JobResultEnum.java
View file @
6e3322fc
...
...
@@ -5,20 +5,22 @@ package com.byit.job.enums;
* @author huangfu
*/
public
enum
JobResultEnum
implements
IEnum
{
SUCCESS
(
"100200"
,
"任务执行成功"
),
FAIL
(
"100500"
,
"任务执行失败"
),
FAIL_TIMEOUT
(
"100502"
,
"超时错误"
),
DISPATCH_SUCCESS
(
"200200"
,
"调度成功"
),
DISPATCH_FAIL
(
"200500"
,
"调度失败"
);
SUCCESS
(
"100200"
,
"任务执行成功"
,
"1"
),
FAIL
(
"100500"
,
"任务执行失败"
,
"2"
),
FAIL_TIMEOUT
(
"100502"
,
"超时错误"
,
"2"
),
DISPATCH_SUCCESS
(
"200200"
,
"调度成功"
,
"1"
),
DISPATCH_FAIL
(
"200500"
,
"调度失败"
,
"2"
);
private
String
code
;
private
String
msg
;
private
String
res
;
JobResultEnum
()
{
}
JobResultEnum
(
String
code
,
String
msg
)
{
JobResultEnum
(
String
code
,
String
msg
,
String
res
)
{
this
.
code
=
code
;
this
.
msg
=
msg
;
this
.
res
=
res
;
}
...
...
@@ -31,4 +33,7 @@ public enum JobResultEnum implements IEnum {
public
String
getMsg
()
{
return
this
.
msg
;
}
public
String
getRes
()
{
return
this
.
res
;
}
}
\ No newline at end of file
byit-myth-core/myth-core-common/src/main/java/com/byit/job/vo/ReturnResult.java
View file @
6e3322fc
...
...
@@ -18,9 +18,9 @@ import java.io.Serializable;
@NoArgsConstructor
public
class
ReturnResult
<
T
>
implements
Serializable
{
public
static
final
long
serialVersionUID
=
1573630876693L
;
public
static
final
ReturnResult
SUCCESS
=
new
ReturnResult
(
null
);
public
static
final
ReturnResult
FAIL
=
new
ReturnResult
(
JobResultEnum
.
FAIL
.
get
Code
(),
JobResultEnum
.
FAIL
.
getMsg
());
public
static
final
ReturnResult
FAIL_TIMEOUT
=
new
ReturnResult
(
JobResultEnum
.
FAIL_TIMEOUT
.
get
Code
(),
JobResultEnum
.
FAIL_TIMEOUT
.
getMsg
());
public
static
final
ReturnResult
<
String
>
SUCCESS
=
new
ReturnResult
<
String
>
(
null
);
public
static
final
ReturnResult
FAIL
=
new
ReturnResult
(
JobResultEnum
.
FAIL
.
get
Res
(),
JobResultEnum
.
FAIL
.
getMsg
());
public
static
final
ReturnResult
FAIL_TIMEOUT
=
new
ReturnResult
(
JobResultEnum
.
FAIL_TIMEOUT
.
get
Res
(),
JobResultEnum
.
FAIL_TIMEOUT
.
getMsg
());
private
String
code
;
private
String
msg
;
private
T
content
;
...
...
@@ -39,8 +39,8 @@ public class ReturnResult<T> implements Serializable {
* 默认就是成功
* @param content
*/
p
ublic
ReturnResult
(
T
content
)
{
this
.
code
=
JobResultEnum
.
SUCCESS
.
get
Code
();
p
rivate
ReturnResult
(
T
content
)
{
this
.
code
=
JobResultEnum
.
SUCCESS
.
get
Res
();
this
.
msg
=
JobResultEnum
.
SUCCESS
.
getMsg
();
this
.
content
=
content
;
}
...
...
byit-myth-executor/myth-exector-plugin/src/main/java/com/byit/rpc/RunJobThread.java
View file @
6e3322fc
...
...
@@ -20,7 +20,7 @@ import java.util.Date;
public
class
RunJobThread
implements
Runnable
{
private
AdminSenPluginDto
adminSenPluginDto
;
public
RunJobThread
(
AdminSenPluginDto
adminSenPluginDto
)
{
RunJobThread
(
AdminSenPluginDto
adminSenPluginDto
)
{
this
.
adminSenPluginDto
=
adminSenPluginDto
;
}
...
...
demo-client/byit-demo-client/src/main/java/com/byit/job/Mains.java
View file @
6e3322fc
...
...
@@ -16,13 +16,13 @@ import java.io.IOException;
public
class
Mains
{
public
static
void
main
(
String
[]
args
)
throws
IOException
{
new
JobRunServerLauncher
(
8888
);
String
plServerUrl
=
"http://10.0.55.23
7
:8888"
;
String
plServerUrl
=
"http://10.0.55.23
8
:8888"
;
String
mythCron
=
"时间"
;
String
routingStrategy
=
LoadBalance
.
ROUND
.
name
();
String
blockingStrategy
=
"阻塞策略"
;
String
callbackToken
=
"dsasadsad"
;
String
gatewayToken
=
"asdsadsa"
;
String
requestIP
=
"10.0.55.23
7
"
;
String
requestIP
=
"10.0.55.23
8
"
;
String
requestPort
=
"8080"
;
String
param
=
"不延迟任务"
;
String
name
=
"addJob"
;
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment