diff --git a/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/corpus/OdsVdqwMessageOTD.java b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/corpus/OdsVdqwMessageOTD.java index eabf97a..c364373 100644 --- a/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/corpus/OdsVdqwMessageOTD.java +++ b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/corpus/OdsVdqwMessageOTD.java @@ -1,10 +1,5 @@ package com.volvo.ai.analytic.center.dto.corpus; -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.volvo.common.core.base.BaseEntity; import lombok.Data; import java.util.Date; @@ -27,6 +22,7 @@ public class OdsVdqwMessageOTD { private String fromUserUnId; private String acceptUserUnId; private String content; + private String fromAcceptUserId; public OdsVdqwMessageOTD() { } diff --git a/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/entity/TtVdqwRecord.java b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/entity/TtVdqwRecord.java new file mode 100644 index 0000000..3158a19 --- /dev/null +++ b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/entity/TtVdqwRecord.java @@ -0,0 +1,64 @@ +package com.volvo.ai.analytic.center.entity; + +import com.baomidou.mybatisplus.annotation.*; +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +import java.util.Date; + +@Data +@Builder +@AllArgsConstructor +@NoArgsConstructor +@TableName("tt_vdqw_record") +public class TtVdqwRecord { + + @TableId(value = "id", type = IdType.AUTO) + private Long id; + + @TableField("from_accept_user_id") + private String from_accept_user_id; + + @TableField("from_user_id") + private String fromUserId; + + @TableField("accept_user_id") + private String acceptUserId; // JSON 字符串 + + @TableField(value="msg_time") + private Date msgTime; + + @TableField("is_deleted") + @TableLogic + private Integer isDeleted; + + @TableField("versions") + @Version + private Integer versions; + + /** + * 创建者 + */ + @TableField("create_by") + private String createBy; + + /** + * 创建时间 + */ + @TableField("create_time") + private Date createTime; + + /** + * 更新者 + */ + @TableField("update_by") + private String updateBy; + + /** + * 更新时间 + */ + @TableField("update_time") + private Date updateTime; +} diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/QiWeiCorpusJob.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/QiWeiCorpusJob.java index 33f20b4..b70c2a5 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/QiWeiCorpusJob.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/QiWeiCorpusJob.java @@ -5,12 +5,17 @@ import com.volvo.common.core.util.ResultMsg; import com.xxl.job.core.context.XxlJobHelper; import com.xxl.job.core.handler.annotation.XxlJob; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.StringUtils; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; +import org.springframework.web.bind.annotation.PostMapping; +import org.springframework.web.bind.annotation.RequestBody; +import org.springframework.web.bind.annotation.RestController; @Slf4j @Component +@RestController public class QiWeiCorpusJob { @Autowired @@ -20,10 +25,14 @@ public class QiWeiCorpusJob { * 企微语料处理 */ @XxlJob("QiWeiCorpusTask") - public ResultMsg qiWeiCorpusTask(String paramJson) { + @PostMapping("QiWeiCorpusTask") + public ResultMsg qiWeiCorpusTask(@RequestBody String paramJson) { try { // 获取任务参数 String param = XxlJobHelper.getJobParam(); + if(StringUtils.isEmpty(param)){ + param = paramJson; + } // 执行业务逻辑 XxlJobHelper.log("任务参数: {}", param); log.info("qiWeiCorpusTask 企微语料查询处理:{}",param); diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/TmOdsVdqwMessagearchivingMapper.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/TmOdsVdqwMessagearchivingMapper.java index dcacc2a..36c8f89 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/TmOdsVdqwMessagearchivingMapper.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/TmOdsVdqwMessagearchivingMapper.java @@ -16,8 +16,8 @@ import java.util.List; @Mapper public interface TmOdsVdqwMessagearchivingMapper extends BaseMapper { - int countOdsVdqwMessageByData(@Param("statTime") String statTime, @Param("endTime") String endTime); - List queryOdsVdqwMessageByData(@Param("statTime") String statTime, @Param("endTime") String endTime, @Param("offset") int offset, @Param("pageSize") int pageSize); + int countOdsVdqwMessageByData(@Param("statTime") String statTime, @Param("endTime") String endTime, @Param("retry") boolean retry); + List queryOdsVdqwMessageByData(@Param("statTime") String statTime, @Param("endTime") String endTime, @Param("offset") int offset, @Param("pageSize") int pageSize, @Param("retry") boolean retry); List queryOdsVdqwMessageByFromUserIdAndAcceptUserId(@Param("statTime") String statTime, @Param("endTime") String endTime, @Param("userIdList") List userIdList); diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/TtVdqwRecordMapper.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/TtVdqwRecordMapper.java new file mode 100644 index 0000000..32453ca --- /dev/null +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/TtVdqwRecordMapper.java @@ -0,0 +1,18 @@ +package com.volvo.ai.analytic.center.mapper; + +import com.baomidou.mybatisplus.core.mapper.BaseMapper; +import com.volvo.ai.analytic.center.entity.TtVdqwRecord; +import org.apache.ibatis.annotations.Mapper; +import org.springframework.stereotype.Repository; + +/** + * @description 电话语料表-同步表 + * @author BEJSON + * @date 2025-03-04 + */ +@Mapper +@Repository +public interface TtVdqwRecordMapper extends BaseMapper { + + +} \ No newline at end of file diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmOdsVdqwMessagearchivingServiceImpl.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmOdsVdqwMessagearchivingServiceImpl.java index 1f469c4..8280860 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmOdsVdqwMessagearchivingServiceImpl.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmOdsVdqwMessagearchivingServiceImpl.java @@ -18,6 +18,7 @@ import com.volvo.ai.analytic.center.feign.RemoteCarModelClient; import com.volvo.ai.analytic.center.mapper.TmOdsVdqwExternalcontactMapper; import com.volvo.ai.analytic.center.mapper.TmOdsVdqwMessagearchivingMapper; import com.volvo.ai.analytic.center.mapper.TmOdsVdqwWorkuserinfoMapper; +import com.volvo.ai.analytic.center.mapper.TtVdqwRecordMapper; import com.volvo.ai.analytic.center.service.*; import com.volvo.ai.analytic.center.utils.ConstantStr; import com.volvo.ai.analytic.center.utils.FlowResultSplitUtil; @@ -32,7 +33,10 @@ import org.springframework.stereotype.Service; import javax.annotation.Resource; import java.time.LocalDate; -import java.util.*; +import java.util.Arrays; +import java.util.HashMap; +import java.util.List; +import java.util.Map; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -58,6 +62,10 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl messageList = tmOdsVdqwMessagearchivingMapper.queryOdsVdqwMessageByData(statTime, endTime, offset, pageSize); + List messageList = tmOdsVdqwMessagearchivingMapper.queryOdsVdqwMessageByData(statTime, endTime, offset, pageSize, retry); // 处理查询到的数据 // 使用 CompletableFuture 并行处理 CompletableFuture[] futures = messageList.stream() @@ -219,6 +229,7 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl + and not exists( + select 1 from tt_vdqw_record t + where t.from_accept_user_id = CONCAT(LEAST(tovm.from_user_id, tovm.accept_user_id),'_',GREATEST(tovm.from_user_id, tovm.accept_user_id)) + AND t.msg_time between #{statTime} and #{endTime} + ) + GROUP BY LEAST(tovm.from_user_id, tovm.accept_user_id), GREATEST(tovm.from_user_id, tovm.accept_user_id) @@ -26,12 +33,20 @@ tovm.from_user_id as fromUserId, tovm.accept_user_id as acceptUserId, tovm.msg_time as msgTime, - tovm.chat_type as chatType + tovm.chat_type as chatType, + CONCAT(LEAST(tovm.from_user_id, tovm.accept_user_id),'_',GREATEST(tovm.from_user_id, tovm.accept_user_id)) as fromAcceptUserId FROM `tm_ods_vdqw_messagearchiving` tovm WHERE tovm.is_deleted = 0 AND tovm.msg_time between #{statTime} and #{endTime} AND tovm.chat_type = 0 + + and not exists( + select 1 from tt_vdqw_record t + where t.from_accept_user_id = CONCAT(LEAST(tovm.from_user_id, tovm.accept_user_id),'_',GREATEST(tovm.from_user_id, tovm.accept_user_id)) + AND t.msg_time between #{statTime} and #{endTime} + ) + GROUP BY LEAST(tovm.from_user_id, tovm.accept_user_id), GREATEST(tovm.from_user_id, tovm.accept_user_id)