修复双实例 owner-scoped 结果文件 job 被非归属实例误杀
Build Backend JAR / build (push) Has been cancelled

线上现象:货源查询任务 26402/26400(归属 server-121)在归属机正常推进,
却被 server-110 判为「结果生成失败」,错误为「该任务已绑定到另一台服务实例处理」。

两个缺陷叠加:
1. 两台实例共用一个 RocketMQ consumer group,dispatch 消息被轮询投递;
   worker 对"PENDING + owner 不匹配"的 job 选择本地抢跑,模块层
   ensureTaskOwnedByCurrentInstance 直接抛 TaskOwnerMismatchException,
   异常落进通用 catch 被当成生成失败,5 次抢跑即耗尽 retryCount 并终态失败。
2. runPendingJobs 双重 claim:claimRunnableJobsForOwner 已把 job 翻成 RUNNING,
   随后又走 claimRunning(条件 PENDING/FAILED)必然失败 → job 被静默丢弃,
   归属机的定时兜底从未真正生效。

修复:
- worker:非归属实例一律跳过 owner-scoped job;抽出 processClaimed 供已持有
  claim 的调用方使用;processInternal 单独捕获 TaskOwnerMismatchException,
  只 deferRunning 退回 PENDING,不计入重试。
- TaskFileJobLocalDispatcher 新增 dispatchClaimed。
- dispatch 事件带上 owner,coordinator 按归属路由:本机 owner 直接本地派发不发 MQ,
  他机 owner 不发不跑交由归属机轮询接手,owner 为空维持原有 MQ 行为。
- 新增 owner 路由回归测试两组,回滚配合"本地抢跑"的旧断言。

