email分配发送任务,email发送

This commit is contained in:
万佳 2021-09-17 10:52:07 +08:00
parent a9224761cc
commit 75754b5246
17 changed files with 236 additions and 65 deletions

View File

@ -2,6 +2,7 @@ package cn.org.gitlink.notification.common.utils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.PropertySource;
import org.springframework.mail.javamail.JavaMailSender;
import org.springframework.mail.javamail.MimeMessageHelper;
@ -14,6 +15,9 @@ import javax.mail.internet.MimeMessage;
@Service(value = "GitlinkMailUtils")
public class EmailUtils {
@Value("${spring.mail.username}")
private String emailFrom;
@Autowired
private JavaMailSender mailSender;
@ -21,6 +25,7 @@ public class EmailUtils {
MimeMessage mail = mailSender.createMimeMessage();
MimeMessageHelper helper = new MimeMessageHelper(mail, true);
helper.setTo(recipient);
helper.setFrom(emailFrom);
helper.setSubject(subject);
helper.setText(content, true);
mailSender.send(mail);

26
db/gns-email.sql Normal file
View File

@ -0,0 +1,26 @@
DROP TABLE IF EXISTS `gitlink_email_jobs`;
CREATE TABLE `gitlink_email_jobs` (
`id` int(11) NOT NULL AUTO_INCREMENT,
`sender` int(11) NOT NULL COMMENT '发送者id',
`emails` text NOT NULL COMMENT '收件人全部邮件地址',
`subject` varchar(500) DEFAULT NULL COMMENT '邮件主题',
`content` text NOT NULL COMMENT '邮件内容',
`created_at` datetime DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`dispatched_at` datetime DEFAULT NULL COMMENT '处理时间',
`dispatched_status` int(11) DEFAULT '-1' COMMENT '发送状态:-1 未处理,1 处理成功,2 处理失败',
PRIMARY KEY (`id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb3;
DROP TABLE IF EXISTS `gitlink_email_send_records`;
CREATE TABLE `gitlink_email_send_records` (
`id` int(11) NOT NULL AUTO_INCREMENT,
`email` varchar(500) NOT NULL DEFAULT '' COMMENT '收件人邮件地址',
`job_id` int(11) NOT NULL COMMENT '邮件详情id',
`created_at` datetime DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`sent_at` datetime DEFAULT NULL COMMENT '发送时间',
`status` int(11) DEFAULT '-1' COMMENT '发送状态:-1 未发送,1 发送成功,2 发送失败',
PRIMARY KEY (`id`),
KEY `index_on_email_and_status` (`email`,`status`),
KEY `index_on_status` (`status`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb3;

View File

@ -1,14 +1,20 @@
package cn.org.gitlink.notification.executor.service.email;
import cn.org.gitlink.notification.common.utils.EmailUtils;
import cn.org.gitlink.notification.model.dao.entity.EmailJob;
import cn.org.gitlink.notification.model.dao.entity.EmailSendRecord;
import cn.org.gitlink.notification.model.dao.entity.vo.NewEmailJobVo;
import cn.org.gitlink.notification.model.service.notification.EmailJobsService;
import cn.org.gitlink.notification.model.service.notification.EmailSendRecordsService;
import com.alibaba.fastjson.JSONObject;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import javax.mail.MessagingException;
import java.beans.Transient;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;
@ -24,14 +30,18 @@ public class EmailService {
@Autowired
private EmailSendRecordsService emailSendRecordsService;
@Autowired
private EmailUtils emailUtils;
/**
* @Description: 处理邮件发送任务根据emails添加到邮件发送记录表中
* @Param platform 平台编码
* @Param dispatchNumber 待处理发送任务列表数量
* 处理邮件发送任务根据emails添加到邮件发送记录表中
* @param platform 平台编码
* @param dispatchNumber 待处理发送任务列表数量
* @return: void
* @Author: wanjia
* @Date: 2021/9/15
*/
@Transactional
public void DispatchEmailJobs(String platform, Integer dispatchNumber) {
//获取指定数量待处理列表
List<EmailJob> emailJobList = new ArrayList<>();
@ -59,14 +69,16 @@ public class EmailService {
}
/**
* @Description: 发送邮件
* @Param platform 平台编码
* @Param sentNumber 一次发送数量
* 发送邮件
* @param platform 平台编码
* @param sentNumber 一次发送数量
* @return: void
* @Author: wanjia
* @Date: 2021/9/15
*/
@Transactional
public void sendEmail(String platform, Integer sentNumber) {
//获取待发送列表
List<EmailSendRecord> emailSendRecordList = new ArrayList<>();
try {
emailSendRecordList = emailSendRecordsService.getRecordsByStatus(platform, -1, sentNumber);
@ -74,14 +86,18 @@ public class EmailService {
logger.error("获取未发送邮件列表失败:\n" + e);
}
//发送邮件
Boolean flag = null;
for (EmailSendRecord unSentEmailSendRecord : emailSendRecordList) {
//todo 发邮件
// 取到一条 emailSendRecord 调用发邮件的Util返回发送结果赋值给flag
unSentEmailSendRecord.setSentAt(new Date());
unSentEmailSendRecord.setStatus(flag ? 1 : 2);
flag = false;
try {
emailUtils.sendMail(unSentEmailSendRecord.getSubject(), unSentEmailSendRecord.getEmail(), unSentEmailSendRecord.getContent());
flag = true;
unSentEmailSendRecord.setSentAt(new Date());
unSentEmailSendRecord.setStatus(flag ? 1 : 2);
} catch (MessagingException e) {
logger.error("发送邮件失败email: " + unSentEmailSendRecord.getEmail() + "\n" + e);
}
}
//更新emailSendRecord状态
try {

View File

@ -1,20 +1,24 @@
package cn.org.gitlink.notification.executor.service.jobhandler;
import cn.org.gitlink.notification.common.utils.KafkaUtil;
import cn.org.gitlink.notification.model.dao.entity.vo.NewEmailJobVo;
import cn.org.gitlink.notification.model.service.notification.EmailJobsService;
import com.alibaba.fastjson.JSONObject;
import org.apache.kafka.clients.admin.NewTopic;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.annotation.KafkaHandler;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
import java.util.Arrays;
@Component
@Configuration
//todo 邮件的topics和groupId待处理
//@KafkaListener(topics = "${spring.kafka.consumer.topic}", groupId = "${spring.kafka.consumer.group_id}")
@KafkaListener(topics = "${spring.kafka.consumer.topic_email}", groupId = "${spring.kafka.consumer.group_id_email}")
public class EmailJobsListener {
private Logger logger = LogManager.getLogger(EmailJobsListener.class);
@ -22,11 +26,28 @@ public class EmailJobsListener {
@Autowired
private EmailJobsService emailJobsService;
@Value("${spring.kafka.producer.topic_new_email_remind}")
private String gitlinkNewEmailRemindTopic;
@Value("${spring.kafka.producer.partitions}")
private Integer partitions;
@Value("${spring.kafka.producer.replication_factor}")
private Short replicationFactor;
@Autowired
private KafkaUtil kafkaUtil;
@KafkaHandler
public void messageHandler(String message) {
try {
NewEmailJobVo newEmailJobVo = JSONObject.parseObject(message, NewEmailJobVo.class);
emailJobsService.sendEmail(newEmailJobVo);
Boolean flag = emailJobsService.sendEmail(newEmailJobVo);
//if the message is inserted successfully, send a new email-job message to kafka
if (flag){
kafkaUtil.createTopic(Arrays.asList(new NewTopic(gitlinkNewEmailRemindTopic, partitions, replicationFactor)));
kafkaUtil.sendMessage(gitlinkNewEmailRemindTopic, JSONObject.toJSONString(newEmailJobVo));
}
} catch (Exception e) {
logger.error(e);
}

View File

@ -0,0 +1,39 @@
package cn.org.gitlink.notification.executor.service.jobhandler;
import cn.org.gitlink.notification.executor.service.email.EmailService;
import cn.org.gitlink.notification.model.dao.entity.vo.NewEmailJobVo;
import com.alibaba.fastjson.JSONObject;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.annotation.KafkaHandler;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
@Component
@Configuration
@KafkaListener(topics = "${spring.kafka.consumer.topic_new_email_remind}", groupId = "${spring.kafka.consumer.group_id_email}")
public class EmailSentListener {
private Logger logger = LogManager.getLogger(EmailSentListener.class);
@Autowired
private EmailService emailService;
@KafkaHandler
public void messageHandler(String message) {
NewEmailJobVo newEmailJobVo = JSONObject.parseObject(message, NewEmailJobVo.class);
//处理邮件任务
emailService.DispatchEmailJobs(newEmailJobVo.getPlatform(), 1);
//发送邮件
emailService.sendEmail(newEmailJobVo.getPlatform(), 100);
}
@KafkaHandler(isDefault = true)
public void defaultHandler(Object object) {
logger.error("Unknow object received. ->" + object);
}//end of method
}

View File

@ -11,13 +11,20 @@ spring:
client_id: gitlink_producer_01
retries: 5
batch_size: 16384
replication_factor: 1
partitions: 3
topic_new_email_remind: topic-gitlink-new-email-remind
consumer:
bootstrap_servers: kafka1:9092,kafka2:9092
client_id: gitlink_consumer
group_id: group-gitlink-notification
group_id_email: group-gitlink-email
group_id_new_email_remind: group-gitlink-new-email-remind
auto_offset_reset: earliest
max_poll_records: 100
topic: topic-gitlink-notification
topic_email: topic-gitlink-email
topic_new_email_remind: topic-gitlink-new-email-remind
enable-auto-commit: true
#自动提交的时间间隔
auto-commit-interval: 1S

View File

@ -1,10 +1,12 @@
spring.mail.host=smtp.exmail.qq.com
spring.mail.default-encoding=utf-8
spring.mail.port=465
spring.mail.username=xxx
spring.mail.password=xxxx
spring.mail.username=gitlink@barats.cn
spring.mail.password=
spring.mail.properties[mail.smtp.auth]=true
spring.mail.properties[mail.smtp.connectiontimeout]=10000
spring.mail.properties[mail.smtp.timeout]=10000
spring.mail.properties[mail.smtp.writetimeout]=10000
spring.mail.properties[mail.smtp.starttls.enable]=true
spring.mail.properties[mail.smtp.starttls.enable]=true
spring.mail.properties[mail.smtp.socketFactory.port]=465
spring.mail.properties[mail.smtp.socketFactory.class]=javax.net.ssl.SSLSocketFactory

View File

@ -16,6 +16,7 @@ services:
- ${SQL_SCRIPT_PATH}/xxl-job-structure.sql:/docker-entrypoint-initdb.d/0000.sql
- ${SQL_SCRIPT_PATH}/gns-notification.sql:/docker-entrypoint-initdb.d/0001.sql
- ${SQL_SCRIPT_PATH}/hehui-gns-notification.sql:/docker-entrypoint-initdb.d/0002.sql
- ${SQL_SCRIPT_PATH}/gns-email.sql:/docker-entrypoint-initdb.d/0003.sql
command: --character-set-server=utf8mb4 --collation-server=utf8mb4_unicode_ci
ports:
- ${MYSQL_LOCAL_PORT}:3306

View File

@ -1,6 +1,7 @@
package cn.org.gitlink.notification.model.dao.entity;
import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableField;
import com.baomidou.mybatisplus.annotation.TableId;
import com.fasterxml.jackson.annotation.JsonFormat;
import com.fasterxml.jackson.annotation.JsonProperty;
@ -25,6 +26,12 @@ public class EmailSendRecord {
private Integer status;
@TableField(exist = false)
private String content;
@TableField(exist = false)
private String subject;
public Integer getId() {
return id;
}
@ -72,4 +79,20 @@ public class EmailSendRecord {
public void setStatus(Integer status) {
this.status = status;
}
public String getContent() {
return content;
}
public void setContent(String content) {
this.content = content;
}
public String getSubject() {
return subject;
}
public void setSubject(String subject) {
this.subject = subject;
}
}

View File

@ -10,9 +10,9 @@ import java.util.List;
public interface EmailJobsService extends IService<EmailJob> {
/**
* @Description: 添加发送邮件任务
* 添加发送邮件任务
*
* @Param newEmailJobVo
* @param newEmailJobVo
* @return: boolean
* @Author: wanjia
* @Date: 2021/9/13
@ -20,11 +20,11 @@ public interface EmailJobsService extends IService<EmailJob> {
boolean sendEmail(NewEmailJobVo newEmailJobVo) throws Exception;
/**
* @Description: 获取所有未处理邮件任务列表
* 获取所有未处理邮件任务列表
*
* @Param platform 平台编码
* @Param status 发送状态-1 未处理,1 处理成功,2 处理失败
* @Param size 列表大小
* @param platform 平台编码
* @param dispatchedStatus 发送状态-1 未处理,1 处理成功,2 处理失败
* @param size 列表大小
* @return: List<EmailJob>
* @Author: wanjia
* @Date: 2021/9/13
@ -32,12 +32,12 @@ public interface EmailJobsService extends IService<EmailJob> {
List<EmailJob> getEmailJobsByDispatchedStatus(String platform,Integer dispatchedStatus, Integer size) throws Exception;
/**
* @Description: 更新邮件任务状态
* 更新邮件任务状态
*
* @Param platform 平台编码
* @Param emailJobId 邮件任务id
* @Param dispatchedAt 处理时间
* @Param dispatchedStatus 处理状态 -1 未处理,1 处理成功,2 处理失败
* @param platform 平台编码
* @param emailJobId 邮件任务id
* @param dispatchedAt 处理时间
* @param dispatchedStatus 处理状态 -1 未处理,1 处理成功,2 处理失败
* @return: int
* @Author: wanjia
* @Date: 2021/9/13

View File

@ -8,11 +8,11 @@ import java.util.List;
public interface EmailSendRecordsService extends IService<EmailSendRecord> {
/**
* @Description: 新增邮件任务到邮件发送记录表中
* 新增邮件任务到邮件发送记录表中
*
* @Param platform 平台编码
* @Param emails 邮件地址eg: w@163.com,j@163.com
* @Param jobId 邮件任务id
* @param platform 平台编码
* @param emails 邮件地址eg: w@163.com,j@163.com
* @param jobId 邮件任务id
* @return: boolean
* @Author: wanjia
* @Date: 2021/9/13
@ -20,11 +20,11 @@ public interface EmailSendRecordsService extends IService<EmailSendRecord> {
boolean newEmailSendRecords(String platform, String emails, Integer jobId) throws Exception;
/**
* @Description: 获取发送记录列表
* 获取发送记录列表
*
* @Param platform 平台编码
* @Param status 邮件发送记录状态
* @Param size 列表数量
* @param platform 平台编码
* @param status 邮件发送记录状态 -1未发送 1发送成功 2发送失败
* @param size 列表数量
* @return: List<EmailSendRecord>
* @Author: wanjia
* @Date: 2021/9/15
@ -32,10 +32,10 @@ public interface EmailSendRecordsService extends IService<EmailSendRecord> {
List<EmailSendRecord> getRecordsByStatus(String platform, Integer status, Integer size) throws Exception;
/**
* @Description: 变更邮件发送记录状态
* 变更邮件发送记录状态
*
* @Param platform 平台编码
* @Param List<EmailSendRecord> 待更新列表
* @param platform 平台编码
* @param emailSendRecordList 待更新列表
* @return: int
* @Author: wanjia
* @Date: 2021/9/13

View File

@ -127,6 +127,9 @@
</update>
<select id="getEmailJobsByDispatchedStatus" resultMap="BaseResultMap">
select * from ${platform}_email_jobs
where dispatched_status = #{dispatchedStatus} limit #{size}
where dispatched_status = #{dispatchedStatus} order by id desc
<if test="size != null">
limit #{size}
</if>
</select>
</mapper>

View File

@ -118,32 +118,37 @@
</foreach >
</insert >
<select id="getRecordsByStatus" resultMap="BaseResultMap">
select * from ${platform}_email_send_records
where status = #{status} limit #{size}
select a.id, a.email, a.job_id, a.created_at, a.sent_at, a.status, b.subject, b.content
from ${platform}_email_send_records a
left join ${platform}_email_jobs b on a.job_id = b.id
where status = #{status} order by id desc
<if test="size != null">
limit #{size}
</if>
</select>
<update id="updateEmailSendRecordsBatch" parameterType="list">
update ${platform}_email_send_records
<trim prefix="set" suffixOverrides=",">
<trim prefix="sentAt =case" suffix="end,">
<foreach collection="list" item="i" index="index">
<if test="i.sentAt!=null">
when id=#{i.id} then #{i.sentAt}
<trim prefix="sent_at =case" suffix="end,">
<foreach collection="list" item="item" index="index">
<if test="item.sentAt!=null">
when id=#{item.id} then #{item.sentAt}
</if>
</foreach>
</trim>
<trim prefix=" status =case" suffix="end,">
<foreach collection="list" item="i" index="index">
<if test="i.status!=null">
when id=#{i.id} then #{i.status}
<foreach collection="list" item="item" index="index">
<if test="item.status!=null">
when id=#{item.id} then #{item.status}
</if>
</foreach>
</trim>
</trim>
where
<foreach collection="list" separator="or" item="i" index="index" >
id=#{i.id}
where id in
<foreach collection="list" item="item" separator="," open="(" close=")">
#{item.id}
</foreach>
</update>
</mapper>

View File

@ -0,0 +1,10 @@
spring.mail.host=smtp.exmail.qq.com
spring.mail.default-encoding=utf-8
spring.mail.port=465
spring.mail.username=gitlink@barats.cn
spring.mail.password=
spring.mail.properties[mail.smtp.auth]=true
spring.mail.properties[mail.smtp.connectiontimeout]=10000
spring.mail.properties[mail.smtp.timeout]=10000
spring.mail.properties[mail.smtp.writetimeout]=10000
spring.mail.properties[mail.smtp.starttls.enable]=true

View File

@ -5,27 +5,24 @@ import cn.org.gitlink.notification.common.response.DataPacketUtil;
import cn.org.gitlink.notification.common.response.ResponseData;
import cn.org.gitlink.notification.common.utils.KafkaUtil;
import cn.org.gitlink.notification.common.utils.ValidatorUtils;
import cn.org.gitlink.notification.model.dao.entity.EmailJob;
import cn.org.gitlink.notification.model.dao.entity.vo.NewEmailJobParamsVo;
import cn.org.gitlink.notification.model.dao.entity.vo.NewEmailJobVo;
import cn.org.gitlink.notification.model.dao.entity.vo.NewSysNotificationVo;
import cn.org.gitlink.notification.model.service.notification.EmailJobsService;
import cn.org.gitlink.notification.model.service.notification.EmailSendRecordsService;
import com.alibaba.fastjson.JSONObject;
import io.swagger.annotations.ApiOperation;
import io.swagger.annotations.ApiParam;
import org.apache.kafka.clients.admin.NewTopic;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.springframework.beans.BeanUtils;
import org.springframework.beans.BeansException;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Configuration;
import org.springframework.validation.BindingResult;
import org.springframework.validation.annotation.Validated;
import org.springframework.web.bind.annotation.*;
import java.util.Date;
import java.util.List;
import java.util.Arrays;
import java.util.Map;
@RestController
@ -33,11 +30,17 @@ import java.util.Map;
@Configuration
public class EmailJobsController {
//todo yml文件中配置topic后替换
private static final String GITLINK_EMAIL_TOPIC = "topic_gitlink_email";
private Logger logger = LogManager.getLogger(EmailJobsController.class);
@Value("${spring.kafka.producer.topic_email}")
private String gitlinkEmailTopic;
@Value("${spring.kafka.producer.partitions}")
private Integer partitions;
@Value("${spring.kafka.producer.replication_factor}")
private Short replicationFactor;
@Autowired
private KafkaUtil kafkaUtil;
@ -69,9 +72,8 @@ public class EmailJobsController {
newEmailJobVo.setPlatform(platform);
try {
//todo creatTopic
kafkaUtil.sendMessage(GITLINK_EMAIL_TOPIC, JSONObject.toJSONString(newEmailJobVo));
kafkaUtil.createTopic(Arrays.asList(new NewTopic(gitlinkEmailTopic, partitions, replicationFactor)));
kafkaUtil.sendMessage(gitlinkEmailTopic, JSONObject.toJSONString(newEmailJobVo));
return DataPacketUtil.jsonSuccessResult();
} catch (Exception e) {
logger.error(e);

View File

@ -14,6 +14,7 @@ spring:
replication_factor: 1
partitions: 3
topic: topic-gitlink-notification
topic_email: topic-gitlink-email
redis:
database: 0

View File

@ -0,0 +1,10 @@
spring.mail.host=smtp.exmail.qq.com
spring.mail.default-encoding=utf-8
spring.mail.port=465
spring.mail.username=gitlink@barats.cn
spring.mail.password=
spring.mail.properties[mail.smtp.auth]=true
spring.mail.properties[mail.smtp.connectiontimeout]=10000
spring.mail.properties[mail.smtp.timeout]=10000
spring.mail.properties[mail.smtp.writetimeout]=10000
spring.mail.properties[mail.smtp.starttls.enable]=true