From 75754b5246fffb96f157b3d6f38cea433837fd10 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=B8=87=E4=BD=B3?= Date: Fri, 17 Sep 2021 10:52:07 +0800 Subject: [PATCH] =?UTF-8?q?email=E5=88=86=E9=85=8D=E5=8F=91=E9=80=81?= =?UTF-8?q?=E4=BB=BB=E5=8A=A1=EF=BC=8Cemail=E5=8F=91=E9=80=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../notification/common/utils/EmailUtils.java | 5 +++ db/gns-email.sql | 26 ++++++++++++ .../executor/service/email/EmailService.java | 40 +++++++++++++------ .../service/jobhandler/EmailJobsListener.java | 27 +++++++++++-- .../service/jobhandler/EmailSentListener.java | 39 ++++++++++++++++++ .../main/resources/application.yml.example | 7 ++++ executor/src/main/resources/mail.properties | 8 ++-- middleware/services.yml | 1 + .../model/dao/entity/EmailSendRecord.java | 23 +++++++++++ .../notification/EmailJobsService.java | 22 +++++----- .../notification/EmailSendRecordsService.java | 22 +++++----- .../mapper/notification/EmailJobsMapper.xml | 5 ++- .../notification/EmailSendRecordsMapper.xml | 29 ++++++++------ reader/src/main/resources/mail.properties | 10 +++++ .../controller/EmailJobsController.java | 26 ++++++------ .../main/resources/application.yml.example | 1 + writer/src/main/resources/mail.properties | 10 +++++ 17 files changed, 236 insertions(+), 65 deletions(-) create mode 100644 db/gns-email.sql create mode 100644 executor/src/main/java/cn/org/gitlink/notification/executor/service/jobhandler/EmailSentListener.java create mode 100644 reader/src/main/resources/mail.properties create mode 100644 writer/src/main/resources/mail.properties diff --git a/common/src/main/java/cn/org/gitlink/notification/common/utils/EmailUtils.java b/common/src/main/java/cn/org/gitlink/notification/common/utils/EmailUtils.java index 3ff3a1a..2525f45 100644 --- a/common/src/main/java/cn/org/gitlink/notification/common/utils/EmailUtils.java +++ b/common/src/main/java/cn/org/gitlink/notification/common/utils/EmailUtils.java @@ -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); diff --git a/db/gns-email.sql b/db/gns-email.sql new file mode 100644 index 0000000..6b9b8ae --- /dev/null +++ b/db/gns-email.sql @@ -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; \ No newline at end of file diff --git a/executor/src/main/java/cn/org/gitlink/notification/executor/service/email/EmailService.java b/executor/src/main/java/cn/org/gitlink/notification/executor/service/email/EmailService.java index eb53126..8306fe8 100644 --- a/executor/src/main/java/cn/org/gitlink/notification/executor/service/email/EmailService.java +++ b/executor/src/main/java/cn/org/gitlink/notification/executor/service/email/EmailService.java @@ -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 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 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 { diff --git a/executor/src/main/java/cn/org/gitlink/notification/executor/service/jobhandler/EmailJobsListener.java b/executor/src/main/java/cn/org/gitlink/notification/executor/service/jobhandler/EmailJobsListener.java index c9eda84..85c8489 100644 --- a/executor/src/main/java/cn/org/gitlink/notification/executor/service/jobhandler/EmailJobsListener.java +++ b/executor/src/main/java/cn/org/gitlink/notification/executor/service/jobhandler/EmailJobsListener.java @@ -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); } diff --git a/executor/src/main/java/cn/org/gitlink/notification/executor/service/jobhandler/EmailSentListener.java b/executor/src/main/java/cn/org/gitlink/notification/executor/service/jobhandler/EmailSentListener.java new file mode 100644 index 0000000..f848c24 --- /dev/null +++ b/executor/src/main/java/cn/org/gitlink/notification/executor/service/jobhandler/EmailSentListener.java @@ -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 +} diff --git a/executor/src/main/resources/application.yml.example b/executor/src/main/resources/application.yml.example index afc99a2..635550e 100644 --- a/executor/src/main/resources/application.yml.example +++ b/executor/src/main/resources/application.yml.example @@ -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 diff --git a/executor/src/main/resources/mail.properties b/executor/src/main/resources/mail.properties index ba4b2cd..497e2a3 100644 --- a/executor/src/main/resources/mail.properties +++ b/executor/src/main/resources/mail.properties @@ -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 \ No newline at end of file +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 diff --git a/middleware/services.yml b/middleware/services.yml index 9c2bf38..700fa04 100644 --- a/middleware/services.yml +++ b/middleware/services.yml @@ -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 diff --git a/model/src/main/java/cn/org/gitlink/notification/model/dao/entity/EmailSendRecord.java b/model/src/main/java/cn/org/gitlink/notification/model/dao/entity/EmailSendRecord.java index 8363f19..337f18e 100644 --- a/model/src/main/java/cn/org/gitlink/notification/model/dao/entity/EmailSendRecord.java +++ b/model/src/main/java/cn/org/gitlink/notification/model/dao/entity/EmailSendRecord.java @@ -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; + } } diff --git a/model/src/main/java/cn/org/gitlink/notification/model/service/notification/EmailJobsService.java b/model/src/main/java/cn/org/gitlink/notification/model/service/notification/EmailJobsService.java index a984385..e1585e0 100644 --- a/model/src/main/java/cn/org/gitlink/notification/model/service/notification/EmailJobsService.java +++ b/model/src/main/java/cn/org/gitlink/notification/model/service/notification/EmailJobsService.java @@ -10,9 +10,9 @@ import java.util.List; public interface EmailJobsService extends IService { /** - * @Description: 添加发送邮件任务 + * 添加发送邮件任务 * - * @Param newEmailJobVo + * @param newEmailJobVo * @return: boolean * @Author: wanjia * @Date: 2021/9/13 @@ -20,11 +20,11 @@ public interface EmailJobsService extends IService { 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 * @Author: wanjia * @Date: 2021/9/13 @@ -32,12 +32,12 @@ public interface EmailJobsService extends IService { List 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 diff --git a/model/src/main/java/cn/org/gitlink/notification/model/service/notification/EmailSendRecordsService.java b/model/src/main/java/cn/org/gitlink/notification/model/service/notification/EmailSendRecordsService.java index b82eda5..daabe07 100644 --- a/model/src/main/java/cn/org/gitlink/notification/model/service/notification/EmailSendRecordsService.java +++ b/model/src/main/java/cn/org/gitlink/notification/model/service/notification/EmailSendRecordsService.java @@ -8,11 +8,11 @@ import java.util.List; public interface EmailSendRecordsService extends IService { /** - * @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 { 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 * @Author: wanjia * @Date: 2021/9/15 @@ -32,10 +32,10 @@ public interface EmailSendRecordsService extends IService { List getRecordsByStatus(String platform, Integer status, Integer size) throws Exception; /** - * @Description: 变更邮件发送记录状态 + * 变更邮件发送记录状态 * - * @Param platform 平台编码 - * @Param List 待更新列表 + * @param platform 平台编码 + * @param emailSendRecordList 待更新列表 * @return: int * @Author: wanjia * @Date: 2021/9/13 diff --git a/model/src/main/resources/mapper/notification/EmailJobsMapper.xml b/model/src/main/resources/mapper/notification/EmailJobsMapper.xml index 4f70cf1..11ce5b3 100644 --- a/model/src/main/resources/mapper/notification/EmailJobsMapper.xml +++ b/model/src/main/resources/mapper/notification/EmailJobsMapper.xml @@ -127,6 +127,9 @@ \ No newline at end of file diff --git a/model/src/main/resources/mapper/notification/EmailSendRecordsMapper.xml b/model/src/main/resources/mapper/notification/EmailSendRecordsMapper.xml index 9200882..7b13f60 100644 --- a/model/src/main/resources/mapper/notification/EmailSendRecordsMapper.xml +++ b/model/src/main/resources/mapper/notification/EmailSendRecordsMapper.xml @@ -118,32 +118,37 @@ update ${platform}_email_send_records - - - - when id=#{i.id} then #{i.sentAt} + + + + when id=#{item.id} then #{item.sentAt} - - - when id=#{i.id} then #{i.status} + + + when id=#{item.id} then #{item.status} - where - - id=#{i.id} + where id in + + #{item.id} \ No newline at end of file diff --git a/reader/src/main/resources/mail.properties b/reader/src/main/resources/mail.properties new file mode 100644 index 0000000..4e500ea --- /dev/null +++ b/reader/src/main/resources/mail.properties @@ -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 \ No newline at end of file diff --git a/writer/src/main/java/cn/org/gitlink/notification/writer/controller/EmailJobsController.java b/writer/src/main/java/cn/org/gitlink/notification/writer/controller/EmailJobsController.java index 1c8a7c0..0bed2ad 100644 --- a/writer/src/main/java/cn/org/gitlink/notification/writer/controller/EmailJobsController.java +++ b/writer/src/main/java/cn/org/gitlink/notification/writer/controller/EmailJobsController.java @@ -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); diff --git a/writer/src/main/resources/application.yml.example b/writer/src/main/resources/application.yml.example index 5009c43..d72ad0e 100644 --- a/writer/src/main/resources/application.yml.example +++ b/writer/src/main/resources/application.yml.example @@ -14,6 +14,7 @@ spring: replication_factor: 1 partitions: 3 topic: topic-gitlink-notification + topic_email: topic-gitlink-email redis: database: 0 diff --git a/writer/src/main/resources/mail.properties b/writer/src/main/resources/mail.properties new file mode 100644 index 0000000..4e500ea --- /dev/null +++ b/writer/src/main/resources/mail.properties @@ -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 \ No newline at end of file