Commit 0c38bcbb by guominglei

Merge remote-tracking branch 'origin/developer' into developer

parents 5a46c038 96f234b1
......@@ -37,7 +37,7 @@ public class MythJobScheduler implements InitializingBean, DisposableBean {
public void destroy() throws Exception {
this.jobScheduleHelper.doStop();
this.logScanHelper.doStop();
runRecordingScanHelper.doStop();
this.runRecordingScanHelper.doStop();
}
/**
......@@ -49,7 +49,7 @@ public class MythJobScheduler implements InitializingBean, DisposableBean {
//启用扫描线程
this.jobScheduleHelper.start();
this.logScanHelper.start();
runRecordingScanHelper.start();
this.runRecordingScanHelper.start();
}
}
......@@ -47,6 +47,7 @@ public interface RunRecordingService {
*/
int updateRunRecordingByFlowIdAndRunId(RunRecording record);
/**
* 根据ID删除
* @param recordingId
......
......@@ -32,12 +32,10 @@ public class LogScanHelper {
private DataSource dataSource;
private final JobTaskRunLogService jobTaskRunLogService;
private final RunRecordingService runRecordingService;
private final EmailAlarmService emailAlarmService;
@Autowired
public LogScanHelper(JobTaskRunLogService jobTaskRunLogService, RunRecordingService runRecordingService, EmailAlarmService emailAlarmService) {
public LogScanHelper(JobTaskRunLogService jobTaskRunLogService, RunRecordingService runRecordingService) {
this.jobTaskRunLogService = jobTaskRunLogService;
this.runRecordingService = runRecordingService;
this.emailAlarmService = emailAlarmService;
}
@Autowired
......@@ -75,6 +73,7 @@ public class LogScanHelper {
//扫描虚节点线程
virtualNodeScanThread = new Thread(() ->{
dateAligned(5000);
log.info("---------------------【com.byit.thread.LogScanHelper#virtualNodeScanMethod】init success------------------------");
while (!virtualNodeScanIsStop){
Connection conn = null;
Boolean connAutoCommit = null;
......@@ -160,8 +159,9 @@ public class LogScanHelper {
*/
private void notAlarmedNodeScanMethod(){
notAlarmedNodeScanThread = new Thread(() ->{
dateAligned(5000);
log.info("---------------------【com.byit.thread.LogScanHelper#notAlarmedNodeScanMethod】init success------------------------");
while (!notAlarmedNodeScanIsStop){
dateAligned(5000);
Connection conn = null;
Boolean connAutoCommit = null;
PreparedStatement preparedStatement = null;
......@@ -283,7 +283,7 @@ public class LogScanHelper {
try {
TimeUnit.MILLISECONDS.sleep(waitTime - System.currentTimeMillis()%1000);
} catch (InterruptedException e) {
e.printStackTrace( );
log.warn("----------------【线程被中断】-----------------------");
}
}
}
......@@ -9,6 +9,9 @@ import com.byit.model.RunRecording;
import com.byit.service.EmailAlarmService;
import com.byit.service.JobTaskRunLogService;
import com.byit.service.RunRecordingService;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
......@@ -19,8 +22,10 @@ import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.util.ArrayList;
import java.util.List;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
/**
* 运行记录扫描线程
......@@ -39,6 +44,9 @@ public class RunRecordingScanHelper {
* 失败时告警
*/
public static final String FAILURE_DONE = "2";
/**
* 成功时告警
*/
public static final String SUCCESS_DONE = "3";
private DataSource dataSource;
......@@ -58,7 +66,9 @@ public class RunRecordingScanHelper {
}
public void start(){
System.out.println("------------------------");
runRecordingThread = new Thread(() ->{
dateAligned(5000);
log.info("--------------------【com.byit.thread.RunRecordingScanThread#start】init success----------------------");
while (!runRecordingThreadStop){
dateAligned(5000);
......@@ -115,13 +125,10 @@ public class RunRecordingScanHelper {
}
}
}
runRecordingThread.setDaemon(true);
runRecordingThread.setName("myth-job#【RunRecordingScanThread】# start");
runRecordingThread.start();
});
runRecordingThread.setDaemon(true);
runRecordingThread.setName("myth-job#【RunRecordingScanThread】# start");
runRecordingThread.start();
}
private void scanRunRec(){
......@@ -135,33 +142,27 @@ public class RunRecordingScanHelper {
//设置为完成时告警
case WHEN_DONE:
if(RunRecordingEnum.FLOW_STATUS_IS_END.getCode().equals(runRecording.getFlowStatus())){
saveEmailAlarms(emailAlarms,runRecording);
saveEmailAlarms(runRecording);
}
break;
//失败时告警
case FAILURE_DONE:
if(RunRecordingEnum.RUN_FLOW_FAILURE.getCode().equals(runRecording.getFlowRunResult())
|| RunRecordingEnum.RUN_FLOW_RE_FAILURE.getCode().equals(runRecording.getFlowRunResult())){
saveEmailAlarms(emailAlarms,runRecording);
saveEmailAlarms(runRecording);
}
break;
//成功时告警
case SUCCESS_DONE:
if(RunRecordingEnum.RUN_FLOW_SUCCESS.getCode().equals(runRecording.getFlowRunResult())
|| RunRecordingEnum.RUN_FLOW_RE_SUCCESS.getCode().equals(runRecording.getFlowRunResult())){
saveEmailAlarms(emailAlarms,runRecording);
saveEmailAlarms(runRecording);
}
break;
default:
break;
}
});
if(CollectionUtil.isNotEmpty(emailAlarms)){
int saveCount = emailAlarmService.saveEmailAlarms(emailAlarms);
}else{
dateAligned(20000);
}
}
}
......@@ -192,15 +193,15 @@ public class RunRecordingScanHelper {
try {
TimeUnit.MILLISECONDS.sleep(waitTime - System.currentTimeMillis()%1000);
} catch (InterruptedException e) {
e.printStackTrace( );
log.warn("----------------【线程被中断】-----------------------");
}
}
private void saveEmailAlarms(List<EmailAlarm> emailAlarms,RunRecording runRecording){
private void saveEmailAlarms(RunRecording runRecording){
List<JobTaskRunLogWithBLOBs> jobTaskRunLogByFlowIdAndRunId = jobTaskRunLogService.findJobTaskRunLogWithBLOBsByFlowIdAndRunId(runRecording.getFlowId(), runRecording.getRunId());
String flowName = runRecording.getFlowName();
String senContentHtml = runMsgHtml(jobTaskRunLogByFlowIdAndRunId, flowName);
EmailAlarm build = EmailAlarm.builder()
EmailAlarm emailAlarm = EmailAlarm.builder()
.flowId(runRecording.getFlowId())
.flowName(flowName)
.runId(runRecording.getRunId())
......@@ -212,7 +213,11 @@ public class RunRecordingScanHelper {
.flowRes(runRecording.getFlowRunResult())
.alarmTitle(flowName)
.build();
emailAlarms.add(build);
//保存邮箱
emailAlarmService.saveEmailAlarm(emailAlarm);
//修改为已告警
runRecording.setIsAlarm("0");
runRecordingService.updateRunRecordingById(runRecording);
}
private String runMsgHtml(List<JobTaskRunLogWithBLOBs> jobTaskRunLogs,String title){
......@@ -257,4 +262,4 @@ public class RunRecordingScanHelper {
public void setDataSource(DataSource dataSource) {
this.dataSource = dataSource;
}
}
}
\ No newline at end of file
package com.byit.test;
/**
* @author huangfu
*/
public class ExploringSynchronized implements Runnable {
private static final String LOCK_MARK = "LOCK_MARK";
/**
* 共享资源(临界资源)
*/
static int i=0;
public void add(){
synchronized (LOCK_MARK){
i++;
}
}
@Override
public void run() {
for (int j = 0; j < 100000; j++) {
add();
}
}
public static void main(String[] args) throws InterruptedException {
Thread t1 = new Thread(new ExploringSynchronized());
Thread t2 = new Thread(new ExploringSynchronized());
t1.start();
t2.start();
//join 主线程需要等待子线程完成后在结束
t1.join();
t2.join();
System.out.println(i);
}
}
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