Ver Fonte

解耦,修改邮件与业务的交互逻辑

chenkq há 1 dia atrás
pai
commit
29ef3f5474

+ 13 - 13
java/storlead-api/src/main/resources/application-test.yml

@@ -21,7 +21,7 @@ spring:
     multipart:
       max-file-size: 20MB
       max-request-size: 20MB
-  ## quartz摰𡁏𧒄隞餃𦛚,��鍂�唳旿摨𤘪䲮?
+  ## quartz摰𡁏𧒄隞餃𦛚,��鍂�唳旿摨𤘪䲮嚙�?
   #  quartz:
   #    job-store-type: jdbc
   #json �園𡢿�喟�銝�頧祆揢
@@ -32,7 +32,7 @@ spring:
     proxy-target-class: true
   #�滨蔭freemarker
   freemarker:
-    # 霈曄蔭璅⊥踎�𡒊��?
+    # 霈曄蔭璅⊥踎�𡒊��?
     suffix: .ftl
     # 霈曄蔭��﹝蝐餃�
     content-type: text/html
@@ -44,7 +44,7 @@ spring:
     # 霈曄蔭ftl��辣頝臬�
     template-loader-path:
       - classpath:/templates
-  # 霈曄蔭�蹱���隞嗉楝敺��js,css? #redis �滨蔭
+  # 霈曄蔭�蹱���隞嗉楝敺��js,css嚙�? #redis �滨蔭
   redis:
     host: test1.storlead.com
     port: 59394
@@ -64,13 +64,13 @@ spring:
     exclude: com.alibaba.druid.spring.boot.autoconfigure.DruidDataSourceAutoConfigure
   datasource:
     dynamic:
-      druid: # �典�druid��㺭嚗𣬚�憭折����澆�暺䁅恕靽脲�銝��氬�?�啣歇�舀�����啣�銝?銝齿�璆𡁜鉄銋劐�閬�僚霈曄蔭)
+      druid: # �典�druid��㺭嚗𣬚�憭折����澆�暺䁅恕靽脲�銝��湛蕭?�啣歇�舀�����啣�嚙�?銝齿�璆𡁜鉄銋劐�閬�僚霈曄蔭)
         # 餈墧𦻖瘙删��滨蔭靽⊥�
-        # �嘥��硋之撠𧶏���撠𧶏���?
+        # �嘥��硋之撠𧶏���撠𧶏���嚙�?
         initial-size: 5
         min-idle: 5
         maxActive: 20
-        # �滨蔭�瑕�餈墧𦻖蝑匧�頞�𧒄��𧒄�?
+        # �滨蔭�瑕�餈墧𦻖蝑匧�頞�𧒄��𧒄�?
         maxWait: 60000
         # �滨蔭�湧�憭帋��滩�銵䔶�甈⊥�瘚页�璉�瘚钅�閬���剔�蝛粹𤦭餈墧𦻖嚗��雿齿糓瘥怎�
         timeBetweenEvictionRunsMillis: 60000
@@ -80,10 +80,10 @@ spring:
         testWhileIdle: true
         testOnBorrow: false
         testOnReturn: false
-        # �枏�PSCache嚗�僎銝娍�摰𡁏�銝芾��乩�PSCache��之?
+        # �枏�PSCache嚗�僎銝娍�摰𡁏�銝芾��乩�PSCache��之嚙�?
         poolPreparedStatements: true
         maxPoolPreparedStatementPerConnectionSize: 20
-        # �滨蔭�烐綉蝏蠘恣�行⏛��ilters嚗�縧�匧��烐綉�屸𢒰sql�䭾�蝏蠘恣嚗?wall'�其��脩�憓?
+        # �滨蔭�烐綉蝏蠘恣�行⏛��ilters嚗�縧�匧��烐綉�屸𢒰sql�䭾�蝏蠘恣嚙�?wall'�其��脩�嚙�?
         filters: stat,wall,slf4j
         # �朞�connectProperties撅墧�扳䔉�枏�mergeSql�蠘�嚗𥟇�SQL霈啣�
         connectionProperties: druid.stat.mergeSql\=true;druid.stat.slowSqlMillis\=5000
