feat(archive): 优化客户昵称补全机制并添加限流和状态管理

- 增加异步默认线程池,避免批量补充客户资料时线程过多问题
- 为客户详情接口新增原始响应获取方法,支持错误码区分临时和永久失败
- 客户消息处理时,跳过已删除或拉黑的不活跃客户,避免无效接口调用
- 实现客户昵称补全重试节流,5分钟内同一客户只尝试一次补全
- 标记企微返回84061错误的客户为不活跃状态,停止无意义的昵称补充
- 支持客户状态恢复逻辑,接收新消息时重置为活跃状态
- 增加定时任务每30分钟补全昵称,针对首次失败遗留的空或ID昵称客户
- 内部调用异步方法时改为通过自注入代理,确保@Async生效
- 修复昵称补全过程中异常处理逻辑,防止失败导致线程中断问题
This commit is contained in:
jqb 2026-08-11 18:17:45 +08:00
parent 2342451983
commit f49b6d3437
5 changed files with 202 additions and 9 deletions

View File

@ -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 {

View File

@ -182,6 +182,30 @@ public class WeComApiClient {
}
}
/**
* 获取客户详情(保留原始响应,含 errcode)
* 供需要区分永久失败(如84061)与临时失败的调用方使用
* 网络异常时返回空 Map
*/
public Map<String, Object> 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<String> response = restTemplate.getForEntity(url, String.class);
Map<String, Object> 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();
}
}
/**
* 获取客户列表(分页)
*/

View File

@ -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<Customer>().in(Customer::getCustomerId, userIds));
Map<String, String> customerNameMap = customers.stream()
.collect(Collectors.toMap(Customer::getCustomerId, c -> c.getName() != null ? c.getName() : c.getCustomerId(), (a, b) -> a));
// 已删除/拉黑客户(如84061)不再实时调企微接口
Set<String> inactiveIds = customers.stream()
.filter(archiveContactService::isCustomerInactive)
.map(Customer::getCustomerId)
.collect(Collectors.toSet());
Map<String, String> 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 {

View File

@ -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<String, Long> fillAttemptTimes = new ConcurrentHashMap<>();
private static final long FILL_RETRY_INTERVAL_MS = 5 * 60 * 1000L;
/**
* 不活跃客户状态:已删除/已拉黑的客户无法从企微获取昵称,不做补全重试
*/
private static final Set<String> 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<String, Object> detail = weComApiClient.getExternalContact(customerId);
if (detail == null || detail.isEmpty()) return;
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return;
}
doFillCustomerDetail(customerId);
}
Map<String, Object> contact = (Map<String, Object>) 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<String, Object> 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<String, Object> contact = (Map<String, Object>) 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<Map<String, Object>> follows = (List<Map<String, Object>>) detail.get("follow_user");
List<Map<String, Object>> follows = (List<Map<String, Object>>) 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<Customer> candidates = customerMapper.selectList(
new LambdaQueryWrapper<Customer>()
.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 清零客户今日消息数
*/

View File

@ -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<String, Object> detail = weComApiClient.getExternalContact(customerId);
if (detail == null || detail.isEmpty()) {