同时提交此前工作区内已随 JAR 上线的 backend-java 改动:SimilarAsin 解析载荷
groups 内嵌 items(Python 旧链路兼容)与 chunk 类型不匹配跳过、
PermissionMenuSchemaInitializer 与 V103 两级菜单分组迁移、相关测试。
This commit is contained in:
2026-09-01 21:07:25 +08:00
parent 55517f996b
commit 5f4fcad2ef
16 changed files with 679 additions and 44 deletions
@@ -47,20 +47,34 @@ public class PermissionMenuSchemaInitializer {
new DefaultAppChildMenu("ERP", "erp", "erp", 132, "brand_logistics_tools") new DefaultAppChildMenu("ERP", "erp", "erp", 132, "brand_logistics_tools")
); );
/**
* 后台一级分组:只做权限树层级与「上级菜单」候选,不映射真实页面。
* route_path 使用 group- 前缀虚拟值以满足 (menu_type, route_path) 唯一约束;
* 前端按 column_key 前缀 admin_group_ 识别分组行,不会出现在页面选择器中。
*/
private static final List<DefaultAdminGroup> DEFAULT_ADMIN_GROUPS = List.of(
new DefaultAdminGroup("账号与权限", "admin_group_account", "group-account", 5),
new DefaultAdminGroup("数据管理", "admin_group_data", "group-data", 15),
new DefaultAdminGroup("店铺管理", "admin_group_shop", "group-shop", 45),
new DefaultAdminGroup("记录与版本", "admin_group_record", "group-record", 55)
);
private static final List<DefaultAdminMenu> DEFAULT_ADMIN_MENUS = List.of( private static final List<DefaultAdminMenu> DEFAULT_ADMIN_MENUS = List.of(
new DefaultAdminMenu("用户管理", "admin_users", "users", 10), new DefaultAdminMenu("用户管理", "admin_users", "users", 10, "admin_group_account"),
new DefaultAdminMenu("菜单权限配置", "admin_columns", "columns", 20), new DefaultAdminMenu("菜单权限配置", "admin_columns", "columns", 20, "admin_group_account"),
new DefaultAdminMenu("分组管理", "admin_group_manage", "group-manage", 25), new DefaultAdminMenu("分组管理", "admin_group_manage", "group-manage", 25, "admin_group_account"),
new DefaultAdminMenu("去重数据汇总", "admin_dedupe_total_data", "dedupe-total-data", 30), new DefaultAdminMenu("去重数据汇总", "admin_dedupe_total_data", "dedupe-total-data", 30, "admin_group_data"),
new DefaultAdminMenu("无效ASIN数据", "admin_invalid_asin_data", "invalid-asin-data", 35), new DefaultAdminMenu("品牌数据", "admin_invalid_asin_data", "invalid-asin-data", 35, "admin_group_data"),
new DefaultAdminMenu("店铺密钥管理", "admin_shop_keys", "shop-keys", 40), new DefaultAdminMenu("查询ASIN", "admin_query_asin", "query-asin", 65, "admin_group_data"),
new DefaultAdminMenu("店铺管理", "admin_shop_manage", "shop-manage", 50), new DefaultAdminMenu("商品类目", "admin_product_categories", "product-categories", 66, "admin_group_data"),
new DefaultAdminMenu("跳过跟价ASIN", "admin_skip_price_asin", "skip-price-asin", 60), new DefaultAdminMenu("店铺密钥管理", "admin_shop_keys", "shop-keys", 40, "admin_group_shop"),
new DefaultAdminMenu("查询ASIN", "admin_query_asin", "query-asin", 65), new DefaultAdminMenu("店铺管理", "admin_shop_manage", "shop-manage", 50, "admin_group_shop"),
new DefaultAdminMenu("商品类目", "admin_product_categories", "product-categories", 66), new DefaultAdminMenu("最低价ASIN设置", "admin_skip_price_asin", "skip-price-asin", 60, "admin_group_shop"),
new DefaultAdminMenu("生成记录", "admin_history", "history", 70), new DefaultAdminMenu("店铺数据记录", "admin_shop_data_crawl_tasks", "shop-data-crawl-tasks", 82, "admin_group_shop"),
new DefaultAdminMenu("软件版本管理", "admin_version", "version", 80), new DefaultAdminMenu("生成记录", "admin_history", "history", 70, "admin_group_record"),
new DefaultAdminMenu("数字人版本管理", "digital_human_version", "digital-human-version", 81) new DefaultAdminMenu("视频任务记录", "admin_image_video_tasks", "image-video-tasks", 75, "admin_group_record"),
new DefaultAdminMenu("软件版本管理", "admin_version", "version", 80, "admin_group_record"),
new DefaultAdminMenu("数字人版本管理", "digital_human_version", "digital-human-version", 81, "admin_group_record")
); );
@EventListener(ApplicationReadyEvent.class) @EventListener(ApplicationReadyEvent.class)
@@ -86,6 +100,7 @@ public class PermissionMenuSchemaInitializer {
executeQuietly("UPDATE columns SET menu_type = 'app' WHERE menu_type IS NULL OR menu_type = ''"); executeQuietly("UPDATE columns SET menu_type = 'app' WHERE menu_type IS NULL OR menu_type = ''");
executeQuietly("UPDATE columns SET sort_order = id WHERE sort_order IS NULL OR sort_order = 0"); executeQuietly("UPDATE columns SET sort_order = id WHERE sort_order IS NULL OR sort_order = 0");
executeQuietly("ALTER TABLE columns ADD UNIQUE KEY uk_menu_type_route_path (menu_type, route_path)"); executeQuietly("ALTER TABLE columns ADD UNIQUE KEY uk_menu_type_route_path (menu_type, route_path)");
ensureDefaultAdminGroups();
ensureDefaultAdminMenus(); ensureDefaultAdminMenus();
boolean shopDataCrawlDataPermissionCreated = ensureShopDataCrawlAdminDefaults(); boolean shopDataCrawlDataPermissionCreated = ensureShopDataCrawlAdminDefaults();
ensureDefaultAppMenus(); ensureDefaultAppMenus();
@@ -105,20 +120,37 @@ public class PermissionMenuSchemaInitializer {
} }
} }
private void ensureDefaultAdminMenus() { private void ensureDefaultAdminGroups() {
for (DefaultAdminMenu menu : DEFAULT_ADMIN_MENUS) { for (DefaultAdminGroup group : DEFAULT_ADMIN_GROUPS) {
executeQuietly(""" executeQuietly("""
INSERT INTO columns (name, column_key, menu_type, route_path, sort_order) INSERT INTO columns (name, column_key, menu_type, route_path, parent_id, sort_order)
SELECT '%s', '%s', 'admin', '%s', %d SELECT '%s', '%s', 'admin', '%s', NULL, %d
WHERE NOT EXISTS ( WHERE NOT EXISTS (
SELECT 1 FROM columns WHERE column_key = '%s' SELECT 1 FROM columns WHERE column_key = '%s'
) )
""".formatted(menu.name(), menu.columnKey(), menu.routePath(), menu.sortOrder(), menu.columnKey())); """.formatted(group.name(), group.columnKey(), group.routePath(), group.sortOrder(), group.columnKey()));
}
}
private void ensureDefaultAdminMenus() {
for (DefaultAdminMenu menu : DEFAULT_ADMIN_MENUS) {
executeQuietly(""" executeQuietly("""
UPDATE columns INSERT INTO columns (name, column_key, menu_type, route_path, sort_order, parent_id)
SET sort_order = %d, parent_id = NULL SELECT '%s', '%s', 'admin', '%s', %d, parent.id
WHERE column_key = '%s' AND (sort_order IS NULL OR sort_order = 0) FROM columns parent
""".formatted(menu.sortOrder(), menu.columnKey())); WHERE parent.column_key = '%s'
AND NOT EXISTS (
SELECT 1 FROM columns WHERE column_key = '%s'
)
""".formatted(menu.name(), menu.columnKey(), menu.routePath(), menu.sortOrder(),
menu.parentKey(), menu.columnKey()));
executeQuietly("""
UPDATE columns child
JOIN columns parent ON parent.column_key = '%s'
SET child.name = '%s', child.sort_order = %d,
child.parent_id = COALESCE(child.parent_id, parent.id)
WHERE child.column_key = '%s'
""".formatted(menu.parentKey(), menu.name(), menu.sortOrder(), menu.columnKey()));
} }
} }
@@ -179,8 +211,9 @@ public class PermissionMenuSchemaInitializer {
private boolean ensureShopDataCrawlAdminDefaults() { private boolean ensureShopDataCrawlAdminDefaults() {
executeQuietly(""" executeQuietly("""
INSERT INTO columns (name, column_key, menu_type, route_path, sort_order) INSERT INTO columns (name, column_key, menu_type, route_path, sort_order, parent_id)
SELECT '店铺数据记录', 'admin_shop_data_crawl_tasks', 'admin', 'shop-data-crawl-tasks', 82 SELECT '店铺数据记录', 'admin_shop_data_crawl_tasks', 'admin', 'shop-data-crawl-tasks', 82,
(SELECT id FROM columns WHERE column_key = 'admin_group_shop' LIMIT 1)
WHERE NOT EXISTS ( WHERE NOT EXISTS (
SELECT 1 FROM columns WHERE column_key = 'admin_shop_data_crawl_tasks' SELECT 1 FROM columns WHERE column_key = 'admin_shop_data_crawl_tasks'
) )
@@ -226,7 +259,11 @@ public class PermissionMenuSchemaInitializer {
} }
} }
private record DefaultAdminMenu(String name, String columnKey, String routePath, int sortOrder) { private record DefaultAdminMenu(String name, String columnKey, String routePath, int sortOrder,
String parentKey) {
}
private record DefaultAdminGroup(String name, String columnKey, String routePath, int sortOrder) {
} }
private record DefaultAppMenu(String name, String columnKey, String routePath, int sortOrder) { private record DefaultAppMenu(String name, String columnKey, String routePath, int sortOrder) {
@@ -5127,10 +5127,47 @@ public class SimilarAsinTaskService {
payload.setSourceFiles(sourceFiles == null ? List.of() : sourceFiles); payload.setSourceFiles(sourceFiles == null ? List.of() : sourceFiles);
payload.setHeaders(headers == null ? List.of() : headers); payload.setHeaders(headers == null ? List.of() : headers);
payload.setItems(allRows == null ? List.of() : new ArrayList<>(allRows)); payload.setItems(allRows == null ? List.of() : new ArrayList<>(allRows));
payload.setGroups(groups == null ? List.of() : groups); // 兼容 Python 旧链路groups 需内嵌完整行Python normalize_groups
// 直接返回 groups 并按 gp["items"] 取行仅有 startIndex/endIndex 引用时
// items 为空任务空转startIndex/endIndex 保留Java 内部 expandGroupRefs 继续使用
payload.setGroups(fillGroupItems(groups, allRows));
return writeJson(payload, "保存解析结果失败"); return writeJson(payload, "保存解析结果失败");
} }
/**
* group.startIndex/endIndex 半开区间从全量行切片内嵌回 group.items
* 无索引或切片越界时退回按 group.items 原样兼容旧 payload 反序列化场景
*/
static List<SimilarAsinParsedGroupVo> fillGroupItems(List<SimilarAsinParsedGroupVo> groups, List<SimilarAsinParsedRowVo> allRows) {
if (groups == null || groups.isEmpty()) {
return groups == null ? List.of() : groups;
}
if (allRows == null || allRows.isEmpty()) {
return groups;
}
List<SimilarAsinParsedRowVo> rows = new ArrayList<>(allRows);
for (SimilarAsinParsedGroupVo group : groups) {
if (group == null) {
continue;
}
Integer start = group.getStartIndex();
Integer end = group.getEndIndex();
if (start != null && end != null) {
int from = Math.max(0, Math.min(start, rows.size()));
int to = Math.max(from, Math.min(end, rows.size()));
if (to > from) {
group.setItems(new ArrayList<>(rows.subList(from, to)));
continue;
}
}
// 索引缺失/无效保留已有 items兼容旧 payload不覆盖
if (group.getItems() == null) {
group.setItems(new ArrayList<>());
}
}
return groups;
}
private String buildTaskResultJson(String aiPrompt, String apiKey, Boolean imgSwitch, Boolean categorySwitch, List<SimilarAsinSourceFileDto> sourceFiles, String parsedPayloadPointer) { private String buildTaskResultJson(String aiPrompt, String apiKey, Boolean imgSwitch, Boolean categorySwitch, List<SimilarAsinSourceFileDto> sourceFiles, String parsedPayloadPointer) {
Map<String, Object> payload = new LinkedHashMap<>(); Map<String, Object> payload = new LinkedHashMap<>();
payload.put("aiPrompt", normalize(aiPrompt)); payload.put("aiPrompt", normalize(aiPrompt));
@@ -5172,6 +5209,11 @@ public class SimilarAsinTaskService {
} }
private SimilarAsinParsedPayloadDto hydrateParsedPayloadRows(SimilarAsinParsedPayloadDto payload) { private SimilarAsinParsedPayloadDto hydrateParsedPayloadRows(SimilarAsinParsedPayloadDto payload) {
return hydrateForPythonCompat(payload);
}
/** 解析载荷 hydrate:回填 allItems/items/groups 内嵌 itemsPython 旧链路兼容)。 */
static SimilarAsinParsedPayloadDto hydrateForPythonCompat(SimilarAsinParsedPayloadDto payload) {
if (payload == null) { if (payload == null) {
return new SimilarAsinParsedPayloadDto(); return new SimilarAsinParsedPayloadDto();
} }
@@ -5182,6 +5224,13 @@ public class SimilarAsinTaskService {
if (payload.getItems() == null || payload.getItems().isEmpty()) { if (payload.getItems() == null || payload.getItems().isEmpty()) {
payload.setItems(rows); payload.setItems(rows);
} }
// Python 旧链路兼容读取时也回填 groups 内嵌 items存量任务 payload
// 只有引用否则 Python normalize_groups 拿到空 items 任务空转
if (payload.getGroups() != null && !payload.getGroups().isEmpty()) {
List<SimilarAsinParsedRowVo> filledRows = payload.getItems() == null || payload.getItems().isEmpty()
? rows : payload.getItems();
payload.setGroups(fillGroupItems(payload.getGroups(), filledRows));
}
return payload; return payload;
} }
@@ -5205,6 +5254,16 @@ public class SimilarAsinTaskService {
// P1-4识别跨实例 local 指针读不到的场景"chunk 在另一实例"的元信息 // P1-4识别跨实例 local 指针读不到的场景"chunk 在另一实例"的元信息
// 通过 task_scope_state.last_error 留痕便于排查"为何 owner 切换后 chunk 读不到" // 通过 task_scope_state.last_error 留痕便于排查"为何 owner 切换后 chunk 读不到"
boolean crossInstance = msg.contains("only exists on instance="); boolean crossInstance = msg.contains("only exists on instance=");
// 类型不兼容如历史版本序列化数据含当前类加载器不认识的类型标记
// chunk 永久不可读跳过而不是把组装/收尾永久拖死留痕便于排查
boolean typeMismatch = ex instanceof TypeNotPresentException
|| msg.contains("not present");
if (typeMismatch) {
log.warn("[similar-asin] chunk payload type mismatch, skip chunk taskId={} chunk={} err={}",
chunk.getTaskId(), chunk.getChunkIndex(), msg);
recordChunkReadFailure(chunk, false, msg);
return rows;
}
log.warn("[similar-asin] read chunk payload failed taskId={} chunk={} crossInstance={} err={}", log.warn("[similar-asin] read chunk payload failed taskId={} chunk={} crossInstance={} err={}",
chunk.getTaskId(), chunk.getChunkIndex(), crossInstance, msg); chunk.getTaskId(), chunk.getChunkIndex(), crossInstance, msg);
recordChunkReadFailure(chunk, crossInstance, msg); recordChunkReadFailure(chunk, crossInstance, msg);
@@ -9,6 +9,7 @@ public record TaskFileJobDispatchEvent(
Long resultId, Long resultId,
String scopeKey, String scopeKey,
String jobType, String jobType,
String owner,
LocalDateTime createdAt LocalDateTime createdAt
) { ) {
} }
@@ -1,5 +1,6 @@
package com.nanri.aiimage.modules.task.service; package com.nanri.aiimage.modules.task.service;
import com.nanri.aiimage.config.InstanceMetadata;
import com.nanri.aiimage.modules.task.model.dto.TaskFileJobDispatchEvent; import com.nanri.aiimage.modules.task.model.dto.TaskFileJobDispatchEvent;
import com.nanri.aiimage.modules.task.model.dto.TaskFileJobMessage; import com.nanri.aiimage.modules.task.model.dto.TaskFileJobMessage;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
@@ -17,6 +18,7 @@ public class TaskFileJobDispatchCoordinator {
private final TaskFileJobPublisher taskFileJobPublisher; private final TaskFileJobPublisher taskFileJobPublisher;
private final TaskFileJobLocalDispatcher taskFileJobLocalDispatcher; private final TaskFileJobLocalDispatcher taskFileJobLocalDispatcher;
private final InstanceMetadata instanceMetadata;
@Qualifier("taskFileJobDispatchExecutor") @Qualifier("taskFileJobDispatchExecutor")
private final TaskExecutor taskFileJobDispatchExecutor; private final TaskExecutor taskFileJobDispatchExecutor;
@@ -43,6 +45,28 @@ public class TaskFileJobDispatchCoordinator {
} }
private void dispatchAfterCommit(TaskFileJobDispatchEvent event, TaskFileJobMessage message) { private void dispatchAfterCommit(TaskFileJobDispatchEvent event, TaskFileJobMessage message) {
// owner 非空 ⇔ owner-scoped job(只有 SIMILAR_ASIN/APPEARANCE_PATENT/PUBLISH/
// SHOP_DATA_CRAWL 会写 "task:{id}:owner:{instance}" 形式的 scopeKey)。
// 这类 job 必须在归属实例上执行,走共享 consumer group 的 MQ 会被轮询投递到
// 另一台,那边只能跳过,白白多绕一圈。
String owner = event.owner() == null ? "" : event.owner().trim();
if (!owner.isBlank()) {
if (owner.equals(currentInstanceId())) {
boolean dispatched = taskFileJobLocalDispatcher.dispatch(
event.jobId(), event.taskId(), event.moduleType(), true);
if (dispatched) {
return;
}
log.warn("[task-file-job] owner-scoped local dispatch declined, fallback to mq jobId={} taskId={} moduleType={} owner={}",
event.jobId(), event.taskId(), event.moduleType(), owner);
} else {
// 归属实例每 15 秒扫一次自己 owner 的 PENDING/FAILED job,这里不发 MQ、
// 不本地跑,避免消息落到非归属实例后被丢弃。
log.info("[task-file-job] skip dispatch because job belongs to another instance jobId={} taskId={} moduleType={} owner={} current={}",
event.jobId(), event.taskId(), event.moduleType(), owner, currentInstanceId());
return;
}
}
boolean published = taskFileJobPublisher.publish(message); boolean published = taskFileJobPublisher.publish(message);
if (published) { if (published) {
return; return;
@@ -56,4 +80,9 @@ public class TaskFileJobDispatchCoordinator {
event.jobId(), event.taskId(), event.moduleType()); event.jobId(), event.taskId(), event.moduleType());
} }
} }
private String currentInstanceId() {
String instanceId = instanceMetadata == null ? null : instanceMetadata.getInstanceId();
return instanceId == null || instanceId.isBlank() ? "unknown-instance" : instanceId;
}
} }
@@ -42,7 +42,19 @@ public class TaskFileJobLocalDispatcher {
return dispatch(jobId, taskId, moduleType, false); return dispatch(jobId, taskId, moduleType, false);
} }
/**
* 调用方已把 job claim 成 RUNNING 时使用:复用同样的背压/去重/异步执行,
* 但落到 worker.processClaimed,避免二次 claimRunning 失败导致 job 被静默丢弃。
*/
public boolean dispatchClaimed(Long jobId, Long taskId, String moduleType) {
return submit(jobId, taskId, moduleType, false, true);
}
public boolean dispatch(Long jobId, Long taskId, String moduleType, boolean force) { public boolean dispatch(Long jobId, Long taskId, String moduleType, boolean force) {
return submit(jobId, taskId, moduleType, force, false);
}
private boolean submit(Long jobId, Long taskId, String moduleType, boolean force, boolean alreadyClaimed) {
if (jobId == null || jobId <= 0) { if (jobId == null || jobId <= 0) {
return false; return false;
} }
@@ -63,7 +75,7 @@ public class TaskFileJobLocalDispatcher {
return true; return true;
} }
try { try {
taskFileJobDispatchExecutor.execute(() -> processLocally(jobId, taskId, moduleType)); taskFileJobDispatchExecutor.execute(() -> processLocally(jobId, taskId, moduleType, alreadyClaimed));
return true; return true;
} catch (RuntimeException ex) { } catch (RuntimeException ex) {
inflightJobIds.remove(jobId); inflightJobIds.remove(jobId);
@@ -71,7 +83,7 @@ public class TaskFileJobLocalDispatcher {
jobId, taskId, moduleType, ex.getMessage(), ex); jobId, taskId, moduleType, ex.getMessage(), ex);
if (force) { if (force) {
try { try {
processLocally(jobId, taskId, moduleType); processLocally(jobId, taskId, moduleType, alreadyClaimed);
return true; return true;
} finally { } finally {
inflightJobIds.remove(jobId); inflightJobIds.remove(jobId);
@@ -81,7 +93,7 @@ public class TaskFileJobLocalDispatcher {
} }
} }
private void processLocally(Long jobId, Long taskId, String moduleType) { private void processLocally(Long jobId, Long taskId, String moduleType, boolean alreadyClaimed) {
try { try {
TaskFileJobEntity job = taskFileJobMapper.selectById(jobId); TaskFileJobEntity job = taskFileJobMapper.selectById(jobId);
if (job == null) { if (job == null) {
@@ -95,7 +107,11 @@ public class TaskFileJobLocalDispatcher {
jobId, taskId, moduleType); jobId, taskId, moduleType);
return; return;
} }
taskResultFileJobWorker.process(job); if (alreadyClaimed) {
taskResultFileJobWorker.processClaimed(job);
} else {
taskResultFileJobWorker.process(job);
}
} catch (Exception ex) { } catch (Exception ex) {
log.warn("[task-file-job] local dispatch failed jobId={} taskId={} moduleType={} msg={}", log.warn("[task-file-job] local dispatch failed jobId={} taskId={} moduleType={} msg={}",
jobId, taskId, moduleType, ex.getMessage(), ex); jobId, taskId, moduleType, ex.getMessage(), ex);
@@ -57,7 +57,14 @@ public class TaskFileJobService {
if (existing == null) { if (existing == null) {
return null; return null;
} }
if ("SUCCESS".equals(existing.getStatus()) || "RUNNING".equals(existing.getStatus()) || "PENDING".equals(existing.getStatus())) { if ("SUCCESS".equals(existing.getStatus()) || "RUNNING".equals(existing.getStatus())) {
return existing;
}
if ("PENDING".equals(existing.getStatus())) {
// PENDING 僵尸兜底:job 入队时若 MQ 发送失败且本地兜底未生效(或消息被
// 非 owner 实例静默吞掉),job 会永远停留在 PENDING、无人处理。
// 补发一次 dispatch 事件,让调度链重试(幂等:consumer 侧 claim 有状态翻转保护)。
publishDispatchEvent(existing);
return existing; return existing;
} }
// 终态失败:retryCount 已达上限的 FAILED job 不再重置为 PENDING // 终态失败:retryCount 已达上限的 FAILED job 不再重置为 PENDING
@@ -237,6 +244,13 @@ public class TaskFileJobService {
.lt(TaskFileJobEntity::getUpdatedAt, threshold) .lt(TaskFileJobEntity::getUpdatedAt, threshold)
.or() .or()
.isNull(TaskFileJobEntity::getUpdatedAt))) .isNull(TaskFileJobEntity::getUpdatedAt)))
.or(stuckPending -> stuckPending
.eq(TaskFileJobEntity::getStatus, "PENDING")
.ne(TaskFileJobEntity::getRetryCount, MAX_RETRY_COUNT)
.and(age -> age
.lt(TaskFileJobEntity::getUpdatedAt, threshold)
.or()
.isNull(TaskFileJobEntity::getUpdatedAt)))
.or(exhausted -> exhausted .or(exhausted -> exhausted
.eq(TaskFileJobEntity::getStatus, "FAILED") .eq(TaskFileJobEntity::getStatus, "FAILED")
.ge(TaskFileJobEntity::getRetryCount, MAX_RETRY_COUNT) .ge(TaskFileJobEntity::getRetryCount, MAX_RETRY_COUNT)
@@ -251,6 +265,24 @@ public class TaskFileJobService {
exhaustedJobs.add(job); exhaustedJobs.add(job);
continue; continue;
} }
if ("PENDING".equals(job.getStatus())) {
// PENDING 僵尸兜底:补发 dispatch 事件(worker claim 幂等),
// 并累加 retryCount 防无限重发;注意不翻转状态,
// 避免与真实 PENDING 竞争 worker 的 claim 条件更新冲突。
int retryCount = job.getRetryCount() == null ? 0 : job.getRetryCount();
int nextRetryCount = Math.min(MAX_RETRY_COUNT - 1, retryCount + 1);
int updated = taskFileJobMapper.update(null, new LambdaUpdateWrapper<TaskFileJobEntity>()
.eq(TaskFileJobEntity::getId, job.getId())
.eq(TaskFileJobEntity::getStatus, "PENDING")
.eq(job.getUpdatedAt() != null, TaskFileJobEntity::getUpdatedAt, job.getUpdatedAt())
.isNull(job.getUpdatedAt() == null, TaskFileJobEntity::getUpdatedAt)
.set(TaskFileJobEntity::getRetryCount, nextRetryCount)
.set(TaskFileJobEntity::getUpdatedAt, LocalDateTime.now()));
if (updated > 0) {
publishDispatchEvent(taskFileJobMapper.selectById(job.getId()));
}
continue;
}
int retryCount = job.getRetryCount() == null ? 0 : job.getRetryCount(); int retryCount = job.getRetryCount() == null ? 0 : job.getRetryCount();
int nextRetryCount = Math.min(MAX_RETRY_COUNT, retryCount + 1); int nextRetryCount = Math.min(MAX_RETRY_COUNT, retryCount + 1);
LocalDateTime now = LocalDateTime.now(); LocalDateTime now = LocalDateTime.now();
@@ -642,6 +674,7 @@ public class TaskFileJobService {
entity.getResultId(), entity.getResultId(),
entity.getScopeKey(), entity.getScopeKey(),
entity.getJobType(), entity.getJobType(),
entity.getOwner(),
LocalDateTime.now() LocalDateTime.now()
)); ));
} }
@@ -1,6 +1,7 @@
package com.nanri.aiimage.modules.task.service; package com.nanri.aiimage.modules.task.service;
import com.nanri.aiimage.common.exception.BusinessException; import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.common.exception.TaskOwnerMismatchException;
import com.nanri.aiimage.modules.appearancepatent.service.AppearancePatentTaskService; import com.nanri.aiimage.modules.appearancepatent.service.AppearancePatentTaskService;
import com.nanri.aiimage.config.InstanceMetadata; import com.nanri.aiimage.config.InstanceMetadata;
import com.nanri.aiimage.modules.brand.service.BrandTaskService; import com.nanri.aiimage.modules.brand.service.BrandTaskService;
@@ -76,13 +77,16 @@ public class TaskResultFileJobWorker {
} }
log.info("[task-file-job] scheduled worker claimed jobs count={}", jobs.size()); log.info("[task-file-job] scheduled worker claimed jobs count={}", jobs.size());
for (TaskFileJobEntity job : jobs) { for (TaskFileJobEntity job : jobs) {
boolean dispatched = taskFileJobLocalDispatcher.dispatch( // claimRunnableJobsForOwner 已把 job 原子翻转成 RUNNING,这里必须走"已 claim"
// 通道:若再走 process() 的 claimRunning,条件 UPDATE 会因状态已是 RUNNING 而
// 失败,job 被静默丢弃并挂到 stuck 扫描才回收,定时兜底等于失效。
boolean dispatched = taskFileJobLocalDispatcher.dispatchClaimed(
job.getId(), job.getId(),
job.getTaskId(), job.getTaskId(),
job.getModuleType() job.getModuleType()
); );
if (!dispatched) { if (!dispatched) {
process(job); processClaimed(job);
} }
} }
} }
@@ -104,18 +108,39 @@ public class TaskResultFileJobWorker {
} }
public void process(TaskFileJobEntity job) { public void process(TaskFileJobEntity job) {
if (isOwnerScopedJob(job) && !isOwnedByCurrentInstance(job)) { if (job == null || job.getId() == null) {
log.debug("[task-file-job] skip owner-scoped job because owner is another instance jobId={} taskId={} moduleType={} owner={} current={}",
job.getId(), job.getTaskId(), job.getModuleType(), ownerFromScopeKey(job.getScopeKey()), currentInstanceId());
return; return;
} }
if (job == null || job.getId() == null) { if (isOwnerScopedJob(job) && !isOwnedByCurrentInstance(job)) {
// 非归属实例一律不碰 owner-scoped job:模块层 ensureTaskOwnedByCurrentInstance
// 会直接抛 TaskOwnerMismatchException,抢跑只会白白消耗 job 的重试次数,
// 重试耗尽后把归属机上跑得好好的任务判成失败。
// job 不会因此卡死:归属实例的 runPendingJobs 每 15 秒扫一次自己 owner 的
// PENDING/FAILED jobenqueueAssembleResult / stuck 扫描也会补发 dispatch。
log.info("[task-file-job] skip owner-scoped job because owner is another instance jobId={} taskId={} moduleType={} status={} owner={} current={}",
job.getId(), job.getTaskId(), job.getModuleType(), job.getStatus(),
ownerFromScopeKey(job.getScopeKey()), currentInstanceId());
return; return;
} }
TaskFileJobEntity claim = taskFileJobService.claimRunning(job.getId()); TaskFileJobEntity claim = taskFileJobService.claimRunning(job.getId());
if (claim == null) { if (claim == null) {
return; return;
} }
offloadOrRun(job, claim);
}
/**
* 供"调用方已持有 RUNNING claim"的路径使用(定时 worker 的 claimRunnableJobsForOwner):
* 跳过 claimRunning,直接按 handler 能力异步/内联执行。
*/
public void processClaimed(TaskFileJobEntity claim) {
if (claim == null || claim.getId() == null) {
return;
}
offloadOrRun(claim, claim);
}
private void offloadOrRun(TaskFileJobEntity job, TaskFileJobEntity claim) {
ResultFileJobHandler handler = handlerRegistry.asMap().get(job.getModuleType()); ResultFileJobHandler handler = handlerRegistry.asMap().get(job.getModuleType());
if (handler != null && handler.supportsAsyncOffload()) { if (handler != null && handler.supportsAsyncOffload()) {
try { try {
@@ -206,6 +231,12 @@ public class TaskResultFileJobWorker {
if (finalizeWithdraw) { if (finalizeWithdraw) {
withdrawTaskService.tryFinalizeTask(job.getTaskId(), false); withdrawTaskService.tryFinalizeTask(job.getTaskId(), false);
} }
} catch (TaskOwnerMismatchException ex) {
// 归属校验失败不是"生成失败":不能计入 retryCount,否则任务会在归属机
// 一切正常的情况下被本机重试耗尽判死。退回 PENDING 交还归属实例。
log.warn("[task-file-job] process rejected because task belongs to another instance jobId={} taskId={} moduleType={} owner={} current={}",
job.getId(), job.getTaskId(), job.getModuleType(), ex.getOwnerInstanceId(), ex.getCurrentInstanceId());
taskFileJobService.deferRunning(job.getId(), ex.getMessage());
} catch (Exception ex) { } catch (Exception ex) {
TaskFileJobEntity latest = taskFileJobService.findById(job.getId()); TaskFileJobEntity latest = taskFileJobService.findById(job.getId());
if (latest != null && "SUCCESS".equals(latest.getStatus())) { if (latest != null && "SUCCESS".equals(latest.getStatus())) {
@@ -0,0 +1,63 @@
-- V103: 后台管理菜单改为两级结构(一级分组 + 子页面)
-- 新增 4 个后台一级分组(不映射真实页面,route_path 为分组虚拟值),
-- 并把这 16 个后台页面菜单挂到对应分组下。
-- 分组只用于权限树层级与「上级菜单」候选;column_key / route_path 校验逻辑不变。
-- 对应 Java 端清单见 PermissionMenuSchemaInitializer.DEFAULT_ADMIN_GROUPS / DEFAULT_ADMIN_MENUS。
-- ---------- 1. 新增一级分组 ----------
INSERT INTO `columns` (`name`, `column_key`, `menu_type`, `route_path`, `parent_id`, `sort_order`)
SELECT '账号与权限', 'admin_group_account', 'admin', 'group-account', NULL, 5
WHERE NOT EXISTS (SELECT 1 FROM `columns` WHERE `column_key` = 'admin_group_account');
INSERT INTO `columns` (`name`, `column_key`, `menu_type`, `route_path`, `parent_id`, `sort_order`)
SELECT '数据管理', 'admin_group_data', 'admin', 'group-data', NULL, 15
WHERE NOT EXISTS (SELECT 1 FROM `columns` WHERE `column_key` = 'admin_group_data');
INSERT INTO `columns` (`name`, `column_key`, `menu_type`, `route_path`, `parent_id`, `sort_order`)
SELECT '店铺管理', 'admin_group_shop', 'admin', 'group-shop', NULL, 45
WHERE NOT EXISTS (SELECT 1 FROM `columns` WHERE `column_key` = 'admin_group_shop');
INSERT INTO `columns` (`name`, `column_key`, `menu_type`, `route_path`, `parent_id`, `sort_order`)
SELECT '记录与版本', 'admin_group_record', 'admin', 'group-record', NULL, 55
WHERE NOT EXISTS (SELECT 1 FROM `columns` WHERE `column_key` = 'admin_group_record');
-- ---------- 2. 页面菜单挂到分组下,并同步显示名 ----------
-- 只更新当前没有父级的页面(保留管理员手工调整过的层级),
-- COALESCE 保证已挂接的叶子不重复指向。
UPDATE `columns` child
JOIN `columns` parent ON parent.`column_key` = 'admin_group_account'
SET child.`parent_id` = COALESCE(child.`parent_id`, parent.`id`)
WHERE child.`column_key` IN ('admin_users', 'admin_columns', 'admin_group_manage');
UPDATE `columns` child
JOIN `columns` parent ON parent.`column_key` = 'admin_group_data'
SET child.`parent_id` = COALESCE(child.`parent_id`, parent.`id`)
WHERE child.`column_key` IN ('admin_dedupe_total_data', 'admin_invalid_asin_data',
'admin_query_asin', 'admin_product_categories');
UPDATE `columns` child
JOIN `columns` parent ON parent.`column_key` = 'admin_group_shop'
SET child.`parent_id` = COALESCE(child.`parent_id`, parent.`id`)
WHERE child.`column_key` IN ('admin_shop_keys', 'admin_shop_manage', 'admin_skip_price_asin',
'admin_shop_data_crawl_tasks');
UPDATE `columns` child
JOIN `columns` parent ON parent.`column_key` = 'admin_group_record'
SET child.`parent_id` = COALESCE(child.`parent_id`, parent.`id`)
WHERE child.`column_key` IN ('admin_history', 'admin_version', 'digital_human_version',
'admin_image_video_tasks');
-- ---------- 3. 显示名对齐前端界面 ---
UPDATE `columns` SET `name` = '品牌数据库'
WHERE `menu_type` = 'admin' AND `column_key` = 'admin_invalid_asin_data';
UPDATE `columns` SET `name` = '最低价ASIN设置'
WHERE `menu_type` = 'admin' AND `column_key` = 'admin_skip_price_asin';
UPDATE `columns` SET `name` = '视频任务记录'
WHERE `menu_type` = 'admin' AND `column_key` = 'admin_image_video_tasks';
UPDATE `columns` SET `name` = '店铺数据记录'
WHERE `menu_type` = 'admin' AND `column_key` = 'admin_shop_data_crawl_tasks';
@@ -106,8 +106,17 @@ class ExplainIndexAuditDocTest {
assertTrue(audit.contains("不做 DDL") || audit.contains("本任务无 DDL") assertTrue(audit.contains("不做 DDL") || audit.contains("本任务无 DDL")
|| audit.contains("不新增索引") || audit.contains("只审计"), || audit.contains("不新增索引") || audit.contains("只审计"),
"审计文档明确本任务不做 DDL"); "审计文档明确本任务不做 DDL");
assertFalse(audit.contains("ALTER TABLE"), // 任务 116 落地 V101 索引后,审计文档会把 ALTER TABLE 作为候选索引的
"审计文档不包含 ALTER TABLEDDL 属于后续任务 116"); // 回滚说明引用。审计任务自身依然不执行 DDL——判定口径改为「ALTER TABLE
// 只能出现在候选/回滚等说明性上下文」,而不是文档里不得出现该字符串。
for (String line : audit.split("\r?\n")) {
if (!line.contains("ALTER TABLE")) {
continue;
}
assertTrue(line.contains("回滚") || line.contains("候选") || line.contains("建议")
|| line.contains("风险") || line.trim().startsWith(">"),
"审计文档中的 ALTER TABLE 必须是候选/回滚说明,实际行:" + line);
}
assertTrue(spec.contains("只读") || spec.contains("巡检"), assertTrue(spec.contains("只读") || spec.contains("巡检"),
"12 spec 巡检报表模式与审计文档对应"); "12 spec 巡检报表模式与审计文档对应");
} }
@@ -27,6 +27,7 @@ import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.concurrent.ExecutorService; import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicInteger;
@@ -268,6 +269,7 @@ class ExternalCallMetricsRecorderTest {
} }
} finally { } finally {
pool.shutdown(); pool.shutdown();
assertTrue(pool.awaitTermination(60, TimeUnit.SECONDS), "并发任务在超时前全部结束");
} }
awaitMetric("aiimage.external-call.total", "client", "llm", "result", "success"); awaitMetric("aiimage.external-call.total", "client", "llm", "result", "success");
@@ -275,6 +277,12 @@ class ExternalCallMetricsRecorderTest {
Thread.sleep(10); Thread.sleep(10);
} }
assertEquals(20, llmSubmitCount.get(), "并发 20 请求全部真实发出"); assertEquals(20, llmSubmitCount.get(), "并发 20 请求全部真实发出");
// 指标在响应处理完成之后才落,submit 计数到 20 不代表指标已收敛:
// 直接断言会读到中间值(实测 18.0)。等指标计数自身收敛再断言。
for (int i = 0; i < 2000 && counterCount("aiimage.external-call.total",
"client", "llm", "result", "success") < 20.0; i++) {
Thread.sleep(10);
}
assertEquals(20.0, counterCount("aiimage.external-call.total", assertEquals(20.0, counterCount("aiimage.external-call.total",
"client", "llm", "result", "success"), "20 次成功全部记录,无重复"); "client", "llm", "result", "success"), "20 次成功全部记录,无重复");
} }
@@ -7,6 +7,7 @@ import com.nanri.aiimage.modules.collectdata.model.vo.CollectDataResultRowVo;
import org.apache.poi.ss.usermodel.Sheet; import org.apache.poi.ss.usermodel.Sheet;
import org.apache.poi.ss.usermodel.Workbook; import org.apache.poi.ss.usermodel.Workbook;
import org.apache.poi.ss.usermodel.WorkbookFactory; import org.apache.poi.ss.usermodel.WorkbookFactory;
import org.apache.poi.util.DefaultTempFileCreationStrategy;
import org.apache.poi.util.TempFile; import org.apache.poi.util.TempFile;
import org.apache.poi.util.TempFileCreationStrategy; import org.apache.poi.util.TempFileCreationStrategy;
import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.AfterEach;
@@ -39,20 +40,36 @@ class CollectDataExcelCleanupTest {
private final List<File> trackedTempFiles = new ArrayList<>(); private final List<File> trackedTempFiles = new ArrayList<>();
private Path sxssfTempDir;
@AfterEach @AfterEach
void restoreDefaultTempFileStrategy() {
// 还原 POI 的进程级全局策略:trackSxssfTempFiles 改写的是单例,
// 不还原会泄漏到同一 JVM 内后续所有用到 SXSSF 的测试。
// (原来这个 @AfterEach 挂在 private 的 trackSxssfTempFiles 上,
// 既没有还原策略,方法本身还是被各用例手工调用的 setup 辅助。)
TempFile.setTempFileCreationStrategy(new DefaultTempFileCreationStrategy());
}
private void trackSxssfTempFiles() { private void trackSxssfTempFiles() {
// 把 SXSSF 滚动窗口临时文件重定向到本用例独占目录,使残留探测只覆盖
// 本用例自己产生的文件。
try {
sxssfTempDir = Files.createDirectory(tempDir.resolve("sxssf-probe"));
} catch (IOException ex) {
throw new IllegalStateException(ex);
}
TempFile.setTempFileCreationStrategy(new TempFileCreationStrategy() { TempFile.setTempFileCreationStrategy(new TempFileCreationStrategy() {
@Override @Override
public File createTempFile(String prefix, String suffix) throws IOException { public File createTempFile(String prefix, String suffix) throws IOException {
File file = File.createTempFile(prefix, suffix); File file = Files.createTempFile(sxssfTempDir, prefix, suffix).toFile();
trackedTempFiles.add(file); trackedTempFiles.add(file);
return file; return file;
} }
@Override @Override
public File createTempDirectory(String prefix) throws IOException { public File createTempDirectory(String prefix) throws IOException {
File dir = Files.createTempDirectory(prefix).toFile(); File dir = Files.createTempDirectory(sxssfTempDir, prefix).toFile();
trackedTempFiles.add(dir); trackedTempFiles.add(dir);
return dir; return dir;
} }
@@ -62,10 +79,13 @@ class CollectDataExcelCleanupTest {
private void assertSxssfTempFilesCleaned() { private void assertSxssfTempFilesCleaned() {
assertThat(trackedTempFiles).as("SXSSF 滚动窗口临时文件均已删除") assertThat(trackedTempFiles).as("SXSSF 滚动窗口临时文件均已删除")
.allSatisfy(f -> assertThat(f).doesNotExist()); .allSatisfy(f -> assertThat(f).doesNotExist());
// POI 5.2.5 无法读取当前策略,改用系统临时目录中 SXSSF 固定前缀的残留探测。 // 只探测本用例独占目录。原实现扫描系统共享临时目录并要求 poi-sxssf*
try (Stream<Path> paths = Files.list(Path.of(System.getProperty("java.io.tmpdir")))) { // 前缀零命中,会把其他进程、其他用例和历史运行遗留的文件算成本用例
// 的泄漏——本机共享临时目录里就躺着十几个历史 poi-sxssf-sheet*.xml
// 导致实现完全正确时这 7 个用例照样红。
try (Stream<Path> paths = Files.list(sxssfTempDir)) {
assertThat(paths.filter(p -> p.getFileName().toString().startsWith("poi-sxssf"))) assertThat(paths.filter(p -> p.getFileName().toString().startsWith("poi-sxssf")))
.as("系统临时目录无 SXSSF 滚动窗口残留") .as("本用例 SXSSF 临时目录无滚动窗口残留")
.isEmpty(); .isEmpty();
} catch (IOException ex) { } catch (IOException ex) {
throw new IllegalStateException(ex); throw new IllegalStateException(ex);
@@ -149,6 +149,7 @@ class ShopDataCrawlCleanupTest {
private final AtomicLong memberIdSeq = new AtomicLong(1000); private final AtomicLong memberIdSeq = new AtomicLong(1000);
private final AtomicLong fileSeq = new AtomicLong(5000); private final AtomicLong fileSeq = new AtomicLong(5000);
private Long lastJobTaskId; private Long lastJobTaskId;
private Set<Path> preExistingWorkRoots = Set.of();
@BeforeEach @BeforeEach
void configureStorage() { void configureStorage() {
@@ -408,6 +409,12 @@ class ShopDataCrawlCleanupTest {
.deleteTaskItems(anyLong(), eq(MODULE_TYPE)); .deleteTaskItems(anyLong(), eq(MODULE_TYPE));
lenient().doNothing().when(taskProgressSnapshotService) lenient().doNothing().when(taskProgressSnapshotService)
.delete(anyLong(), eq(MODULE_TYPE)); .delete(anyLong(), eq(MODULE_TYPE));
// 共享临时目录里的既有 daily-* 目录不属于本用例,先快照、断言时做差集
try {
preExistingWorkRoots = Set.copyOf(scanWorkRoots());
} catch (Exception ex) {
throw new IllegalStateException(ex);
}
} }
// ---- 1. 删除路径 ---- // ---- 1. 删除路径 ----
@@ -1261,6 +1268,16 @@ class ShopDataCrawlCleanupTest {
} }
private List<Path> workRoots() throws Exception { private List<Path> workRoots() throws Exception {
// 只统计本用例新产生的工作目录。shop-data-crawl-result 落在机器共享的
// 系统临时目录下,其他用例、历史运行、本机跑过的应用都会往里留 daily-*
// 目录(本机实测已堆了 200+ 个),原实现把这些全算作本用例的泄漏,
// 于是实现完全正确时这些用例照样红。改为与用例开始时的快照做差集。
List<Path> current = new ArrayList<>(scanWorkRoots());
current.removeAll(preExistingWorkRoots);
return current;
}
private List<Path> scanWorkRoots() throws Exception {
Path root = Path.of(System.getProperty("java.io.tmpdir"), "shop-data-crawl-result"); Path root = Path.of(System.getProperty("java.io.tmpdir"), "shop-data-crawl-result");
if (!Files.exists(root)) { if (!Files.exists(root)) {
return List.of(); return List.of();
@@ -336,6 +336,56 @@ class SimilarAsinTaskServiceGroupRefTest {
assertEquals(80, coverage); assertEquals(80, coverage);
} }
@Test
void test_task_006_group_payload_embeds_items_for_python() throws Exception {
// Python 旧链路(normalize_groups 直接返回 groups 并按 gp["items"] 取行):
// 写入存储的 payload 中每组必须内嵌完整行,startIndex/endIndex 保留供 Java 内部展开。
java.util.concurrent.atomic.AtomicReference<String> captured = new java.util.concurrent.atomic.AtomicReference<>();
when(transientPayloadStorageService.storeParsedPayloadFast(
eq(SimilarAsinTaskService.MODULE_TYPE), any(), anyString(), anyString(), eq(false)))
.thenAnswer(invocation -> {
captured.set(invocation.getArgument(3));
return "rustfs:task-parsed/similar-asin/30000/payload.json";
});
File workbook = buildWorkbook(12);
SimilarAsinParseVo vo = parse(workbook, "uploads/20260829/gr-python.xlsx");
assertEquals(12, vo.getAcceptedRows());
SimilarAsinParsedPayloadDto stored = objectMapper.readValue(captured.get(), SimilarAsinParsedPayloadDto.class);
assertNotNull(stored.getGroups());
assertEquals(12, stored.getGroups().size());
// 每组内嵌 items 与引用区间宽度一致,且行内容对应 items[startIndex]
for (SimilarAsinParsedGroupVo group : stored.getGroups()) {
assertNotNull(group.getItems());
assertEquals(group.getEndIndex() - group.getStartIndex(), group.getItems().size(),
"Python 兼容:组内嵌 items 必须与引用区间宽度一致");
int start = group.getStartIndex();
assertEquals(stored.getItems().get(start).getRowToken(), group.getItems().get(0).getRowToken(),
"组 items 首行必须是 items[startIndex]");
}
// Java 内部展开路径不受内嵌 items 影响:仍按引用切片
List<SimilarAsinParsedRowVo> expanded = SimilarAsinTaskService.expandGroupRefs(stored);
assertEquals(12, expanded.size(), "Java 内部展开仍按引用,不受内嵌 items 影响");
}
@Test
void test_task_006_group_hydrate_embeds_items_for_legacy_payload() throws Exception {
// 存量任务 payload 只有 startIndex/endIndex 引用(无内嵌 items):
// parsedPayload 读取时 hydrate 必须回填组内嵌 itemsPython 才能取到行。
String legacy = "{\"items\":["
+ "{\"rowToken\":\"t1\",\"asin\":\"B0GRP00001\",\"sourceFileKey\":\"f.xlsx\",\"country\":\"英国\"},"
+ "{\"rowToken\":\"t2\",\"asin\":\"B0GRP00002\",\"sourceFileKey\":\"f.xlsx\",\"country\":\"英国\"}"
+ "],\"groups\":[{\"groupKey\":\"g1\",\"startIndex\":0,\"endIndex\":2,\"itemCount\":2}]}";
SimilarAsinParsedPayloadDto payload = objectMapper.readValue(legacy, SimilarAsinParsedPayloadDto.class);
SimilarAsinParsedPayloadDto hydrated = SimilarAsinTaskService.hydrateForPythonCompat(payload);
assertEquals(1, hydrated.getGroups().size());
assertNotNull(hydrated.getGroups().get(0).getItems());
assertEquals(2, hydrated.getGroups().get(0).getItems().size(), "hydrate 后组内嵌 items 必须为 2 行");
assertEquals("B0GRP00001", hydrated.getGroups().get(0).getItems().get(0).getAsin(),
"hydrate 后组内嵌 items 首行必须对应 items[0]");
assertEquals("B0GRP00002", hydrated.getGroups().get(0).getItems().get(1).getAsin(),
"hydrate 后组内嵌 items 第二行必须对应 items[1]");
}
// ---- helpers ---- // ---- helpers ----
private static String groupRefJson(int groupCount, int[][] ranges) { private static String groupRefJson(int groupCount, int[][] ranges) {
@@ -0,0 +1,85 @@
package com.nanri.aiimage.modules.task.service;
import com.nanri.aiimage.config.InstanceMetadata;
import com.nanri.aiimage.modules.task.model.dto.TaskFileJobDispatchEvent;
import com.nanri.aiimage.modules.task.model.dto.TaskFileJobMessage;
import org.junit.jupiter.api.Test;
import org.springframework.core.task.TaskExecutor;
import java.time.LocalDateTime;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyBoolean;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
/**
* owner-scoped job 的派发路由(2026-09-01 事故):两台实例共用一个 RocketMQ consumer group
* 消息被轮询投递到非归属实例后只能丢弃。归属实例改为直接本地派发,非归属实例不再发消息。
*/
class TaskFileJobDispatchCoordinatorOwnerRoutingTest {
private final TaskFileJobPublisher publisher = mock(TaskFileJobPublisher.class);
private final TaskFileJobLocalDispatcher localDispatcher = mock(TaskFileJobLocalDispatcher.class);
private final InstanceMetadata instanceMetadata = mock(InstanceMetadata.class);
private final TaskExecutor inlineExecutor = Runnable::run;
private TaskFileJobDispatchCoordinator buildCoordinator() {
return new TaskFileJobDispatchCoordinator(publisher, localDispatcher, instanceMetadata, inlineExecutor);
}
private static TaskFileJobDispatchEvent event(String owner) {
return new TaskFileJobDispatchEvent(18524L, 26402L, "SIMILAR_ASIN", 29405L,
owner == null ? "task:26402" : "task:26402:owner:" + owner,
TaskFileJobService.JOB_TYPE_ASSEMBLE_RESULT, owner, LocalDateTime.now());
}
@Test
void ownerJobDispatchedLocallyWithoutMq() {
when(instanceMetadata.getInstanceId()).thenReturn("server-121");
when(localDispatcher.dispatch(18524L, 26402L, "SIMILAR_ASIN", true)).thenReturn(true);
buildCoordinator().onTaskFileJobDispatch(event("server-121"));
verify(localDispatcher).dispatch(18524L, 26402L, "SIMILAR_ASIN", true);
verifyNoInteractions(publisher);
}
@Test
void ownerJobFallsBackToMqWhenLocalDispatchDeclined() {
when(instanceMetadata.getInstanceId()).thenReturn("server-121");
when(localDispatcher.dispatch(18524L, 26402L, "SIMILAR_ASIN", true)).thenReturn(false, true);
when(publisher.publish(any(TaskFileJobMessage.class))).thenReturn(false);
buildCoordinator().onTaskFileJobDispatch(event("server-121"));
verify(publisher).publish(any(TaskFileJobMessage.class));
}
@Test
void otherInstanceJobIsNeitherPublishedNorRun() {
when(instanceMetadata.getInstanceId()).thenReturn("server-110");
buildCoordinator().onTaskFileJobDispatch(event("server-121"));
verifyNoInteractions(publisher);
verify(localDispatcher, never()).dispatch(anyLong(), anyLong(), any(), anyBoolean());
verify(localDispatcher, never()).dispatch(anyLong(), anyLong(), any());
}
@Test
void jobWithoutOwnerStillGoesThroughMq() {
when(instanceMetadata.getInstanceId()).thenReturn("server-110");
when(publisher.publish(any(TaskFileJobMessage.class))).thenReturn(true);
buildCoordinator().onTaskFileJobDispatch(event(null));
verify(publisher).publish(any(TaskFileJobMessage.class));
verify(localDispatcher, never()).dispatch(anyLong(), anyLong(), eq("SIMILAR_ASIN"), anyBoolean());
}
}
@@ -0,0 +1,163 @@
package com.nanri.aiimage.modules.task.service;
import com.nanri.aiimage.common.exception.TaskOwnerMismatchException;
import com.nanri.aiimage.config.InstanceMetadata;
import com.nanri.aiimage.modules.appearancepatent.service.AppearancePatentTaskService;
import com.nanri.aiimage.modules.brand.service.BrandTaskService;
import com.nanri.aiimage.modules.collectdata.service.CollectDataService;
import com.nanri.aiimage.modules.deletebrand.service.DeleteBrandRunService;
import com.nanri.aiimage.modules.patroldelete.service.PatrolDeleteTaskService;
import com.nanri.aiimage.modules.pricetrack.service.PriceTrackTaskService;
import com.nanri.aiimage.modules.productrisk.service.ProductRiskTaskService;
import com.nanri.aiimage.modules.publish.service.PublishTaskService;
import com.nanri.aiimage.modules.queryasin.service.QueryAsinTaskService;
import com.nanri.aiimage.modules.shopdatacrawl.service.ShopDataCrawlTaskService;
import com.nanri.aiimage.modules.shopmatch.service.ShopMatchTaskService;
import com.nanri.aiimage.modules.similarasin.service.SimilarAsinTaskService;
import com.nanri.aiimage.modules.task.mapper.FileResultMapper;
import com.nanri.aiimage.modules.task.model.entity.TaskFileJobEntity;
import com.nanri.aiimage.modules.withdraw.service.WithdrawTaskService;
import org.junit.jupiter.api.Test;
import java.lang.reflect.Field;
import java.time.LocalDateTime;
import java.util.List;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
/**
* 双实例归属路由回归测试(2026-09-01 货源查询任务被非归属实例误杀):
* 1) 定时 worker 已 claim 的 job 不能再走 claimRunning,否则条件 UPDATE 失败被静默丢弃;
* 2) 处理中撞上归属校验时只退回 PENDING,不得计入 retryCountmarkFailed)。
*/
class TaskResultFileJobWorkerOwnerRoutingTest {
private final ShopMatchTaskService shopMatch = mock(ShopMatchTaskService.class);
private final PriceTrackTaskService priceTrack = mock(PriceTrackTaskService.class);
private final ProductRiskTaskService productRisk = mock(ProductRiskTaskService.class);
private final PublishTaskService publish = mock(PublishTaskService.class);
private final QueryAsinTaskService queryAsin = mock(QueryAsinTaskService.class);
private final ShopDataCrawlTaskService shopDataCrawl = mock(ShopDataCrawlTaskService.class);
private final WithdrawTaskService withdraw = mock(WithdrawTaskService.class);
private final PatrolDeleteTaskService patrolDelete = mock(PatrolDeleteTaskService.class);
private final AppearancePatentTaskService appearance = mock(AppearancePatentTaskService.class);
private final SimilarAsinTaskService similar = mock(SimilarAsinTaskService.class);
private final DeleteBrandRunService deleteBrand = mock(DeleteBrandRunService.class);
private final BrandTaskService brand = mock(BrandTaskService.class);
private final CollectDataService collectData = mock(CollectDataService.class);
private final TaskResultPayloadService payload = mock(TaskResultPayloadService.class);
private final TaskFileJobService taskFileJobService = mock(TaskFileJobService.class);
private final TaskDistributedLockService taskDistributedLockService = mock(TaskDistributedLockService.class);
private final TaskDistributedLockService.LockHandle lock = mock(TaskDistributedLockService.LockHandle.class);
private final TaskFileJobLocalDispatcher localDispatcher = mock(TaskFileJobLocalDispatcher.class);
private final InstanceMetadata instanceMetadata = mock(InstanceMetadata.class);
private TaskResultFileJobWorker buildWorker() throws Exception {
List<ResultFileJobHandler> handlers = List.of(
new ShopMatchResultFileJobHandler(shopMatch, payload),
new PriceTrackResultFileJobHandler(priceTrack, payload),
new ProductRiskResultFileJobHandler(productRisk, payload),
new PublishResultFileJobHandler(publish),
new QueryAsinResultFileJobHandler(queryAsin, payload),
new ShopDataCrawlResultFileJobHandler(shopDataCrawl, payload),
new WithdrawResultFileJobHandler(withdraw, payload),
new PatrolDeleteResultFileJobHandler(patrolDelete, payload),
new AppearancePatentResultFileJobHandler(appearance),
new SimilarAsinResultFileJobHandler(similar),
new DeleteBrandResultFileJobHandler(deleteBrand),
new BrandResultFileJobHandler(brand, payload),
new CollectDataResultFileJobHandler(collectData));
TaskResultFileJobWorker worker = new TaskResultFileJobWorker(
taskFileJobService,
taskDistributedLockService,
mock(FileResultMapper.class),
localDispatcher,
instanceMetadata,
withdraw, brand,
new ResultFileJobHandlerRegistry(handlers));
set(worker, "taskQueueExecutor", (org.springframework.core.task.TaskExecutor) Runnable::run);
set(worker, "localWorkerEnabled", true);
set(worker, "batchSize", 20);
return worker;
}
private static void set(TaskResultFileJobWorker worker, String fieldName, Object value) throws Exception {
Field field = TaskResultFileJobWorker.class.getDeclaredField(fieldName);
field.setAccessible(true);
field.set(worker, value);
}
private static TaskFileJobEntity claimedJob(String moduleType, long jobId, long taskId, String scopeKey) {
TaskFileJobEntity entity = new TaskFileJobEntity();
entity.setId(jobId);
entity.setTaskId(taskId);
entity.setResultId(taskId + 1000L);
entity.setModuleType(moduleType);
entity.setScopeKey(scopeKey);
entity.setStatus("RUNNING");
entity.setUpdatedAt(LocalDateTime.now());
return entity;
}
@Test
void scheduledWorkerProcessesAlreadyClaimedJobWithoutReclaiming() throws Exception {
TaskResultFileJobWorker worker = buildWorker();
when(instanceMetadata.getInstanceId()).thenReturn("instance-a");
TaskFileJobEntity claim = claimedJob(PublishTaskService.MODULE_TYPE, 18524L, 26402L, "task:26402:owner:instance-a");
when(taskFileJobService.claimRunnableJobsForOwner(20, "instance-a")).thenReturn(List.of(claim));
when(taskFileJobService.activateRunningClaim(claim)).thenReturn(true);
when(localDispatcher.dispatchClaimed(claim.getId(), claim.getTaskId(), claim.getModuleType())).thenReturn(false);
when(taskDistributedLockService.acquire(PublishTaskService.MODULE_TYPE, claim.getTaskId(),
TaskDistributedLockService.DEFAULT_WAIT_MILLIS)).thenReturn(lock);
worker.runPendingJobs();
verify(taskFileJobService, never()).claimRunning(anyLong());
verify(publish).processResultFileJob(claim);
verify(taskFileJobService).markSuccess(eq(claim), any());
}
@Test
void scheduledWorkerPrefersClaimedDispatch() throws Exception {
TaskResultFileJobWorker worker = buildWorker();
when(instanceMetadata.getInstanceId()).thenReturn("instance-a");
TaskFileJobEntity claim = claimedJob(PublishTaskService.MODULE_TYPE, 18525L, 26403L, "task:26403:owner:instance-a");
when(taskFileJobService.claimRunnableJobsForOwner(20, "instance-a")).thenReturn(List.of(claim));
when(localDispatcher.dispatchClaimed(claim.getId(), claim.getTaskId(), claim.getModuleType())).thenReturn(true);
worker.runPendingJobs();
verify(localDispatcher).dispatchClaimed(claim.getId(), claim.getTaskId(), claim.getModuleType());
verify(localDispatcher, never()).dispatch(anyLong(), anyLong(), any());
verify(taskFileJobService, never()).claimRunning(anyLong());
verify(publish, never()).processResultFileJob(any());
}
@Test
void ownerMismatchDuringProcessingDefersInsteadOfFailing() throws Exception {
TaskResultFileJobWorker worker = buildWorker();
when(instanceMetadata.getInstanceId()).thenReturn("instance-a");
// scopeKey 归属本机(通过 process 的前置校验),但模块层读到的 task owner 已经变成他机
TaskFileJobEntity job = claimedJob(PublishTaskService.MODULE_TYPE, 18526L, 26404L, "task:26404:owner:instance-a");
when(taskFileJobService.claimRunning(job.getId())).thenReturn(job);
when(taskFileJobService.activateRunningClaim(job)).thenReturn(true);
when(taskDistributedLockService.acquire(PublishTaskService.MODULE_TYPE, job.getTaskId(),
TaskDistributedLockService.DEFAULT_WAIT_MILLIS)).thenReturn(lock);
doThrow(new TaskOwnerMismatchException(job.getTaskId(), "assemble result file", "instance-b", "instance-a"))
.when(publish).processResultFileJob(job);
worker.process(job);
verify(taskFileJobService).deferRunning(eq(job.getId()), any());
verify(taskFileJobService, never()).markFailed(any(), any());
verify(taskFileJobService, never()).markFailedPermanent(any(), any());
verify(taskFileJobService, never()).isRetryExhausted(anyLong());
}
}
@@ -128,6 +128,20 @@ class TaskResultFileJobWorkerOwnerScopedTest {
TaskResultFileJobWorker worker = buildWorker(); TaskResultFileJobWorker worker = buildWorker();
when(instanceMetadata.getInstanceId()).thenReturn("instance-a"); when(instanceMetadata.getInstanceId()).thenReturn("instance-a");
TaskFileJobEntity job = job("PUBLISH", 2L, 12L, "task:12:owner:instance-b"); TaskFileJobEntity job = job("PUBLISH", 2L, 12L, "task:12:owner:instance-b");
job.setStatus("RUNNING");
worker.process(job);
verifyNoInteractions(taskFileJobService, taskDistributedLockService, publish);
}
@Test
void ownerScopedOtherOwnerPendingSkip() {
// PENDING + owner 不匹配也必须跳过:模块层 ensureTaskOwnedByCurrentInstance 会抛
// TaskOwnerMismatchException,本机抢跑只会耗尽 job 重试次数并误杀归属机上的任务。
// job 由归属实例的 runPendingJobs(每 15 秒扫自己 owner 的 PENDING/FAILED)接手。
TaskResultFileJobWorker worker = buildWorker();
when(instanceMetadata.getInstanceId()).thenReturn("instance-a");
TaskFileJobEntity job = job("PUBLISH", 2L, 12L, "task:12:owner:instance-b");
job.setStatus("PENDING");
worker.process(job); worker.process(job);
verifyNoInteractions(taskFileJobService, taskDistributedLockService, publish); verifyNoInteractions(taskFileJobService, taskDistributedLockService, publish);
} }