@@ -127,7 +127,7 @@ spring:
 #mybatis plus 霈曄蔭
 mybatis-plus:
   mapper-locations: classpath*:/mapper/*Mapper.xml,classpath*:/mapper/*/*Mapper.xml
-  # 摰硺��急�嚗��銝?package �券�堒噡�𤥁����瑕��?
+  # 摰硺��急�嚗��嚙�?package �券�堒噡�𤥁����瑕�嚙�?
   type-aliases-package: com.storlead.tems.modules.*.entity
   type-enums-package:
     #  configuration:
@@ -138,9 +138,9 @@ mybatis-plus:
     # �喲𡡒MP3.0�芸蒂��anner
     banner: false
     db-config:
-      #銝駁睸蝐餃�  0:"�唳旿摨𨧻D�芸�",1:"霂亦掩�衤蛹�芾挽蝵桐蜓�桃掩�?, 2:"�冽�颲枏�ID",3:"�典��臭�ID (�啣�蝐餃��臭�ID)", 4:"�典��臭�ID UUID",5:"摮㛖泵銝脣�撅��臭�ID (idWorker ���蝚虫葡銵函內)";
+      #銝駁睸蝐餃�  0:"�唳旿摨𨧻D�芸�",1:"霂亦掩�衤蛹�芾挽蝵桐蜓�桃掩�?, 2:"�冽�颲枏�ID",3:"�典��臭�ID (�啣�蝐餃��臭�ID)", 4:"�典��臭�ID UUID",5:"摮㛖泵銝脣�撅��臭�ID (idWorker ���蝚虫葡銵函內)";
       id-type: 4
-      # 暺䁅恕�唳旿摨栞”銝见�蝥踹𦶢�?
+      # 暺䁅恕�唳旿摨栞”銝见�蝥踹𦶢�?
       table-underline: true
     #configuration:
     # 餈嗘葵�滨蔭隡𡁜��扯���ql�枏㫲�箸䔉嚗�銁撘��烐�瘚贝���𧒄�坔虾隞亦鍂
@@ -190,8 +190,8 @@ storlead:
   mail:
     remote:
       enabled: true
-      scheme: http
-      host: 127.0.0.1
+      scheme: https
+      host: email.test.storlead.com
       port: 18090
       context-path: /router/rest
       token-header: token

+ 10 - 0
java/storlead-mail/storlead-mail-biz/src/main/java/com/storlead/mail/mapper/MktMailStatusMapper.java

@@ -0,0 +1,10 @@
+package com.storlead.mail.mapper;
+
+import com.storlead.mail.mapper.support.MailBaseMapper;
+import com.storlead.mail.pojo.entity.MktMailStatusEntity;
+
+/**
+ * 营销邮件状态 Mapper(邮件库)。
+ */
+public interface MktMailStatusMapper extends MailBaseMapper<MktMailStatusEntity> {
+}

+ 50 - 0
java/storlead-mail/storlead-mail-biz/src/main/java/com/storlead/mail/service/impl/MktMailStatusServiceImpl.java

