2 Commits ff85e528e0 ... bf886e128d

Autore SHA1 Messaggio Data
  chenkq bf886e128d Merge remote-tracking branch 'origin/master' 1 giorno fa
  chenkq 29ef3f5474 解耦,修改邮件与业务的交互逻辑 1 giorno fa

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

@@ -191,8 +191,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)
@@ -188,7 +175,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;
         }