From f21369094a1de79ccc3293c3eb497fab6fe75b92 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=B8=87=E4=BD=B3?= Date: Wed, 29 Sep 2021 09:38:58 +0800 Subject: [PATCH 1/2] =?UTF-8?q?kafka=E7=BA=BF=E7=A8=8B=E6=95=B0=E4=BF=AE?= =?UTF-8?q?=E6=94=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- executor/src/main/resources/application.yml.example | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/executor/src/main/resources/application.yml.example b/executor/src/main/resources/application.yml.example index b82de5c..1a4d3df 100644 --- a/executor/src/main/resources/application.yml.example +++ b/executor/src/main/resources/application.yml.example @@ -28,7 +28,7 @@ spring: topic_email: topic-gitlink-email topic_new_email_remind: topic-gitlink-new-email-remind listener: - concurrency: 3 + concurrency: 1 ack-mode: record redis: From 65a0c6a3a1946aed731eb0417b29db1fd7c60804 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=B8=87=E4=BD=B3?= Date: Fri, 1 Oct 2021 00:31:20 +0800 Subject: [PATCH 2/2] =?UTF-8?q?EmailService=E5=A4=84=E7=90=86=E6=96=B9?= =?UTF-8?q?=E5=BC=8F=E5=8F=98=E6=9B=B4=EF=BC=8C=E6=A0=B9=E6=8D=AEemailJobI?= =?UTF-8?q?d=E5=A4=84=E7=90=86=E5=8D=95=E4=B8=AAemailJob?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../executor/service/email/EmailService.java | 12 ++++++------ .../service/jobhandler/EmailJobsListener.java | 5 +++-- .../service/jobhandler/EmailSentListener.java | 4 ++-- .../model/dao/entity/vo/NewEmailJobVo.java | 10 ++++++++++ .../model/dao/mapper/EmailJobsMapper.java | 2 +- .../model/dao/mapper/EmailSendRecordsMapper.java | 2 +- .../model/service/notification/EmailJobsService.java | 8 ++++---- .../notification/EmailSendRecordsService.java | 4 ++-- .../notification/impl/EmailJobsServiceImpl.java | 9 +++++---- .../impl/EmailSendRecordsServiceImpl.java | 4 ++-- .../mapper/notification/EmailJobsMapper.xml | 7 ++----- .../mapper/notification/EmailSendRecordsMapper.xml | 5 +---- 12 files changed, 39 insertions(+), 33 deletions(-) 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 ae533a0..5755bd6 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 @@ -35,17 +35,17 @@ public class EmailService { * 处理邮件发送任务,根据emails添加到邮件发送记录表中 * * @param platform 平台编码 - * @param dispatchNumber 待处理发送任务列表数量 + * @param emailJobId 邮件任务id * @return: void * @Author: wanjia * @Date: 2021/9/15 */ @Transactional - public void DispatchEmailJobs(String platform, Integer dispatchNumber) { + public void DispatchEmailJobs(String platform, Integer emailJobId) { //获取指定数量待处理列表 List emailJobList = new ArrayList<>(); try { - emailJobList = emailJobsService.getEmailJobsByDispatchedStatus(platform, NotificationSystemConstant.EMAIL_JOB_NOT_DISPATCHED, dispatchNumber); + emailJobList = emailJobsService.getEmailJobsByDispatchedStatus(platform, NotificationSystemConstant.EMAIL_JOB_NOT_DISPATCHED, emailJobId); } catch (Exception e) { logger.error("获取未处理邮件任务列表失败:\n" + e); } @@ -71,17 +71,17 @@ public class EmailService { * 发送邮件 * * @param platform 平台编码 - * @param sentNumber 一次发送数量 + * @param emailJobId 邮件任务id * @return: void * @Author: wanjia * @Date: 2021/9/15 */ @Transactional - public void sendEmail(String platform, Integer sentNumber) { + public void sendEmail(String platform, Integer emailJobId) { //获取待发送列表 List emailSendRecordList = new ArrayList<>(); try { - emailSendRecordList = emailSendRecordsService.getRecordsByStatus(platform, NotificationSystemConstant.EMAIL_UNSENT_RECORD, sentNumber); + emailSendRecordList = emailSendRecordsService.getRecordsByStatus(platform, NotificationSystemConstant.EMAIL_UNSENT_RECORD, emailJobId); } catch (Exception e) { logger.error("获取未发送邮件列表失败:\n" + e); } 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 a6795ac..bf9f290 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 @@ -35,9 +35,10 @@ public class EmailJobsListener { public void messageHandler(String message) { try { NewEmailJobVo newEmailJobVo = JSONObject.parseObject(message, NewEmailJobVo.class); - Boolean flag = emailJobsService.createEmailJob(newEmailJobVo); + int flag = emailJobsService.createEmailJob(newEmailJobVo); //if the message is inserted successfully, send a new email-job message to kafka - if (flag){ + if (flag > 0){ + newEmailJobVo.setId(flag); kafkaUtil.sendMessage(gitlinkNewEmailRemindTopic, JSONObject.toJSONString(newEmailJobVo)); } } catch (Exception 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 index f848c24..55a9e76 100644 --- 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 @@ -26,9 +26,9 @@ public class EmailSentListener { NewEmailJobVo newEmailJobVo = JSONObject.parseObject(message, NewEmailJobVo.class); //处理邮件任务 - emailService.DispatchEmailJobs(newEmailJobVo.getPlatform(), 1); + emailService.DispatchEmailJobs(newEmailJobVo.getPlatform(), newEmailJobVo.getId()); //发送邮件 - emailService.sendEmail(newEmailJobVo.getPlatform(), 100); + emailService.sendEmail(newEmailJobVo.getPlatform(), newEmailJobVo.getId()); } diff --git a/model/src/main/java/cn/org/gitlink/notification/model/dao/entity/vo/NewEmailJobVo.java b/model/src/main/java/cn/org/gitlink/notification/model/dao/entity/vo/NewEmailJobVo.java index 9781bfc..702abe4 100644 --- a/model/src/main/java/cn/org/gitlink/notification/model/dao/entity/vo/NewEmailJobVo.java +++ b/model/src/main/java/cn/org/gitlink/notification/model/dao/entity/vo/NewEmailJobVo.java @@ -13,6 +13,8 @@ public class NewEmailJobVo { //平台编码 private String platform; + private Integer id; + @ApiModelProperty(value = "邮件发送者", required = true) @NotNull(message = "邮件发送者不能为空") @Range(min = Integer.MIN_VALUE, max = Integer.MAX_VALUE) @@ -40,6 +42,14 @@ public class NewEmailJobVo { this.platform = platform; } + public Integer getId() { + return id; + } + + public void setId(Integer id) { + this.id = id; + } + public Integer getSender() { return sender; } diff --git a/model/src/main/java/cn/org/gitlink/notification/model/dao/mapper/EmailJobsMapper.java b/model/src/main/java/cn/org/gitlink/notification/model/dao/mapper/EmailJobsMapper.java index 2a60965..18cc815 100644 --- a/model/src/main/java/cn/org/gitlink/notification/model/dao/mapper/EmailJobsMapper.java +++ b/model/src/main/java/cn/org/gitlink/notification/model/dao/mapper/EmailJobsMapper.java @@ -19,5 +19,5 @@ public interface EmailJobsMapper extends BaseMapper { int updateByPrimaryKey(@Param("platform") String platform, @Param("record") EmailJob record); - List getEmailJobsByDispatchedStatus(@Param("platform") String platform, @Param("dispatchedStatus") Integer dispatchedStatus, @Param("size") Integer size); + List getEmailJobsByDispatchedStatus(@Param("platform") String platform, @Param("dispatchedStatus") Integer dispatchedStatus, @Param("id") Integer emailJobId); } diff --git a/model/src/main/java/cn/org/gitlink/notification/model/dao/mapper/EmailSendRecordsMapper.java b/model/src/main/java/cn/org/gitlink/notification/model/dao/mapper/EmailSendRecordsMapper.java index 5f950be..19e6d27 100644 --- a/model/src/main/java/cn/org/gitlink/notification/model/dao/mapper/EmailSendRecordsMapper.java +++ b/model/src/main/java/cn/org/gitlink/notification/model/dao/mapper/EmailSendRecordsMapper.java @@ -22,7 +22,7 @@ public interface EmailSendRecordsMapper extends BaseMapper { //批量插入邮件发送任务记录 int insertEmailSendRecordBatch(@Param("platform") String platform,@Param("list") List emailSendRecordList); - List getRecordsByStatus(@Param("platform") String platform, @Param("status") Integer status, @Param("size") Integer sentNumber); + List getRecordsByStatus(@Param("platform") String platform, @Param("status") Integer status, @Param("jobId") Integer emailJobId); //批量更新邮件发送任务记录 int updateEmailSendRecordsBatch(@Param("platform") String platform, @Param("list") List emailSendRecordList); 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 1ca4abe..b157c16 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 @@ -13,23 +13,23 @@ public interface EmailJobsService extends IService { * 添加发送邮件任务 * * @param newEmailJobVo - * @return: boolean + * @return: int * @Author: wanjia * @Date: 2021/9/13 */ - boolean createEmailJob(NewEmailJobVo newEmailJobVo) throws Exception; + int createEmailJob(NewEmailJobVo newEmailJobVo) throws Exception; /** * 获取所有未处理邮件任务列表 * * @param platform 平台编码 * @param dispatchedStatus 发送状态:-1 未处理,1 处理成功,2 处理失败 - * @param size 列表大小 + * @param emailJobId 邮件任务id * @return: List * @Author: wanjia * @Date: 2021/9/13 */ - List getEmailJobsByDispatchedStatus(String platform,Integer dispatchedStatus, Integer size) throws Exception; + List getEmailJobsByDispatchedStatus(String platform,Integer dispatchedStatus, Integer emailJobId) throws Exception; /** * 更新邮件任务状态 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 0de3039..66bdd4c 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 @@ -23,12 +23,12 @@ public interface EmailSendRecordsService extends IService { * * @param platform 平台编码 * @param status 邮件发送记录状态 -1未发送 1发送成功 2发送失败 - * @param size 列表数量 + * @param emailJobId 邮件任务id * @return: List * @Author: wanjia * @Date: 2021/9/15 */ - List getRecordsByStatus(String platform, Integer status, Integer size) throws Exception; + List getRecordsByStatus(String platform, Integer status, Integer emailJobId) throws Exception; /** * 变更邮件发送记录状态 diff --git a/model/src/main/java/cn/org/gitlink/notification/model/service/notification/impl/EmailJobsServiceImpl.java b/model/src/main/java/cn/org/gitlink/notification/model/service/notification/impl/EmailJobsServiceImpl.java index 95d59d6..6f7ad80 100644 --- a/model/src/main/java/cn/org/gitlink/notification/model/service/notification/impl/EmailJobsServiceImpl.java +++ b/model/src/main/java/cn/org/gitlink/notification/model/service/notification/impl/EmailJobsServiceImpl.java @@ -15,16 +15,17 @@ import java.util.List; public class EmailJobsServiceImpl extends ServiceImpl implements EmailJobsService { @Override - public boolean createEmailJob(NewEmailJobVo newEmailJobVo) { + public int createEmailJob(NewEmailJobVo newEmailJobVo) { String platform = newEmailJobVo.getPlatform(); EmailJob emailJob = new EmailJob(); BeanUtils.copyProperties(newEmailJobVo, emailJob); - return baseMapper.insertSelective(platform, emailJob) > 0; + baseMapper.insertSelective(platform, emailJob); + return newEmailJobVo.getId(); } @Override - public List getEmailJobsByDispatchedStatus(String platform,Integer dispatchedStatus, Integer size) { - return baseMapper.getEmailJobsByDispatchedStatus(platform,dispatchedStatus, size); + public List getEmailJobsByDispatchedStatus(String platform,Integer dispatchedStatus, Integer emailJobId) { + return baseMapper.getEmailJobsByDispatchedStatus(platform,dispatchedStatus, emailJobId); } @Override diff --git a/model/src/main/java/cn/org/gitlink/notification/model/service/notification/impl/EmailSendRecordsServiceImpl.java b/model/src/main/java/cn/org/gitlink/notification/model/service/notification/impl/EmailSendRecordsServiceImpl.java index 4fd254d..8f950c7 100644 --- a/model/src/main/java/cn/org/gitlink/notification/model/service/notification/impl/EmailSendRecordsServiceImpl.java +++ b/model/src/main/java/cn/org/gitlink/notification/model/service/notification/impl/EmailSendRecordsServiceImpl.java @@ -35,7 +35,7 @@ public class EmailSendRecordsServiceImpl extends ServiceImpl getRecordsByStatus(String platform, Integer status, Integer sentNumber){ - return baseMapper.getRecordsByStatus(platform, status, sentNumber); + public List getRecordsByStatus(String platform, Integer status, Integer emailJobId){ + return baseMapper.getRecordsByStatus(platform, status, emailJobId); } } diff --git a/model/src/main/resources/mapper/notification/EmailJobsMapper.xml b/model/src/main/resources/mapper/notification/EmailJobsMapper.xml index 11ce5b3..be90ed3 100644 --- a/model/src/main/resources/mapper/notification/EmailJobsMapper.xml +++ b/model/src/main/resources/mapper/notification/EmailJobsMapper.xml @@ -32,7 +32,7 @@ #{record.subject,jdbcType=VARCHAR}, #{record.content,jdbcType=VARCHAR}, #{record.createdAt,jdbcType=TIMESTAMP}, #{record.dispatchedAt,jdbcType=TIMESTAMP}, #{record.dispatchedStatus,jdbcType=INTEGER}) - + insert into ${platform}_email_jobs @@ -127,9 +127,6 @@ \ 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 b6dfd19..b24ee8f 100644 --- a/model/src/main/resources/mapper/notification/EmailSendRecordsMapper.xml +++ b/model/src/main/resources/mapper/notification/EmailSendRecordsMapper.xml @@ -121,10 +121,7 @@ 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 - - limit #{size} - + where status = #{status} and job_id = #{jobId} update ${platform}_email_send_records