@@ -0,0 +1,50 @@
+package com.storlead.mail.service.impl;
+
+import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
+import com.storlead.framework.common.constant.CommonConstant;
+import com.storlead.mail.mapper.MktMailStatusMapper;
+import com.storlead.mail.pojo.entity.MktMailStatusEntity;
+import com.storlead.mail.service.MktMailStatusService;
+import com.storlead.mail.service.impl.support.MailDataSourceServiceImpl;
+import org.springframework.stereotype.Service;
+import org.springframework.util.CollectionUtils;
+
+import java.util.Collections;
+import java.util.Date;
+import java.util.List;
+
+/**
+ * 邮件库 mkt_mail_status 读写(@DS mail)。
+ */
+@Service
+public class MktMailStatusServiceImpl
+        extends MailDataSourceServiceImpl<MktMailStatusMapper, MktMailStatusEntity>
+        implements MktMailStatusService {
+
+    @Override
+    public List<MktMailStatusEntity> listPendingChanges(Long lastId, int limit) {
+        if (limit <= 0) {
+            return Collections.emptyList();
+        }
+        LambdaQueryWrapper<MktMailStatusEntity> wrapper = new LambdaQueryWrapper<>();
+        wrapper.eq(MktMailStatusEntity::getIsDelete, CommonConstant.DEL_FLAG_0)
+                .apply("update_time > status_updated_at")
+                .gt(lastId != null, MktMailStatusEntity::getId, lastId)
+                .orderByAsc(MktMailStatusEntity::getId)
+                .last("limit " + limit);
+        return list(wrapper);
+    }
+
+    @Override
+    public boolean ackStatusSync(List<Long> mailIds) {
+        if (CollectionUtils.isEmpty(mailIds)) {
+            return true;
+        }
+        LambdaUpdateWrapper<MktMailStatusEntity> update = new LambdaUpdateWrapper<>();
+        update.in(MktMailStatusEntity::getMailId, mailIds)
+                .eq(MktMailStatusEntity::getIsDelete, CommonConstant.DEL_FLAG_0)
+                .set(MktMailStatusEntity::getStatusUpdatedAt, new Date());
+        return update(update);
+    }
+}

+ 73 - 0
java/storlead-mail/storlead-mail-core/src/main/java/com/storlead/mail/pojo/entity/MktMailStatusEntity.java

@@ -0,0 +1,73 @@
+package com.storlead.mail.pojo.entity;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableField;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import com.storlead.framework.mybatis.entity.SysBaseField;
+import io.swagger.annotations.ApiModel;
+import io.swagger.annotations.ApiModelProperty;
+import lombok.Data;
+import lombok.EqualsAndHashCode;
+import lombok.experimental.Accessors;
+import org.springframework.format.annotation.DateTimeFormat;
+import com.fasterxml.jackson.annotation.JsonFormat;
+
+import java.util.Date;
+
+/**
+ * 营销邮件状态(邮件库 mkt_mail_status,供业务库增量同步)。
+ */
+@Data
+@EqualsAndHashCode(callSuper = true)
+@Accessors(chain = true)
+@TableName("mkt_mail_status")
+@ApiModel(value = "MktMailStatusEntity", description = "营销邮件状态")
+public class MktMailStatusEntity extends SysBaseField {
+
+    private static final long serialVersionUID = 1L;
+
+    @TableId(value = "id", type = IdType.AUTO)
+    private Long id;
+
+    @ApiModelProperty("邮件ID,对应 mkt_mail_message.id / market_emails.email_id")
+    @TableField("mail_id")
+    private Long mailId;
+
+    @ApiModelProperty("发送状态:-2草稿 -1定时 0发送中 1成功 2失败")
+    @TableField("send_status")
+    private Integer sendStatus;
+
+    @ApiModelProperty("是否已打开(像素):0否 1是")
+    @TableField("opened")
+    private Integer opened;
+
+    @JsonFormat(timezone = "GMT+8", pattern = "yyyy-MM-dd HH:mm:ss")
+    @DateTimeFormat(pattern = "yyyy-MM-dd HH:mm:ss")
+    @TableField("first_open_time")
+    private Date firstOpenTime;
+
+    @ApiModelProperty("是否已读:0否 1是")
+    @TableField("is_read")
+    private Integer isRead;
+
+    @JsonFormat(timezone = "GMT+8", pattern = "yyyy-MM-dd HH:mm:ss")
+    @DateTimeFormat(pattern = "yyyy-MM-dd HH:mm:ss")
+    @TableField("first_read_time")
+    private Date firstReadTime;
+
+    @ApiModelProperty("是否已回复:0否 1是")
+    @TableField("replied")
+    private Integer replied;
+
+    @JsonFormat(timezone = "GMT+8", pattern = "yyyy-MM-dd HH:mm:ss")
+    @DateTimeFormat(pattern = "yyyy-MM-dd HH:mm:ss")
+    @TableField("first_reply_time")
+    private Date firstReplyTime;
+
+    @ApiModelProperty("状态最后更新时间(第三方增量游标)")
+    @JsonFormat(timezone = "GMT+8", pattern = "yyyy-MM-dd HH:mm:ss")
+    @DateTimeFormat(pattern = "yyyy-MM-dd HH:mm:ss")
+    @TableField("status_updated_at")
+    private Date statusUpdatedAt;
+}

