diff --git a/dolphinscheduler-alert/dolphinscheduler-alert-server/src/main/java/org/apache/dolphinscheduler/alert/AlertSenderService.java b/dolphinscheduler-alert/dolphinscheduler-alert-server/src/main/java/org/apache/dolphinscheduler/alert/AlertSenderService.java index af47913537..455b797760 100644 --- a/dolphinscheduler-alert/dolphinscheduler-alert-server/src/main/java/org/apache/dolphinscheduler/alert/AlertSenderService.java +++ b/dolphinscheduler-alert/dolphinscheduler-alert-server/src/main/java/org/apache/dolphinscheduler/alert/AlertSenderService.java @@ -17,7 +17,6 @@ package org.apache.dolphinscheduler.alert; -import org.apache.commons.collections.CollectionUtils; import org.apache.dolphinscheduler.alert.api.AlertChannel; import org.apache.dolphinscheduler.alert.api.AlertConstants; import org.apache.dolphinscheduler.alert.api.AlertData; @@ -34,16 +33,17 @@ import org.apache.dolphinscheduler.dao.entity.Alert; import org.apache.dolphinscheduler.dao.entity.AlertPluginInstance; import org.apache.dolphinscheduler.remote.command.alert.AlertSendResponseCommand; import org.apache.dolphinscheduler.remote.command.alert.AlertSendResponseResult; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; -import org.springframework.stereotype.Service; + +import org.apache.commons.collections.CollectionUtils; import java.util.ArrayList; -import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Optional; -import java.util.Set; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.stereotype.Service; @Service public final class AlertSenderService extends Thread { @@ -77,19 +77,19 @@ public final class AlertSenderService extends Thread { } } - public void send(List alerts) { for (Alert alert : alerts) { //get alert group from alert - int alertGroupId = alert.getAlertGroupId(); + int alertId = Optional.ofNullable(alert.getId()).orElse(0); + int alertGroupId = Optional.ofNullable(alert.getAlertGroupId()).orElse(0); List alertInstanceList = alertDao.listInstanceByAlertGroupId(alertGroupId); if (CollectionUtils.isEmpty(alertInstanceList)) { logger.error("send alert msg fail,no bind plugin instance."); - alertDao.updateAlert(AlertStatus.EXECUTION_FAILURE, "no bind plugin instance", alert.getId()); + alertDao.updateAlert(AlertStatus.EXECUTION_FAILURE, "no bind plugin instance", alertId); continue; } AlertData alertData = new AlertData(); - alertData.setId(alert.getId()) + alertData.setId(alertId) .setContent(alert.getContent()) .setLog(alert.getLog()) .setTitle(alert.getTitle()) @@ -101,7 +101,7 @@ public final class AlertSenderService extends Thread { AlertResult alertResult = this.alertResultHandler(instance, alertData); if (alertResult != null) { AlertStatus sendStatus = Boolean.parseBoolean(String.valueOf(alertResult.getStatus())) ? AlertStatus.EXECUTION_SUCCESS : AlertStatus.EXECUTION_FAILURE; - alertDao.addAlertSendStatus(sendStatus,alertResult.getMessage(),alert.getId(),instance.getId()); + alertDao.addAlertSendStatus(sendStatus,alertResult.getMessage(),alertId,instance.getId()); if (sendStatus.equals(AlertStatus.EXECUTION_SUCCESS)) { sendSuccessCount++; } @@ -113,7 +113,7 @@ public final class AlertSenderService extends Thread { } else if (sendSuccessCount < alertInstanceList.size()) { alertStatus = AlertStatus.EXECUTION_PARTIAL_SUCCESS; } - alertDao.updateAlert(alertStatus, "", alert.getId()); + alertDao.updateAlert(alertStatus, "", alertId); } } @@ -213,7 +213,8 @@ public final class AlertSenderService extends Thread { } if (!sendWarning) { - logger.info("Alert Plugin {} send ignore warning type not match: plugin warning type is {}, alert data warning type is {}", pluginInstanceName, warningType.getCode(), alertData.getWarnType()); + logger.info("Alert Plugin {} send ignore warning type not match: plugin warning type is {}, alert data warning type is {}", + pluginInstanceName, warningType.getCode(), alertData.getWarnType()); return null; } diff --git a/dolphinscheduler-alert/dolphinscheduler-alert-server/src/main/java/org/apache/dolphinscheduler/alert/AlertServer.java b/dolphinscheduler-alert/dolphinscheduler-alert-server/src/main/java/org/apache/dolphinscheduler/alert/AlertServer.java index 60bd09e6e2..1e17afefd1 100644 --- a/dolphinscheduler-alert/dolphinscheduler-alert-server/src/main/java/org/apache/dolphinscheduler/alert/AlertServer.java +++ b/dolphinscheduler-alert/dolphinscheduler-alert-server/src/main/java/org/apache/dolphinscheduler/alert/AlertServer.java @@ -24,6 +24,11 @@ import org.apache.dolphinscheduler.dao.PluginDao; import org.apache.dolphinscheduler.remote.NettyRemotingServer; import org.apache.dolphinscheduler.remote.command.CommandType; import org.apache.dolphinscheduler.remote.config.NettyServerConfig; + +import java.io.Closeable; + +import javax.annotation.PreDestroy; + import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.boot.WebApplicationType; @@ -33,9 +38,6 @@ import org.springframework.boot.context.event.ApplicationReadyEvent; import org.springframework.context.annotation.ComponentScan; import org.springframework.context.event.EventListener; -import javax.annotation.PreDestroy; -import java.io.Closeable; - @SpringBootApplication @ComponentScan("org.apache.dolphinscheduler") public class AlertServer implements Closeable { diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/AlertDao.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/AlertDao.java index b2d5c60a67..a2d2aacdf8 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/AlertDao.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/AlertDao.java @@ -19,9 +19,9 @@ package org.apache.dolphinscheduler.dao; import org.apache.dolphinscheduler.common.enums.AlertEvent; import org.apache.dolphinscheduler.common.enums.AlertStatus; +import org.apache.dolphinscheduler.common.enums.AlertType; import org.apache.dolphinscheduler.common.enums.AlertWarnLevel; import org.apache.dolphinscheduler.common.enums.WarningType; -import org.apache.dolphinscheduler.common.enums.AlertType; import org.apache.dolphinscheduler.common.utils.JSONUtils; import org.apache.dolphinscheduler.dao.entity.Alert; import org.apache.dolphinscheduler.dao.entity.AlertPluginInstance; @@ -36,21 +36,31 @@ import org.apache.dolphinscheduler.dao.mapper.AlertMapper; import org.apache.dolphinscheduler.dao.mapper.AlertPluginInstanceMapper; import org.apache.dolphinscheduler.dao.mapper.AlertSendStatusMapper; -import org.apache.commons.lang.StringUtils; +import org.apache.commons.codec.digest.DigestUtils; +import org.apache.commons.lang3.StringUtils; import java.util.ArrayList; import java.util.Arrays; import java.util.Date; import java.util.List; +import java.util.Optional; import java.util.stream.Collectors; +import org.joda.time.DateTime; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; +import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; +import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper; import com.google.common.collect.Lists; @Component public class AlertDao { + + @Value("${alert.alarm-suppression.crash:60}") + private Integer crashAlarmSuppression; + @Autowired private AlertMapper alertMapper; @@ -70,6 +80,8 @@ public class AlertDao { * @return add alert result */ public int addAlert(Alert alert) { + String sign = generateSign(alert); + alert.setSign(sign); return alertMapper.insert(alert); } @@ -82,13 +94,28 @@ public class AlertDao { * @return update alert result */ public int updateAlert(AlertStatus alertStatus, String log, int id) { - Alert alert = alertMapper.selectById(id); + Alert alert = new Alert(); + alert.setId(id); alert.setAlertStatus(alertStatus); alert.setUpdateTime(new Date()); alert.setLog(log); return alertMapper.updateById(alert); } + /** + * generate sign for alert + * + * @param alert alert + * @return sign's str + */ + private String generateSign (Alert alert) { + return Optional.of(alert) + .map(Alert::getContent) + .map(DigestUtils::sha1Hex) + .map(String::toLowerCase) + .orElse(StringUtils.EMPTY); + } + /** * add AlertSendStatus * @@ -109,13 +136,13 @@ public class AlertDao { } /** - * MasterServer or WorkerServer stoped + * MasterServer or WorkerServer stopped * * @param alertGroupId alertGroupId * @param host host * @param serverType serverType */ - public void sendServerStopedAlert(int alertGroupId, String host, String serverType) { + public void sendServerStoppedAlert(int alertGroupId, String host, String serverType) { ServerAlertContent serverStopAlertContent = ServerAlertContent.newBuilder(). type(serverType) .host(host) @@ -133,8 +160,11 @@ public class AlertDao { alert.setCreateTime(new Date()); alert.setUpdateTime(new Date()); alert.setAlertType(AlertType.FAULT_TOLERANCE_WARNING); + alert.setSign(generateSign(alert)); // we use this method to avoid insert duplicate alert(issue #5525) - alertMapper.insertAlertWhenServerCrash(alert); + // we modified this method to optimize performance(issue #9174) + Date crashAlarmSuppressionStartTime = DateTime.now().plusMinutes(-crashAlarmSuppression).toDate(); + alertMapper.insertAlertWhenServerCrash(alert, crashAlarmSuppressionStartTime); } /** @@ -178,6 +208,8 @@ public class AlertDao { alert.setContent(content); alert.setCreateTime(new Date()); alert.setUpdateTime(new Date()); + String sign = generateSign(alert); + alert.setSign(sign); alertMapper.insert(alert); } @@ -220,7 +252,9 @@ public class AlertDao { * List alerts that are pending for execution */ public List listPendingAlerts() { - return alertMapper.listAlertByStatus(AlertStatus.WAIT_EXECUTION); + LambdaQueryWrapper wrapper = new QueryWrapper<>(new Alert()).lambda() + .eq(Alert::getAlertStatus, AlertStatus.WAIT_EXECUTION); + return alertMapper.selectList(wrapper); } /** @@ -265,4 +299,8 @@ public class AlertDao { public void setAlertGroupMapper(AlertGroupMapper alertGroupMapper) { this.alertGroupMapper = alertGroupMapper; } + + public void setCrashAlarmSuppression(Integer crashAlarmSuppression) { + this.crashAlarmSuppression = crashAlarmSuppression; + } } diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/entity/Alert.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/entity/Alert.java index 388d294655..4b27a55b5a 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/entity/Alert.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/entity/Alert.java @@ -24,6 +24,7 @@ import org.apache.dolphinscheduler.common.enums.WarningType; import java.util.Date; import java.util.HashMap; import java.util.Map; +import java.util.Objects; import com.baomidou.mybatisplus.annotation.IdType; import com.baomidou.mybatisplus.annotation.TableField; @@ -36,7 +37,12 @@ public class Alert { * primary key */ @TableId(value = "id", type = IdType.AUTO) - private int id; + private Integer id; + /** + * sign + */ + @TableField(value = "sign") + private String sign; /** * title */ @@ -67,11 +73,11 @@ public class Alert { @TableField(value = "log") private String log; - /**g + /** * alertgroup_id */ @TableField("alertgroup_id") - private int alertGroupId; + private Integer alertGroupId; /** * create_time @@ -100,7 +106,7 @@ public class Alert { * process_instance_id */ @TableField("process_instance_id") - private int processInstanceId; + private Integer processInstanceId; /** * alert_type @@ -114,14 +120,22 @@ public class Alert { public Alert() { } - public int getId() { + public Integer getId() { return id; } - public void setId(int id) { + public void setId(Integer id) { this.id = id; } + public String getSign() { + return sign; + } + + public void setSign(String sign) { + this.sign = sign; + } + public String getTitle() { return title; } @@ -154,11 +168,11 @@ public class Alert { this.log = log; } - public int getAlertGroupId() { + public Integer getAlertGroupId() { return alertGroupId; } - public void setAlertGroupId(int alertGroupId) { + public void setAlertGroupId(Integer alertGroupId) { this.alertGroupId = alertGroupId; } @@ -210,11 +224,11 @@ public class Alert { this.processDefinitionCode = processDefinitionCode; } - public int getProcessInstanceId() { + public Integer getProcessInstanceId() { return processInstanceId; } - public void setProcessInstanceId(int processInstanceId) { + public void setProcessInstanceId(Integer processInstanceId) { this.processInstanceId = processInstanceId; } @@ -234,78 +248,40 @@ public class Alert { if (o == null || getClass() != o.getClass()) { return false; } - Alert alert = (Alert) o; - - if (id != alert.id) { - return false; - } - if (alertGroupId != alert.alertGroupId) { - return false; - } - if (!title.equals(alert.title)) { - return false; - } - if (!content.equals(alert.content)) { - return false; - } - if (alertStatus != alert.alertStatus) { - return false; - } - if (!log.equals(alert.log)) { - return false; - } - if (!createTime.equals(alert.createTime)) { - return false; - } - if (warningType != alert.warningType) { - return false; - } - return updateTime.equals(alert.updateTime) && info.equals(alert.info); - + return Objects.equals(id, alert.id) + && Objects.equals(alertGroupId, alert.alertGroupId) + && Objects.equals(sign, alert.sign) + && Objects.equals(title, alert.title) + && Objects.equals(content, alert.content) + && alertStatus == alert.alertStatus + && warningType == alert.warningType + && Objects.equals(log, alert.log) + && Objects.equals(createTime, alert.createTime) + && Objects.equals(updateTime, alert.updateTime) + && Objects.equals(info, alert.info) + ; } @Override public int hashCode() { - int result = id; - result = 31 * result + title.hashCode(); - result = 31 * result + content.hashCode(); - result = 31 * result + alertStatus.hashCode(); - result = 31 * result + warningType.hashCode(); - result = 31 * result + log.hashCode(); - result = 31 * result + alertGroupId; - result = 31 * result + createTime.hashCode(); - result = 31 * result + updateTime.hashCode(); - result = 31 * result + info.hashCode(); - return result; + return Objects.hash(id, sign, title, content, alertStatus, warningType, log, alertGroupId, createTime, updateTime, info); } @Override public String toString() { return "Alert{" - + "id=" - + id - + ", title='" - + title + '\'' - + ", content='" - + content - + '\'' - + ", alertStatus=" - + alertStatus - + ", warningType=" - + warningType - + ", log='" - + log - + '\'' - + ", alertGroupId=" - + alertGroupId - + '\'' - + ", createTime=" - + createTime - + ", updateTime=" - + updateTime - + ", info=" - + info + + "id=" + id + + ", sign='" + sign + '\'' + + ", title='" + title + '\'' + + ", content='" + content + '\'' + + ", alertStatus=" + alertStatus + + ", warningType=" + warningType + + ", log='" + log + '\'' + + ", alertGroupId=" + alertGroupId + + ", createTime=" + createTime + + ", updateTime=" + updateTime + + ", info=" + info + '}'; } } diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/AlertMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/AlertMapper.java index 77786c5a1e..1b6f29d43c 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/AlertMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/AlertMapper.java @@ -17,31 +17,25 @@ package org.apache.dolphinscheduler.dao.mapper; -import org.apache.dolphinscheduler.common.enums.AlertStatus; import org.apache.dolphinscheduler.dao.entity.Alert; +import org.apache.ibatis.annotations.Mapper; import org.apache.ibatis.annotations.Param; -import java.util.List; +import java.util.Date; import com.baomidou.mybatisplus.core.mapper.BaseMapper; /** * alert mapper interface */ +@Mapper public interface AlertMapper extends BaseMapper { - /** - * list alert by status - * @param alertStatus alertStatus - * @return alert list - */ - List listAlertByStatus(@Param("alertStatus") AlertStatus alertStatus); - /** * Insert server crash alert *

This method will ensure that there is at most one unsent alert which has the same content in the database. */ - void insertAlertWhenServerCrash(@Param("alert") Alert alert); + void insertAlertWhenServerCrash(@Param("alert") Alert alert, @Param("crashAlarmSuppressionStartTime") Date crashAlarmSuppressionStartTime); } diff --git a/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/AlertMapper.xml b/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/AlertMapper.xml index 912fba0595..f0f514ab52 100644 --- a/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/AlertMapper.xml +++ b/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/AlertMapper.xml @@ -18,24 +18,13 @@ - - id - , title, content, alert_status, log, - alertgroup_id, create_time, update_time, warning_type - - - insert into t_ds_alert(title, content, alert_status, warning_type, log, alertgroup_id, create_time, update_time) - SELECT #{alert.title}, #{alert.content}, #{alert.alertStatus.code}, #{alert.warningType.code}, #{alert.log}, #{alert.alertGroupId}, - #{alert.createTime}, #{alert.updateTime} + insert into t_ds_alert(sign, title, content, alert_status, warning_type, log, alertgroup_id, create_time, update_time) + SELECT #{alert.sign}, #{alert.title}, #{alert.content}, #{alert.alertStatus.code}, #{alert.warningType.code}, + #{alert.log}, #{alert.alertGroupId}, #{alert.createTime}, #{alert.updateTime} from t_ds_alert - where content = #{alert.content} and alert_status = #{alert.alertStatus.code} + where create_time >= #{crashAlarmSuppressionStartTime} and sign = #{alert.sign} and alert_status = #{alert.alertStatus.code} having count(*) = 0 diff --git a/dolphinscheduler-dao/src/main/resources/sql/dolphinscheduler_h2.sql b/dolphinscheduler-dao/src/main/resources/sql/dolphinscheduler_h2.sql index 7da05b3610..7fd9dc8256 100644 --- a/dolphinscheduler-dao/src/main/resources/sql/dolphinscheduler_h2.sql +++ b/dolphinscheduler-dao/src/main/resources/sql/dolphinscheduler_h2.sql @@ -272,6 +272,7 @@ CREATE TABLE t_ds_alert ( id int(11) NOT NULL AUTO_INCREMENT, title varchar(64) DEFAULT NULL, + sign char(40) NOT NULL DEFAULT '', content text, alert_status tinyint(4) DEFAULT '0', warning_type tinyint(4) DEFAULT '2', @@ -283,7 +284,8 @@ CREATE TABLE t_ds_alert process_definition_code bigint(20) DEFAULT NULL, process_instance_id int(11) DEFAULT NULL, alert_type int(11) DEFAULT NULL, - PRIMARY KEY (id) + PRIMARY KEY (id), + KEY idx_sign (sign) ); -- ---------------------------- diff --git a/dolphinscheduler-dao/src/main/resources/sql/dolphinscheduler_mysql.sql b/dolphinscheduler-dao/src/main/resources/sql/dolphinscheduler_mysql.sql index ca0c563039..051ea1d006 100644 --- a/dolphinscheduler-dao/src/main/resources/sql/dolphinscheduler_mysql.sql +++ b/dolphinscheduler-dao/src/main/resources/sql/dolphinscheduler_mysql.sql @@ -279,6 +279,7 @@ DROP TABLE IF EXISTS `t_ds_alert`; CREATE TABLE `t_ds_alert` ( `id` int(11) NOT NULL AUTO_INCREMENT COMMENT 'key', `title` varchar(64) DEFAULT NULL COMMENT 'title', + `sign` char(40) NOT NULL DEFAULT '' COMMENT 'sign=sha1(content)', `content` text COMMENT 'Message content (can be email, can be SMS. Mail is stored in JSON map, and SMS is string)', `alert_status` tinyint(4) DEFAULT '0' COMMENT '0:wait running,1:success,2:failed', `warning_type` tinyint(4) DEFAULT '2' COMMENT '1 process is successfully, 2 process/task is failed', @@ -291,7 +292,8 @@ CREATE TABLE `t_ds_alert` ( `process_instance_id` int(11) DEFAULT NULL COMMENT 'process_instance_id', `alert_type` int(11) DEFAULT NULL COMMENT 'alert_type', PRIMARY KEY (`id`), - KEY `idx_status` (`alert_status`) USING BTREE + KEY `idx_status` (`alert_status`) USING BTREE, + KEY `idx_sign` (`sign`) USING BTREE ) ENGINE=InnoDB AUTO_INCREMENT=1 DEFAULT CHARSET=utf8; -- ---------------------------- diff --git a/dolphinscheduler-dao/src/main/resources/sql/dolphinscheduler_postgresql.sql b/dolphinscheduler-dao/src/main/resources/sql/dolphinscheduler_postgresql.sql index 61119fb0b3..9e9c64d089 100644 --- a/dolphinscheduler-dao/src/main/resources/sql/dolphinscheduler_postgresql.sql +++ b/dolphinscheduler-dao/src/main/resources/sql/dolphinscheduler_postgresql.sql @@ -208,6 +208,7 @@ DROP TABLE IF EXISTS t_ds_alert; CREATE TABLE t_ds_alert ( id int NOT NULL , title varchar(64) DEFAULT NULL , + sign varchar(40) NOT NULL DEFAULT '', content text , alert_status int DEFAULT '0' , warning_type int DEFAULT '2' , @@ -220,9 +221,11 @@ CREATE TABLE t_ds_alert ( process_instance_id int DEFAULT NULL , alert_type int DEFAULT NULL , PRIMARY KEY (id) -) ; +); +comment on column t_ds_alert.sign is 'sign=sha1(content)'; create index idx_status on t_ds_alert (alert_status); +create index idx_sign on t_ds_alert (sign); -- -- Table structure for table t_ds_alertgroup diff --git a/dolphinscheduler-dao/src/main/resources/sql/upgrade/2.0.6_schema/mysql/dolphinscheduler_ddl.sql b/dolphinscheduler-dao/src/main/resources/sql/upgrade/2.0.6_schema/mysql/dolphinscheduler_ddl.sql index 45f8acd4da..33dcab5031 100644 --- a/dolphinscheduler-dao/src/main/resources/sql/upgrade/2.0.6_schema/mysql/dolphinscheduler_ddl.sql +++ b/dolphinscheduler-dao/src/main/resources/sql/upgrade/2.0.6_schema/mysql/dolphinscheduler_ddl.sql @@ -36,3 +36,24 @@ d// delimiter ; CALL uc_dolphin_T_t_ds_resources_R_full_name; DROP PROCEDURE uc_dolphin_T_t_ds_resources_R_full_name; + +-- uc_dolphin_T_t_ds_alert_R_sign +drop PROCEDURE if EXISTS uc_dolphin_T_t_ds_alert_R_sign; +delimiter d// +CREATE PROCEDURE uc_dolphin_T_t_ds_alert_R_sign() + BEGIN + IF NOT EXISTS (SELECT 1 FROM information_schema.COLUMNS + WHERE TABLE_NAME='t_ds_alert' + AND TABLE_SCHEMA=(SELECT DATABASE()) + AND COLUMN_NAME='sign') + THEN + ALTER TABLE `t_ds_alert` ADD COLUMN `sign` char(40) NOT NULL DEFAULT '' COMMENT 'sign=sha1(content)' after `id`; + ALTER TABLE `t_ds_alert` ADD INDEX `idx_sign` (`sign`) USING BTREE; + END IF; +END; + +d// + +delimiter ; +CALL uc_dolphin_T_t_ds_alert_R_sign; +DROP PROCEDURE uc_dolphin_T_t_ds_alert_R_sign; diff --git a/dolphinscheduler-dao/src/main/resources/sql/upgrade/2.0.6_schema/postgresql/dolphinscheduler_ddl.sql b/dolphinscheduler-dao/src/main/resources/sql/upgrade/2.0.6_schema/postgresql/dolphinscheduler_ddl.sql index 14a20fcd8e..71de54fe23 100644 --- a/dolphinscheduler-dao/src/main/resources/sql/upgrade/2.0.6_schema/postgresql/dolphinscheduler_ddl.sql +++ b/dolphinscheduler-dao/src/main/resources/sql/upgrade/2.0.6_schema/postgresql/dolphinscheduler_ddl.sql @@ -32,6 +32,10 @@ BEGIN --- alter column EXECUTE 'ALTER TABLE ' || quote_ident(v_schema) ||'.t_ds_resources ALTER COLUMN full_name Type varchar(128)'; + --- add column + EXECUTE 'ALTER TABLE ' || quote_ident(v_schema) ||'.t_ds_alert ADD COLUMN IF NOT EXISTS sign varchar(40) NOT NULL DEFAULT '''' '; + EXECUTE 'comment on column ' || quote_ident(v_schema) ||'.t_ds_alert.sign is ''sign=sha1(content)'''; + return 'Success!'; exception when others then ---Raise EXCEPTION '(%)',SQLERRM; diff --git a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/AlertDaoTest.java b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/AlertDaoTest.java index 820ac793be..fe3efbf8de 100644 --- a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/AlertDaoTest.java +++ b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/AlertDaoTest.java @@ -68,12 +68,12 @@ public class AlertDaoTest { } @Test - public void testSendServerStopedAlert() { + public void testSendServerStoppedAlert() { int alertGroupId = 1; String host = "127.0.0.998165432"; String serverType = "Master"; - alertDao.sendServerStopedAlert(alertGroupId, host, serverType); - alertDao.sendServerStopedAlert(alertGroupId, host, serverType); + alertDao.sendServerStoppedAlert(alertGroupId, host, serverType); + alertDao.sendServerStoppedAlert(alertGroupId, host, serverType); long count = alertDao.listPendingAlerts() .stream() .filter(alert -> alert.getContent().contains(host)) diff --git a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/AlertMapperTest.java b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/AlertMapperTest.java index b92d47e34e..cd4b87e0a2 100644 --- a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/AlertMapperTest.java +++ b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/AlertMapperTest.java @@ -28,8 +28,9 @@ import org.apache.dolphinscheduler.common.utils.DateUtils; import org.apache.dolphinscheduler.dao.BaseDaoTest; import org.apache.dolphinscheduler.dao.entity.Alert; +import org.apache.commons.codec.digest.DigestUtils; + import java.util.HashMap; -import java.util.List; import java.util.Map; import org.junit.Test; @@ -99,26 +100,6 @@ public class AlertMapperTest extends BaseDaoTest { assertNull(actualAlert); } - /** - * test list alert by status - */ - @Test - public void testListAlertByStatus() { - Integer count = 10; - AlertStatus waitExecution = AlertStatus.WAIT_EXECUTION; - - Map expectedAlertMap = createAlertMap(count, waitExecution); - - List actualAlerts = alertMapper.listAlertByStatus(waitExecution); - - for (Alert actualAlert : actualAlerts) { - Alert expectedAlert = expectedAlertMap.get(actualAlert.getId()); - if (expectedAlert != null) { - assertEquals(expectedAlert, actualAlert); - } - } - } - /** * create alert map * @@ -153,9 +134,11 @@ public class AlertMapperTest extends BaseDaoTest { * @return alert */ private Alert createAlert(AlertStatus alertStatus) { + String content = "[{'type':'WORKER','host':'192.168.xx.xx','event':'server down','warning level':'serious'}]"; Alert alert = new Alert(); alert.setTitle("test alert"); - alert.setContent("[{'type':'WORKER','host':'192.168.xx.xx','event':'server down','warning level':'serious'}]"); + alert.setContent(content); + alert.setSign(DigestUtils.sha1Hex(content)); alert.setAlertStatus(alertStatus); alert.setWarningType(WarningType.FAILURE); alert.setLog("success"); diff --git a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/registry/ServerNodeManager.java b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/registry/ServerNodeManager.java index 0666ca82a0..a6599b6032 100644 --- a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/registry/ServerNodeManager.java +++ b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/registry/ServerNodeManager.java @@ -246,7 +246,7 @@ public class ServerNodeManager implements InitializingBean { String group = parseGroup(path); Collection currentNodes = registryClient.getWorkerGroupNodesDirectly(group); syncWorkerGroupNodes(group, currentNodes); - alertDao.sendServerStopedAlert(1, path, "WORKER"); + alertDao.sendServerStoppedAlert(1, path, "WORKER"); } else if (type == Type.UPDATE) { logger.debug("worker group node : {} update, data: {}", path, data); String group = parseGroup(path); @@ -296,7 +296,7 @@ public class ServerNodeManager implements InitializingBean { if (type.equals(Type.REMOVE)) { logger.info("master node : {} down.", path); updateMasterNodes(); - alertDao.sendServerStopedAlert(1, path, "MASTER"); + alertDao.sendServerStoppedAlert(1, path, "MASTER"); } } catch (Exception ex) { logger.error("MasterNodeListener capture data change and get data failed.", ex);