Browse Source

完课福利重复发送修复

xw 3 days ago
parent
commit
c660f0190b

+ 29 - 19
fs-qw-task/src/main/java/com/fs/app/taskService/impl/SopLogsTaskServiceImpl.java

@@ -2691,7 +2691,6 @@ public class SopLogsTaskServiceImpl implements SopLogsTaskService {
 
     // 处理单个批次的方法
     private void processBatch(List<FsCourseWatchLog> batch) {
-        List<FsCourseWatchLog> finishLogsToUpdate = new ArrayList<>();
         List<QwSopLogs> sopLogsToInsert = new ArrayList<>();
         log.info("开始执行处理批次方法-数量:{}",batch.size());
         for (FsCourseWatchLog finishLog : batch) {
@@ -2715,12 +2714,9 @@ public class SopLogsTaskServiceImpl implements SopLogsTaskService {
                 // 查询完课模板信息
                 FsCourseFinishTemp finishTemp = fsCourseFinishTempMapper.selectFsCourseFinishTempByCompanyId(finishLog.getCompanyUserId(),finishLog.getCompanyId(), finishLog.getVideoId());
 
-                // 设置 finishLog 为已发送状态,并加入批量更新列表
-                finishLog.setSendFinishMsg(1);
-                finishLogsToUpdate.add(finishLog);
-
                 if (finishTemp == null) {
 //                    log.error("完课模板不存在: " + finishLog.getQwUserId() + ", " + finishLog.getVideoId());
+                    markSendFinishMsgIfPending(finishLog.getLogId());
                     continue;
                 }
 
@@ -2728,15 +2724,22 @@ public class SopLogsTaskServiceImpl implements SopLogsTaskService {
                 QwSopLogs sopLogs = buildSopLogs(finishLog, externalContact, finishTemp);
                 if (sopLogs == null) {
                     log.error("生成完课发送记录为空-:{}", finishLog.getQwExternalContactId());
+                    markSendFinishMsgIfPending(finishLog.getLogId());
                     continue;
                 }
 
-                // 如果客户状态有效,则加入批量插入列表
-                if (isValidExternalContact(externalContact)) {
-                    sopLogsToInsert.add(sopLogs);
-                } else {
+                if (!isValidExternalContact(externalContact)) {
                     log.info("完课消息-客户信息有误,不生成完课消息: {}", finishLog.getQwExternalContactId());
+                    markSendFinishMsgIfPending(finishLog.getLogId());
+                    continue;
+                }
+
+                if (!tryMarkSendFinishMsg(finishLog.getLogId())) {
+                    log.info("完课记录已被其他任务处理,跳过: logId={}", finishLog.getLogId());
+                    continue;
                 }
+
+                sopLogsToInsert.add(sopLogs);
 //                try {
 //                    fsUserCompanyBindService.finish(externalContact.getFsUserId(), externalContact.getQwUserId(), externalContact.getCompanyUserId(), finishLog);
 //                }catch (Exception e){
@@ -2747,16 +2750,6 @@ public class SopLogsTaskServiceImpl implements SopLogsTaskService {
             }
         }
 
-        // 批量更新和插入
-        if (!finishLogsToUpdate.isEmpty()) {
-            try {
-                fsCourseWatchLogMapper.batchUpdateWatchLogSendMsg(finishLogsToUpdate);
-                log.info("批量更新 finishLog 成功,数量: {}", finishLogsToUpdate.size());
-            } catch (Exception e) {
-                log.error("批量更新 finishLog 失败", e);
-            }
-        }
-
         if (!sopLogsToInsert.isEmpty()) {
             try {
                 qwSopLogsService.batchInsertQwSopLogs(sopLogsToInsert);
@@ -2768,6 +2761,23 @@ public class SopLogsTaskServiceImpl implements SopLogsTaskService {
         log.info("结束处理批次方法-数量:{}",batch.size());
     }
 
+    /**
+     * 乐观锁标记完课消息已处理,返回是否成功抢占
+     */
+    private boolean tryMarkSendFinishMsg(Long logId) {
+        if (logId == null) {
+            return false;
+        }
+        return fsCourseWatchLogMapper.casUpdateSendFinishMsg(logId) > 0;
+    }
+
+    /**
+     * 无需再发送时,避免定时任务重复捞取同一条看课记录
+     */
+    private void markSendFinishMsgIfPending(Long logId) {
+        tryMarkSendFinishMsg(logId);
+    }
+
     /**
      * 构建 QwSopLogs 对象
      */

+ 6 - 0
fs-service/src/main/java/com/fs/course/mapper/FsCourseWatchLogMapper.java

@@ -524,6 +524,12 @@ public interface FsCourseWatchLogMapper extends BaseMapper<FsCourseWatchLog> {
 
     void batchUpdateWatchLogSendMsg(@Param("list") List<FsCourseWatchLog> logs);
 
+    /**
+     * 乐观锁:仅当尚未发送完课消息时标记为已发送,返回影响行数(0 表示已被其他任务处理)
+     */
+    @Update("UPDATE fs_course_watch_log SET send_finish_msg = 1 WHERE log_id = #{logId} AND (send_finish_msg = 0 OR send_finish_msg IS NULL)")
+    int casUpdateSendFinishMsg(@Param("logId") Long logId);
+
     @Select("SELECT l.qw_external_contact_id,l.log_type,l.qw_user_id,l.create_time ,u.first_time,u.create_time lineTime FROM fs_course_watch_log  l " +
             "LEFT JOIN qw_external_contact u ON u.id=l.qw_external_contact_id  " +
             " WHERE  DATE(l.create_time) = DATE_SUB(CURDATE(), INTERVAL 1 DAY) and l.video_id =#{videoId}")