+ 25 - 0
java/storlead-mail/storlead-mail-spi/src/main/java/com/storlead/mail/service/MktMailStatusService.java

@@ -0,0 +1,25 @@
+package com.storlead.mail.service;
+
+import com.storlead.framework.mybatis.service.MyBaseService;
+import com.storlead.mail.pojo.entity.MktMailStatusEntity;
+
+import java.util.List;
+
+/**
+ * 营销邮件状态(邮件库 mkt_mail_status)。
+ */
+public interface MktMailStatusService extends MyBaseService<MktMailStatusEntity> {
+
+    /**
+     * 待同步:update_time &gt; status_updated_at,按 id 升序分页。
+     *
+     * @param lastId 上一批最后一条 id,首页传 null
+     * @param limit  批量大小
+     */
+    List<MktMailStatusEntity> listPendingChanges(Long lastId, int limit);
+
+    /**
+     * 同步成功后回写 status_updated_at = now()。
+     */
+    boolean ackStatusSync(List<Long> mailIds);
+}

+ 5 - 0
java/storlead-sasa/storlead-trade/pom.xml

@@ -72,6 +72,11 @@
             <groupId>com.storlead.boot</groupId>
             <artifactId>storlead-mail-spi</artifactId>
         </dependency>
+        <!-- 跨库读 mkt_mail_status:Mapper/Service 实现在 mail-biz(@DS mail) -->
+        <dependency>
+            <groupId>com.storlead.boot</groupId>
+            <artifactId>storlead-mail-biz</artifactId>
+        </dependency>
         <dependency>
             <groupId>com.storlead.boot</groupId>
             <artifactId>storlead-knowledge-core</artifactId>

+ 17 - 30
java/storlead-sasa/storlead-trade/src/main/java/com/storlead/trade/service/impl/MarketEmailsStatusSyncServiceImpl.java

@@ -5,10 +5,8 @@ import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
 import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
 import com.storlead.framework.common.constant.CommonConstant;
 import com.storlead.framework.common.constant.DSConstants;
-import com.storlead.framework.common.result.Result;
-import com.storlead.mail.webclient.client.MktMailRemoteClient;
-import com.storlead.mail.webclient.dto.MktMailStatusChangeRemoteDTO;
-import com.storlead.mail.webclient.dto.MktMailStatusSyncQueryRemoteDTO;
+import com.storlead.mail.pojo.entity.MktMailStatusEntity;
+import com.storlead.mail.service.MktMailStatusService;
 import com.storlead.trade.entity.MarketEmailsEntity;
 import com.storlead.trade.service.MarketEmailsService;
 import com.storlead.trade.service.MarketEmailsStatusSyncService;
@@ -21,8 +19,9 @@ import java.util.ArrayList;
 import java.util.List;
 
 /**
- * mkt_mail_status → market_emails:
- * 查待同步(update_time &gt; status_updated_at)→ 按 mail_id 改业务表 → 回写同步时间。
+ * 跨库同步:邮件库 mkt_mail_status → 业务库 market_emails,再回写 status_updated_at。
+ * <p>
+ * 邮件侧走 {@link MktMailStatusService}(@DS mail),业务侧走 trade 数据源;不做跨库事务。
  */
 @Slf4j
 @Service
@@ -37,7 +36,7 @@ public class MarketEmailsStatusSyncServiceImpl implements MarketEmailsStatusSync
     private static final int TRADE_STATUS_FAIL = 2;
 
     @Resource
-    private MktMailRemoteClient mktMailRemoteClient;
+    private MktMailStatusService mktMailStatusService;
 
     @Resource
     private MarketEmailsService marketEmailsService;
