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/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: 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