task-76: JSON owner 查询迁移到显式列并补充任务/状态复合索引

This commit is contained in:
2026-08-30 20:51:07 +08:00
parent e9ba3689ba
commit 04f45bebec
4 changed files with 324 additions and 10 deletions
@@ -17,6 +17,8 @@ public class TaskFileJobEntity {
private String moduleType;
private Long resultId;
private String scopeKey;
/** 从 scope_key 提取的显式 ownertask:1:owner:instance-a → instance-a),支持等值索引查询。 */
private String owner;
private String jobType;
private String status;
private Integer retryCount;
@@ -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<TaskFileJobEntity>()
.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<TaskFileJobEntity> ownerJobs = taskFileJobMapper.selectList(new LambdaQueryWrapper<TaskFileJobEntity>()
.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<TaskFileJobEntity> ownerJobs = taskFileJobMapper.selectList(new LambdaQueryWrapper<TaskFileJobEntity>()
.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_keytask: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;
@@ -0,0 +1,41 @@
-- Task 76biz_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;
@@ -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<TaskFileJobEntity> ownerJobs, List<TaskFileJobEntity> 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<TaskFileJobEntity> 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<TaskFileJobEntity> 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<TaskFileJobEntity> ownerJobs = List.of(ownerJob(1L, "instance-a"));
stubOwnerQueries(ownerJobs, List.of());
stubClaimResults(1);
stubClaimedRows(List.of(runningOwnerJob(1L, "instance-a")));
List<TaskFileJobEntity> claimed = service.claimRunnableJobsForOwner(10, "instance-a");
assertEquals(1, claimed.size(), "owner 任务按等值查询返回");
ArgumentCaptor<LambdaQueryWrapper<TaskFileJobEntity>> 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<TaskFileJobEntity> 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<TaskFileJobEntity> claimed = service.claimRunnableJobsForOwner(10, "instance-a");
assertEquals(List.of(1L, 2L, 9L), claimed.stream().map(TaskFileJobEntity::getId).toList(),
"owner 任务在前、通用任务补足,顺序稳定");
ArgumentCaptor<LambdaQueryWrapper<TaskFileJobEntity>> 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 不重复插入,
// 复用既有 jobFAILED 重置路径同步刷新 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<LambdaUpdateWrapper<TaskFileJobEntity>> 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<TaskFileJobEntity> claimed = service.claimRunnableJobsForOwner(10, " ");
assertTrue(claimed.isEmpty());
ArgumentCaptor<LambdaQueryWrapper<TaskFileJobEntity>> 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<TaskFileJobEntity> 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<LambdaQueryWrapper<TaskFileJobEntity>> 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));
}
}