@@ -48,56 +47,44 @@ public class MarketEmailsStatusSyncServiceImpl implements MarketEmailsStatusSync
         Long lastId = null;
 
         for (int page = 0; page < MAX_PAGES; page++) {
-            // 1. 查 mkt_mail_status:有最新修改(update_time > status_updated_at)
-            MktMailStatusSyncQueryRemoteDTO query = new MktMailStatusSyncQueryRemoteDTO();
-            query.setLimit(BATCH_LIMIT);
-            query.setLastId(lastId);
-
-            Result<List<MktMailStatusChangeRemoteDTO>> result;
+            // 1. 邮件库:update_time > status_updated_at
+            List<MktMailStatusEntity> changes;
             try {
-                result = mktMailRemoteClient.statusChanges(query);
+                changes = mktMailStatusService.listPendingChanges(lastId, BATCH_LIMIT);
             } catch (Exception e) {
                 log.error("查询 mkt_mail_status 待同步数据失败 lastId={}", lastId, e);
                 break;
             }
-            if (result == null || !result.isSuccess()) {
-                log.warn("查询 mkt_mail_status 失败 message={}", result == null ? null : result.getMessage());
-                break;
-            }
-            List<MktMailStatusChangeRemoteDTO> changes = result.getResult();
             if (CollectionUtils.isEmpty(changes)) {
                 break;
             }
 
             List<Long> syncedMailIds = new ArrayList<>();
             Long batchLastId = lastId;
-            for (MktMailStatusChangeRemoteDTO change : changes) {
+            for (MktMailStatusEntity change : changes) {
                 if (change == null || change.getMailId() == null) {
                     continue;
                 }
                 if (change.getId() != null) {
                     batchLastId = change.getId();
                 }
-                // 2. 按邮件 id 修改 market_emails
+                // 2. 业务库:按 email_id = mail_id 更新 market_emails
                 ApplyResult applyResult = applyToMarketEmails(change);
                 if (applyResult == ApplyResult.UPDATED) {
                     totalUpdated++;
                     syncedMailIds.add(change.getMailId());
                 } else if (applyResult == ApplyResult.MATCHED_NO_CHANGE) {
-                    // 业务表已对齐,仍需回写同步时间,避免反复拉取
                     syncedMailIds.add(change.getMailId());
                 }
                 // UNMATCHED:尚无 email_id 关联,不 ack,下次继续等
             }
 
-            // 3. 同步后回写 mkt_mail_status.status_updated_at
+            // 3. 邮件库:回写 status_updated_at
             if (!syncedMailIds.isEmpty()) {
                 try {
-                    Result<?> ack = mktMailRemoteClient.ackStatusSync(syncedMailIds);
-                    if (ack == null || !ack.isSuccess()) {
-                        log.warn("回写 mkt_mail_status.status_updated_at 失败 mailIds={} message={}",
-                                syncedMailIds, ack == null ? null : ack.getMessage());
-                        // ack 失败则不再继续翻页,避免漏确认
+                    boolean ackOk = mktMailStatusService.ackStatusSync(syncedMailIds);
+                    if (!ackOk) {
+                        log.warn("回写 mkt_mail_status.status_updated_at 失败 mailIds={}", syncedMailIds);
                         break;
                     }
                 } catch (Exception e) {
@@ -117,7 +104,7 @@ public class MarketEmailsStatusSyncServiceImpl implements MarketEmailsStatusSync
         return totalUpdated;
     }
 
-    private ApplyResult applyToMarketEmails(MktMailStatusChangeRemoteDTO change) {
+    private ApplyResult applyToMarketEmails(MktMailStatusEntity change) {
         MarketEmailsEntity existing = marketEmailsService.getOne(new LambdaQueryWrapper<MarketEmailsEntity>()
                 .eq(MarketEmailsEntity::getEmailId, change.getMailId())
                 .eq(MarketEmailsEntity::getIsDelete, CommonConstant.DEL_FLAG_0)
@@ -168,7 +155,7 @@ public class MarketEmailsStatusSyncServiceImpl implements MarketEmailsStatusSync
         return TRADE_STATUS_DRAFT;
     }
 
-    private static Integer resolveReader(MktMailStatusChangeRemoteDTO change) {
+    private static Integer resolveReader(MktMailStatusEntity change) {
         if (Integer.valueOf(1).equals(change.getOpened()) || Integer.valueOf(1).equals(change.getIsRead())) {
             return 1;
         }