From 04f45bebec68aaa0e2b40b81cc42e5c05b4a20db Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=BB=84=E8=87=AA=E8=BE=BE?= <980324341@qq.com> Date: Sun, 30 Aug 2026 20:51:07 +0800 Subject: [PATCH] =?UTF-8?q?task-76:=20JSON=20owner=20=E6=9F=A5=E8=AF=A2?= =?UTF-8?q?=E8=BF=81=E7=A7=BB=E5=88=B0=E6=98=BE=E5=BC=8F=E5=88=97=E5=B9=B6?= =?UTF-8?q?=E8=A1=A5=E5=85=85=E4=BB=BB=E5=8A=A1/=E7=8A=B6=E6=80=81?= =?UTF-8?q?=E5=A4=8D=E5=90=88=E7=B4=A2=E5=BC=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../task/model/entity/TaskFileJobEntity.java | 2 + .../task/service/TaskFileJobService.java | 31 ++- .../V96__task_file_job_owner_column_index.sql | 41 +++ .../service/TaskFileJobOwnerColumnTest.java | 260 ++++++++++++++++++ 4 files changed, 324 insertions(+), 10 deletions(-) create mode 100644 backend-java/src/main/resources/db/V96__task_file_job_owner_column_index.sql create mode 100644 backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskFileJobOwnerColumnTest.java diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/task/model/entity/TaskFileJobEntity.java b/backend-java/src/main/java/com/nanri/aiimage/modules/task/model/entity/TaskFileJobEntity.java index 8a6a999b..45853762 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/task/model/entity/TaskFileJobEntity.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/task/model/entity/TaskFileJobEntity.java @@ -17,6 +17,8 @@ public class TaskFileJobEntity { private String moduleType; private Long resultId; private String scopeKey; + /** 从 scope_key 提取的显式 owner(task:1:owner:instance-a → instance-a),支持等值索引查询。 */ + private String owner; private String jobType; private String status; private Integer retryCount; diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskFileJobService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskFileJobService.java index edb1f010..1f26a1b5 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskFileJobService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskFileJobService.java @@ -40,6 +40,7 @@ public class TaskFileJobService { entity.setModuleType(moduleType); entity.setResultId(resultId); entity.setScopeKey(scopeKey); + entity.setOwner(ownerFromScopeKey(scopeKey)); entity.setJobType(JOB_TYPE_ASSEMBLE_RESULT); entity.setStatus("PENDING"); entity.setRetryCount(0); @@ -68,6 +69,7 @@ public class TaskFileJobService { taskFileJobMapper.update(null, new LambdaUpdateWrapper() .eq(TaskFileJobEntity::getId, existing.getId()) .set(TaskFileJobEntity::getScopeKey, scopeKey) + .set(TaskFileJobEntity::getOwner, ownerFromScopeKey(scopeKey)) .set(TaskFileJobEntity::getStatus, "PENDING") .set(TaskFileJobEntity::getErrorMessage, null) .set(TaskFileJobEntity::getUpdatedAt, now) @@ -103,14 +105,13 @@ public class TaskFileJobService { if (normalizedOwner.isBlank()) { return claimRunnableJobs(limit); } - String ownerMarker = ":owner:" + normalizedOwner; int safeLimit = Math.max(1, Math.min(limit, 100)); List ownerJobs = taskFileJobMapper.selectList(new LambdaQueryWrapper() .in(TaskFileJobEntity::getStatus, List.of("PENDING", "FAILED")) .lt(TaskFileJobEntity::getRetryCount, MAX_RETRY_COUNT) .in(TaskFileJobEntity::getModuleType, List.of("APPEARANCE_PATENT", "SIMILAR_ASIN", "PUBLISH", "SHOP_DATA_CRAWL")) - .like(TaskFileJobEntity::getScopeKey, ownerMarker) + .eq(TaskFileJobEntity::getOwner, normalizedOwner) .orderByAsc(TaskFileJobEntity::getUpdatedAt) .last("limit " + safeLimit)); @@ -122,9 +123,7 @@ public class TaskFileJobService { .and(wrapper -> wrapper .notIn(TaskFileJobEntity::getModuleType, List.of("APPEARANCE_PATENT", "SIMILAR_ASIN", "PUBLISH", "SHOP_DATA_CRAWL")) .or() - .isNull(TaskFileJobEntity::getScopeKey) - .or() - .notLike(TaskFileJobEntity::getScopeKey, ":owner:")) + .isNull(TaskFileJobEntity::getOwner)) .orderByAsc(TaskFileJobEntity::getUpdatedAt) .last("limit " + (safeLimit - ownerJobs.size())); candidates.addAll(taskFileJobMapper.selectList(genericWrapper)); @@ -157,14 +156,13 @@ public class TaskFileJobService { if (normalizedOwner.isBlank()) { return listRunnableJobs(limit); } - String ownerMarker = ":owner:" + normalizedOwner; int safeLimit = Math.max(1, Math.min(limit, 100)); List ownerJobs = taskFileJobMapper.selectList(new LambdaQueryWrapper() .in(TaskFileJobEntity::getStatus, List.of("PENDING", "FAILED")) .lt(TaskFileJobEntity::getRetryCount, MAX_RETRY_COUNT) .in(TaskFileJobEntity::getModuleType, List.of("APPEARANCE_PATENT", "SIMILAR_ASIN", "PUBLISH", "SHOP_DATA_CRAWL")) - .like(TaskFileJobEntity::getScopeKey, ownerMarker) + .eq(TaskFileJobEntity::getOwner, normalizedOwner) .orderByAsc(TaskFileJobEntity::getUpdatedAt) .last("limit " + safeLimit)); if (ownerJobs.size() >= safeLimit) { @@ -178,9 +176,7 @@ public class TaskFileJobService { .and(wrapper -> wrapper .notIn(TaskFileJobEntity::getModuleType, List.of("APPEARANCE_PATENT", "SIMILAR_ASIN", "PUBLISH", "SHOP_DATA_CRAWL")) .or() - .isNull(TaskFileJobEntity::getScopeKey) - .or() - .notLike(TaskFileJobEntity::getScopeKey, ":owner:")) + .isNull(TaskFileJobEntity::getOwner)) .orderByAsc(TaskFileJobEntity::getUpdatedAt) .last("limit " + (safeLimit - jobs.size())); jobs.addAll(taskFileJobMapper.selectList(genericWrapper)); @@ -586,6 +582,21 @@ public class TaskFileJobService { .eq(TaskFileJobEntity::getResultId, resultId)); } + private static final String OWNER_MARKER = ":owner:"; + + /** 从 scope_key(task:1:owner:instance-a)提取 owner;无标记或空白返回 null。 */ + static String ownerFromScopeKey(String scopeKey) { + if (scopeKey == null || scopeKey.isBlank()) { + return null; + } + int index = scopeKey.lastIndexOf(OWNER_MARKER); + if (index < 0) { + return null; + } + String owner = scopeKey.substring(index + OWNER_MARKER.length()).trim(); + return owner.isBlank() ? null : owner; + } + private TaskFileJobEntity findJob(Long taskId, String moduleType, Long resultId, String jobType) { if (taskId == null || moduleType == null || moduleType.isBlank() || resultId == null || jobType == null || jobType.isBlank()) { return null; diff --git a/backend-java/src/main/resources/db/V96__task_file_job_owner_column_index.sql b/backend-java/src/main/resources/db/V96__task_file_job_owner_column_index.sql new file mode 100644 index 00000000..6a43f4dd --- /dev/null +++ b/backend-java/src/main/resources/db/V96__task_file_job_owner_column_index.sql @@ -0,0 +1,41 @@ +-- Task 76:biz_task_file_job 增加显式 owner 列并补充 (module_type, owner, status, retry_count) 复合索引。 +-- owner 从 scope_key 的 ":owner:" 标记提取(task:1:owner:instance-a → instance-a), +-- 存量行回填,owner 等值查询不再依赖 LIKE '%:owner:%' 扫描。 + +SET @db_name = DATABASE(); + +SET @col_exists := ( + SELECT COUNT(*) + FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = @db_name + AND TABLE_NAME = 'biz_task_file_job' + AND COLUMN_NAME = 'owner' +); +SET @sql := IF(@col_exists = 0, + 'ALTER TABLE biz_task_file_job ADD COLUMN owner VARCHAR(128) NULL COMMENT ''explicit owner extracted from scope_key'' AFTER scope_key', + 'SELECT 1' +); +PREPARE stmt FROM @sql; +EXECUTE stmt; +DEALLOCATE PREPARE stmt; + +UPDATE biz_task_file_job +SET owner = SUBSTRING(scope_key, LOCATE(':owner:', scope_key) + 7) +WHERE owner IS NULL + AND scope_key IS NOT NULL + AND LOCATE(':owner:', scope_key) > 0; + +SET @idx_exists := ( + SELECT COUNT(*) + FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = @db_name + AND TABLE_NAME = 'biz_task_file_job' + AND INDEX_NAME = 'idx_file_job_module_owner_status_retry' +); +SET @sql := IF(@idx_exists = 0, + 'ALTER TABLE biz_task_file_job ADD INDEX idx_file_job_module_owner_status_retry (module_type, owner, status, retry_count)', + 'SELECT 1' +); +PREPARE stmt FROM @sql; +EXECUTE stmt; +DEALLOCATE PREPARE stmt; diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskFileJobOwnerColumnTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskFileJobOwnerColumnTest.java new file mode 100644 index 00000000..4c703b62 --- /dev/null +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskFileJobOwnerColumnTest.java @@ -0,0 +1,260 @@ +package com.nanri.aiimage.modules.task.service; + +import com.baomidou.mybatisplus.core.MybatisConfiguration; +import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; +import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper; +import com.baomidou.mybatisplus.core.metadata.TableInfoHelper; +import com.nanri.aiimage.modules.task.mapper.TaskFileJobMapper; +import com.nanri.aiimage.modules.task.model.dto.TaskFileJobDispatchEvent; +import com.nanri.aiimage.modules.task.model.entity.TaskFileJobEntity; +import org.apache.ibatis.builder.MapperBuilderAssistant; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.context.ApplicationEventPublisher; +import org.springframework.dao.DuplicateKeyException; + +import java.time.LocalDateTime; +import java.util.ArrayList; +import java.util.List; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.isNull; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** + * Task 76:将 JSON owner 查询迁移到显式列并补充任务/状态复合索引。 + * biz_task_file_job 新增 owner 列(enqueue 时从 scope_key 提取、V96 迁移 + * 回填存量行),owner 任务查询从 LIKE '%:owner:%' 改为 owner = ? 等值查询, + * 新增 (module_type, owner, status, retry_count) 复合索引;无 owner 的通用 + * 任务改为 owner IS NULL 过滤,语义与迁移前一致。 + */ +@ExtendWith(MockitoExtension.class) +class TaskFileJobOwnerColumnTest { + + @Mock private TaskFileJobMapper taskFileJobMapper; + @Mock private ApplicationEventPublisher applicationEventPublisher; + + private TaskFileJobService service; + + @BeforeAll + static void initializeTableInfo() { + TableInfoHelper.initTableInfo( + new MapperBuilderAssistant(new MybatisConfiguration(), ""), + TaskFileJobEntity.class); + } + + @BeforeEach + void setUp() { + service = new TaskFileJobService(taskFileJobMapper, applicationEventPublisher); + } + + private static TaskFileJobEntity ownerJob(Long id, String owner) { + TaskFileJobEntity job = new TaskFileJobEntity(); + job.setId(id); + job.setTaskId(1000L + id); + job.setModuleType("SIMILAR_ASIN"); + job.setScopeKey("task:" + id + ":owner:" + owner); + job.setOwner(owner); + job.setJobType("ASSEMBLE_RESULT"); + job.setStatus("PENDING"); + job.setRetryCount(0); + job.setUpdatedAt(LocalDateTime.now()); + return job; + } + + private static TaskFileJobEntity runningOwnerJob(Long id, String owner) { + TaskFileJobEntity job = ownerJob(id, owner); + job.setStatus("RUNNING"); + return job; + } + + private void stubOwnerQueries(List ownerJobs, List genericJobs) { + when(taskFileJobMapper.selectList(any(LambdaQueryWrapper.class))) + .thenReturn(ownerJobs, genericJobs); + } + + private void stubClaimResults(int... results) { + when(taskFileJobMapper.update(isNull(), any(LambdaUpdateWrapper.class))) + .thenReturn(results[0], java.util.Arrays.stream(results, 1, results.length).boxed().toArray(Integer[]::new)); + } + + private void stubClaimedRows(List claimed) { + when(taskFileJobMapper.selectById(any())).thenAnswer(invocation -> { + Long id = invocation.getArgument(0); + return claimed.stream().filter(job -> job.getId().equals(id)).findFirst().orElse(null); + }); + } + + private TaskFileJobEntity capturedInsert() { + ArgumentCaptor captor = ArgumentCaptor.forClass(TaskFileJobEntity.class); + verify(taskFileJobMapper).insert(captor.capture()); + return captor.getValue(); + } + + @Test + void test_task_076_owner_normal_default_path() { + // 默认路径:enqueue 时从 scope_key 提取 owner 落显式列; + // owner 查询使用 owner = ? 等值条件,不再使用 LIKE。 + TaskFileJobEntity enqueued = service.enqueueAssembleResult(1L, "SIMILAR_ASIN", 2L, "task:1:owner:instance-a"); + assertEquals("instance-a", enqueued.getOwner(), "enqueue 时提取 owner"); + assertEquals("task:1:owner:instance-a", enqueued.getScopeKey(), "scope_key 原样保留"); + assertEquals("instance-a", capturedInsert().getOwner()); + + List ownerJobs = List.of(ownerJob(1L, "instance-a")); + stubOwnerQueries(ownerJobs, List.of()); + stubClaimResults(1); + stubClaimedRows(List.of(runningOwnerJob(1L, "instance-a"))); + + List claimed = service.claimRunnableJobsForOwner(10, "instance-a"); + assertEquals(1, claimed.size(), "owner 任务按等值查询返回"); + + ArgumentCaptor> captor = + ArgumentCaptor.forClass(LambdaQueryWrapper.class); + verify(taskFileJobMapper, times(2)).selectList(captor.capture()); + String ownerSegment = captor.getAllValues().get(0).getCustomSqlSegment(); + assertTrue(ownerSegment.contains("owner = "), "owner 等值查询: " + ownerSegment); + assertFalse(ownerSegment.contains("LIKE"), "不再使用 LIKE 模糊匹配: " + ownerSegment); + } + + @Test + void test_task_076_owner_normal_multiple_items() { + // 批量场景:多个 owner 的任务,按 owner 精确过滤只返回归属任务, + // 通用任务补足仍保留(owner IS NULL),顺序稳定不丢失。 + List ownerJobs = List.of(ownerJob(1L, "instance-a"), ownerJob(2L, "instance-a")); + stubOwnerQueries(ownerJobs, List.of(ownerJob(9L, null))); + stubClaimResults(1, 1, 1); + stubClaimedRows(List.of(runningOwnerJob(1L, "instance-a"), runningOwnerJob(2L, "instance-a"), runningOwnerJob(9L, null))); + + List claimed = service.claimRunnableJobsForOwner(10, "instance-a"); + + assertEquals(List.of(1L, 2L, 9L), claimed.stream().map(TaskFileJobEntity::getId).toList(), + "owner 任务在前、通用任务补足,顺序稳定"); + ArgumentCaptor> captor = + ArgumentCaptor.forClass(LambdaQueryWrapper.class); + verify(taskFileJobMapper, times(2)).selectList(captor.capture()); + assertTrue(captor.getAllValues().get(0).getCustomSqlSegment().contains("owner = "), "owner 精确过滤"); + assertTrue(captor.getAllValues().get(1).getCustomSqlSegment().contains("owner IS NULL"), + "通用任务按 owner IS NULL 过滤"); + } + + @Test + void test_task_076_owner_normal_repeated_operation_is_idempotent() { + // 幂等:同一 (taskId, moduleType, resultId) 重复 enqueue 不重复插入, + // 复用既有 job;FAILED 重置路径同步刷新 owner,不产生重复记录。 + TaskFileJobEntity existing = ownerJob(5L, "instance-a"); + existing.setStatus("FAILED"); + existing.setRetryCount(2); + when(taskFileJobMapper.selectOne(any(LambdaQueryWrapper.class))).thenReturn(existing); + when(taskFileJobMapper.selectById(5L)).thenReturn(existing); + + TaskFileJobEntity result = service.enqueueAssembleResult(5L, "SIMILAR_ASIN", 50L, "task:5:owner:instance-a"); + + assertEquals(5L, result.getId(), "复用既有 job"); + verify(taskFileJobMapper, never()).insert(any(TaskFileJobEntity.class)); + ArgumentCaptor> updateCaptor = + ArgumentCaptor.forClass(LambdaUpdateWrapper.class); + verify(taskFileJobMapper, times(1)).update(isNull(), updateCaptor.capture()); + assertTrue(updateCaptor.getValue().getSqlSet().contains("owner"), "重置路径同步刷新 owner 列"); + } + + @Test + void test_task_076_owner_boundary_empty_input() { + // 空输入:scope_key 为空时 owner 为空;空白 owner 走通用 claim 路径 + // (不发起 owner 等值查询)。 + TaskFileJobEntity enqueued = service.enqueueAssembleResult(1L, "SIMILAR_ASIN", 2L, " "); + assertNull(enqueued.getOwner(), "空 scope_key 不提取 owner"); + assertNull(capturedInsert().getOwner()); + + when(taskFileJobMapper.selectList(any(LambdaQueryWrapper.class))).thenReturn(List.of()); + List claimed = service.claimRunnableJobsForOwner(10, " "); + + assertTrue(claimed.isEmpty()); + ArgumentCaptor> captor = + ArgumentCaptor.forClass(LambdaQueryWrapper.class); + verify(taskFileJobMapper, times(1)).selectList(captor.capture()); + assertFalse(captor.getValue().getCustomSqlSegment().contains("owner"), "空白 owner 不发起 owner 查询"); + } + + @Test + void test_task_076_owner_boundary_single_item() { + // 单元素:单任务单 owner 直接返回,不依赖批量路径。 + stubOwnerQueries(List.of(ownerJob(9L, "instance-a")), List.of()); + stubClaimResults(1); + stubClaimedRows(List.of(runningOwnerJob(9L, "instance-a"))); + + List claimed = service.claimRunnableJobsForOwner(10, "instance-a"); + + assertEquals(List.of(9L), claimed.stream().map(TaskFileJobEntity::getId).toList()); + } + + @Test + void test_task_076_owner_boundary_limit_and_overflow() { + // 上限/超限:limit 非法(0/负/超上限)按现有 clamp 回退,不发生无界取数。 + stubOwnerQueries(List.of(), List.of()); + + service.claimRunnableJobsForOwner(0, "instance-a"); + service.claimRunnableJobsForOwner(200, "instance-a"); + + ArgumentCaptor> captor = + ArgumentCaptor.forClass(LambdaQueryWrapper.class); + verify(taskFileJobMapper, times(4)).selectList(captor.capture()); + assertTrue(captor.getAllValues().get(0).getCustomSqlSegment().contains("limit 1"), "limit=0 回退到 1"); + assertTrue(captor.getAllValues().get(2).getCustomSqlSegment().contains("limit 100"), "limit=200 钳制到 100"); + } + + @Test + void test_task_076_owner_invalid_input_rejected() { + // 非法参数:必填字段缺失时 enqueue 拒绝返回 null 不落库; + // scope_key 不含 owner 标记时 owner 为空但仍可入队。 + assertNull(service.enqueueAssembleResult(null, "SIMILAR_ASIN", 2L, "task:1:owner:instance-a"), + "null taskId 拒绝"); + assertNull(service.enqueueAssembleResult(1L, " ", 2L, "task:1:owner:instance-a"), + "空白 moduleType 拒绝"); + verify(taskFileJobMapper, never()).insert(any(TaskFileJobEntity.class)); + + TaskFileJobEntity noMarker = service.enqueueAssembleResult(1L, "SIMILAR_ASIN", 2L, "task:1:plain"); + assertNull(noMarker.getOwner(), "无 owner 标记不提取 owner"); + assertEquals("task:1:plain", noMarker.getScopeKey()); + } + + @Test + void test_task_076_owner_dependency_failure_releases_resources() { + // 依赖失败:insert 主键冲突(DuplicateKeyException)时回退到既有 FAILED job, + // 走重置路径刷新 owner 并继续发布派发事件;owner 查询抛异常时传播且不产生 claim 翻转。 + TaskFileJobEntity existing = ownerJob(7L, "instance-a"); + existing.setStatus("FAILED"); + existing.setRetryCount(2); + // 首次查询无既有 job → 尝试 insert 触发主键冲突(并发提交), + // 冲突后重新查询命中既有 FAILED job,走重置路径刷新 owner。 + when(taskFileJobMapper.selectOne(any(LambdaQueryWrapper.class))) + .thenReturn(null, existing); + when(taskFileJobMapper.insert(any(TaskFileJobEntity.class))) + .thenThrow(new DuplicateKeyException("duplicate")); + when(taskFileJobMapper.update(isNull(), any(LambdaUpdateWrapper.class))).thenReturn(1); + when(taskFileJobMapper.selectById(7L)).thenReturn(existing); + + TaskFileJobEntity recovered = service.enqueueAssembleResult(7L, "SIMILAR_ASIN", 70L, "task:7:owner:instance-a"); + + assertEquals(7L, recovered.getId(), "冲突后回退到既有 job"); + verify(applicationEventPublisher, times(1)).publishEvent(any(TaskFileJobDispatchEvent.class)); + + when(taskFileJobMapper.selectList(any(LambdaQueryWrapper.class))) + .thenThrow(new RuntimeException("db down")); + assertThrows(RuntimeException.class, () -> service.claimRunnableJobsForOwner(10, "instance-a")); + // 唯一的 update 来自 enqueue 重置路径;claim 在查询阶段中止,未发起任何原子翻转。 + verify(taskFileJobMapper, times(1)).update(isNull(), any(LambdaUpdateWrapper.class)); + } +}