From f49b6d3437b62fbd7025c4c9918a39eb3ac7fd1f Mon Sep 17 00:00:00 2001 From: jqb Date: Tue, 11 Aug 2026 18:17:45 +0800 Subject: [PATCH] =?UTF-8?q?feat(archive):=20=E4=BC=98=E5=8C=96=E5=AE=A2?= =?UTF-8?q?=E6=88=B7=E6=98=B5=E7=A7=B0=E8=A1=A5=E5=85=A8=E6=9C=BA=E5=88=B6?= =?UTF-8?q?=E5=B9=B6=E6=B7=BB=E5=8A=A0=E9=99=90=E6=B5=81=E5=92=8C=E7=8A=B6?= =?UTF-8?q?=E6=80=81=E7=AE=A1=E7=90=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 增加异步默认线程池,避免批量补充客户资料时线程过多问题 - 为客户详情接口新增原始响应获取方法,支持错误码区分临时和永久失败 - 客户消息处理时,跳过已删除或拉黑的不活跃客户,避免无效接口调用 - 实现客户昵称补全重试节流,5分钟内同一客户只尝试一次补全 - 标记企微返回84061错误的客户为不活跃状态,停止无意义的昵称补充 - 支持客户状态恢复逻辑,接收新消息时重置为活跃状态 - 增加定时任务每30分钟补全昵称,针对首次失败遗留的空或ID昵称客户 - 内部调用异步方法时改为通过自注入代理,确保@Async生效 - 修复昵称补全过程中异常处理逻辑,防止失败导致线程中断问题 --- .../archive/ArchiveServiceApplication.java | 17 ++ .../artedu/archive/client/WeComApiClient.java | 24 +++ .../controller/ArchiveMessageController.java | 10 +- .../service/ArchiveContactService.java | 152 +++++++++++++++++- .../service/ArchiveSessionService.java | 8 + 5 files changed, 202 insertions(+), 9 deletions(-) diff --git a/backend/archive-service/src/main/java/com/artedu/archive/ArchiveServiceApplication.java b/backend/archive-service/src/main/java/com/artedu/archive/ArchiveServiceApplication.java index 2d8dd19..47046dd 100644 --- a/backend/archive-service/src/main/java/com/artedu/archive/ArchiveServiceApplication.java +++ b/backend/archive-service/src/main/java/com/artedu/archive/ArchiveServiceApplication.java @@ -22,6 +22,23 @@ public class ArchiveServiceApplication { return new RestTemplate(); } + /** + * @Async 默认执行器:有界线程池,避免批量补充客户资料时创建大量线程 + */ + @Bean("taskExecutor") + public org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor taskExecutor() { + org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor executor = + new org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor(); + executor.setCorePoolSize(2); + executor.setMaxPoolSize(4); + executor.setQueueCapacity(2000); + executor.setThreadNamePrefix("archive-fill-"); + // 队列满时由调用线程执行,起到限流作用,不丢弃任务 + executor.setRejectedExecutionHandler(new java.util.concurrent.ThreadPoolExecutor.CallerRunsPolicy()); + executor.initialize(); + return executor; + } + public static void main(String[] args) { if (args.length > 0 && "worker".equals(args[0])) { try { diff --git a/backend/archive-service/src/main/java/com/artedu/archive/client/WeComApiClient.java b/backend/archive-service/src/main/java/com/artedu/archive/client/WeComApiClient.java index 91f5f07..217ac41 100644 --- a/backend/archive-service/src/main/java/com/artedu/archive/client/WeComApiClient.java +++ b/backend/archive-service/src/main/java/com/artedu/archive/client/WeComApiClient.java @@ -182,6 +182,30 @@ public class WeComApiClient { } } + /** + * 获取客户详情(保留原始响应,含 errcode) + * 供需要区分永久失败(如84061)与临时失败的调用方使用 + * 网络异常时返回空 Map + */ + public Map getExternalContactRaw(String externalUserId) { + String accessToken = getAccessToken(); + if (accessToken == null) return Collections.emptyMap(); + + try { + String url = "https://qyapi.weixin.qq.com/cgi-bin/externalcontact/get?access_token=" + accessToken + "&external_userid=" + externalUserId; + ResponseEntity response = restTemplate.getForEntity(url, String.class); + Map result = JsonUtils.fromJsonMap(response.getBody()); + if (result == null) return Collections.emptyMap(); + if (!Integer.valueOf(0).equals(result.get("errcode"))) { + log.warn("获取客户详情失败: externalUserId={}, resp={}", externalUserId, response.getBody()); + } + return result; + } catch (Exception e) { + log.error("获取客户详情异常: externalUserId={}", externalUserId, e); + return Collections.emptyMap(); + } + } + /** * 获取客户列表(分页) */ diff --git a/backend/archive-service/src/main/java/com/artedu/archive/controller/ArchiveMessageController.java b/backend/archive-service/src/main/java/com/artedu/archive/controller/ArchiveMessageController.java index 8d49c79..c2a885b 100644 --- a/backend/archive-service/src/main/java/com/artedu/archive/controller/ArchiveMessageController.java +++ b/backend/archive-service/src/main/java/com/artedu/archive/controller/ArchiveMessageController.java @@ -67,6 +67,9 @@ public class ArchiveMessageController { @Autowired private WeComApiClient weComApiClient; + @Autowired + private com.artedu.archive.service.ArchiveContactService archiveContactService; + @Value("${wecom.archive.corp-id}") private String corpId; @@ -284,6 +287,11 @@ public class ArchiveMessageController { new LambdaQueryWrapper().in(Customer::getCustomerId, userIds)); Map customerNameMap = customers.stream() .collect(Collectors.toMap(Customer::getCustomerId, c -> c.getName() != null ? c.getName() : c.getCustomerId(), (a, b) -> a)); + // 已删除/拉黑客户(如84061)不再实时调企微接口 + Set inactiveIds = customers.stream() + .filter(archiveContactService::isCustomerInactive) + .map(Customer::getCustomerId) + .collect(Collectors.toSet()); Map nameMap = new java.util.HashMap<>(); nameMap.putAll(staffNameMap); @@ -294,7 +302,7 @@ public class ArchiveMessageController { for (String uid : userIds) { String currentName = nameMap.get(uid); boolean isExternal = isExternalUserId(uid); - boolean needFetch = isExternal && (currentName == null || currentName.equals(uid)); + boolean needFetch = isExternal && !inactiveIds.contains(uid) && (currentName == null || currentName.equals(uid)); log.info("处理uid={}, currentName={}, isExternal={}, needFetch={}", uid, currentName, isExternal, needFetch); if (needFetch) { try { diff --git a/backend/archive-service/src/main/java/com/artedu/archive/service/ArchiveContactService.java b/backend/archive-service/src/main/java/com/artedu/archive/service/ArchiveContactService.java index 156a623..805c138 100644 --- a/backend/archive-service/src/main/java/com/artedu/archive/service/ArchiveContactService.java +++ b/backend/archive-service/src/main/java/com/artedu/archive/service/ArchiveContactService.java @@ -14,6 +14,7 @@ import com.baomidou.mybatisplus.extension.plugins.pagination.Page; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; +import org.springframework.context.annotation.Lazy; import org.springframework.scheduling.annotation.Async; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; @@ -22,6 +23,7 @@ import java.time.Instant; import java.time.LocalDate; import java.time.ZoneId; import java.util.*; +import java.util.concurrent.ConcurrentHashMap; import java.util.stream.Collectors; /** @@ -47,6 +49,23 @@ public class ArchiveContactService { @Autowired private WeComApiClient weComApiClient; + /** + * 自注入代理:类内 this 调用不走 Spring 代理,@Async 会失效,需通过代理对象调用 + */ + @Autowired + @Lazy + private ArchiveContactService self; + + /** + * 昵称补全重试节流:customerId -> 上次尝试时间戳,避免消息高峰期高频调用企微API + */ + private final Map fillAttemptTimes = new ConcurrentHashMap<>(); + private static final long FILL_RETRY_INTERVAL_MS = 5 * 60 * 1000L; + /** + * 不活跃客户状态:已删除/已拉黑的客户无法从企微获取昵称,不做补全重试 + */ + private static final Set INACTIVE_STATUSES = new HashSet<>(Arrays.asList("DELETED", "BLOCKED")); + /** * 处理消息,提取并维护员工/客户信息 */ @@ -142,7 +161,20 @@ public class ArchiveContactService { log.info("自动创建客户记录: customerId={}", customerId); // 异步补充资料 - asyncFillCustomerDetail(customerId); + self.asyncFillCustomerDetail(customerId); + } else if (isExternalUserId(customerId)) { + // 曾标记删除/拉黑的客户又有新消息,说明好友关系已恢复,重置状态 + if (INACTIVE_STATUSES.contains(existing.getStatus())) { + existing.setStatus("ACTIVE"); + log.info("客户状态重置为ACTIVE(收到新消息): customerId={}", customerId); + } + // 昵称为空或仍是ID(首次补充失败的历史遗留),借新消息到达的机会重试补充 + String currentName = existing.getName(); + boolean nameMissing = currentName == null || currentName.trim().isEmpty() || currentName.equals(customerId); + if (nameMissing && allowFillAttempt(customerId)) { + log.info("客户昵称缺失,借新消息触发重新补充: customerId={}", customerId); + self.asyncFillCustomerDetail(customerId); + } } // 更新统计 @@ -246,11 +278,48 @@ public class ArchiveContactService { public void asyncFillCustomerDetail(String customerId) { try { Thread.sleep(100); - Map detail = weComApiClient.getExternalContact(customerId); - if (detail == null || detail.isEmpty()) return; + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return; + } + doFillCustomerDetail(customerId); + } - Map contact = (Map) detail.get("external_contact"); - if (contact == null) return; + /** + * 补全重试节流:同一客户 5 分钟内只允许尝试一次 + */ + private boolean allowFillAttempt(String customerId) { + long now = System.currentTimeMillis(); + Long last = fillAttemptTimes.get(customerId); + if (last != null && now - last < FILL_RETRY_INTERVAL_MS) { + return false; + } + // 防止缓存无限增长 + if (fillAttemptTimes.size() > 10000) { + fillAttemptTimes.clear(); + } + fillAttemptTimes.put(customerId, now); + return true; + } + + /** + * 同步补充客户资料(昵称、头像等),返回是否补充成功 + */ + private boolean doFillCustomerDetail(String customerId) { + try { + Map result = weComApiClient.getExternalContactRaw(customerId); + if (result == null || result.isEmpty()) return false; + + Object errcodeObj = result.get("errcode"); + // 84061 = not external contact(已删除好友/拉黑等):标记客户状态为DELETED,停止重试 + if (errcodeObj instanceof Number && ((Number) errcodeObj).intValue() == 84061) { + markCustomerInactive(customerId, "DELETED"); + return false; + } + if (!Integer.valueOf(0).equals(errcodeObj)) return false; + + Map contact = (Map) result.get("external_contact"); + if (contact == null) return false; String name = (String) contact.get("name"); String avatar = (String) contact.get("avatar"); @@ -259,7 +328,7 @@ public class ArchiveContactService { String unionid = (String) contact.get("unionid"); // follow_info 中的 add_time - List> follows = (List>) detail.get("follow_user"); + List> follows = (List>) result.get("follow_user"); Long addTime = null; if (follows != null && !follows.isEmpty()) { Object at = follows.get(0).get("createtime"); @@ -282,8 +351,33 @@ public class ArchiveContactService { customerMapper.updateById(customer); log.info("补充客户资料成功: customerId={}, name={}", customerId, name); } + return true; } catch (Exception e) { log.warn("补充客户资料失败: customerId={}", customerId, e); + return false; + } + } + + /** + * 客户是否处于不活跃状态(DELETED/BLOCKED),无需补全昵称 + */ + public boolean isCustomerInactive(Customer customer) { + return customer != null && INACTIVE_STATUSES.contains(customer.getStatus()); + } + + /** + * 标记客户为不活跃状态(如企微返回84061已删除好友),停止无意义的昵称补全重试 + */ + private void markCustomerInactive(String customerId, String status) { + try { + Customer customer = customerMapper.selectByCustomerId(customerId, corpId); + if (customer != null && !status.equals(customer.getStatus())) { + customer.setStatus(status); + customerMapper.updateById(customer); + log.info("客户状态标记为{}(84061),停止昵称补全重试: customerId={}", status, customerId); + } + } catch (Exception e) { + log.warn("更新客户状态异常: customerId={}", customerId, e); } } @@ -485,7 +579,7 @@ public class ArchiveContactService { // 异步补充所有客户详情 for (String customerId : processed) { - asyncFillCustomerDetail(customerId); + self.asyncFillCustomerDetail(customerId); } log.info("同步客户完成,共 {} 人", count); @@ -507,7 +601,7 @@ public class ArchiveContactService { log.info("开始批量补充客户资料,共 {} 人", customers.size()); for (Customer customer : customers) { try { - asyncFillCustomerDetail(customer.getCustomerId()); + self.asyncFillCustomerDetail(customer.getCustomerId()); count++; // 稍微延迟避免触发频率限制 if (count % 100 == 0) { @@ -524,6 +618,48 @@ public class ArchiveContactService { return count; } + /** + * 定时兜底补全:每 30 分钟扫描昵称为空或仍为 ID 的客户(首次补充失败遗留),重新从企微获取 + */ + @Scheduled(fixedDelay = 30 * 60 * 1000, initialDelay = 3 * 60 * 1000) + public void repairMissingCustomerNames() { + try { + List candidates = customerMapper.selectList( + new LambdaQueryWrapper() + .eq(Customer::getCorpId, corpId) + .and(w -> w.isNull(Customer::getName) + .or().eq(Customer::getName, "") + .or().apply("name = customer_id")) + // 排除已删除/拉黑客户(如84061),避免无限重试 + .notIn(Customer::getStatus, INACTIVE_STATUSES) + ); + if (candidates == null || candidates.isEmpty()) { + return; + } + log.info("定时补全客户昵称开始,待处理 {} 人", candidates.size()); + int repaired = 0; + for (Customer customer : candidates) { + String customerId = customer.getCustomerId(); + if (!isExternalUserId(customerId)) continue; + try { + if (doFillCustomerDetail(customerId)) { + repaired++; + } + // 延迟避免触发企微API频率限制 + Thread.sleep(300); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + break; + } catch (Exception e) { + log.warn("定时补全客户昵称失败: customerId={}", customerId, e); + } + } + log.info("定时补全客户昵称完成,修复 {} 人 / 共 {} 人", repaired, candidates.size()); + } catch (Exception e) { + log.error("定时补全客户昵称任务失败", e); + } + } + /** * 每天凌晨 00:05 清零客户今日消息数 */ diff --git a/backend/archive-service/src/main/java/com/artedu/archive/service/ArchiveSessionService.java b/backend/archive-service/src/main/java/com/artedu/archive/service/ArchiveSessionService.java index b90680a..64ecdb9 100644 --- a/backend/archive-service/src/main/java/com/artedu/archive/service/ArchiveSessionService.java +++ b/backend/archive-service/src/main/java/com/artedu/archive/service/ArchiveSessionService.java @@ -48,6 +48,9 @@ public class ArchiveSessionService { @Autowired private WeComApiClient weComApiClient; + @Autowired + private ArchiveContactService archiveContactService; + private static final java.time.format.DateTimeFormatter TIME_FORMATTER = java.time.format.DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"); @@ -101,6 +104,11 @@ public class ArchiveSessionService { if (!isExternalUserId(customerId) || (name != null && !name.trim().isEmpty() && !customerId.equals(name))) { continue; } + // 已删除/拉黑客户(如84061)不再实时调企微接口 + Customer dbCustomer = customerMapper.selectByCustomerId(customerId, corpId); + if (archiveContactService.isCustomerInactive(dbCustomer)) { + continue; + } try { Map detail = weComApiClient.getExternalContact(customerId); if (detail == null || detail.isEmpty()) {