diff --git a/backend-java/pom.xml b/backend-java/pom.xml index acf726ba..a57f4908 100644 --- a/backend-java/pom.xml +++ b/backend-java/pom.xml @@ -183,7 +183,7 @@ org.apache.maven.plugins maven-surefire-plugin - -XX:+EnableDynamicAgentLoading -Xshare:off + -XX:+EnableDynamicAgentLoading -Xshare:off -Xmx1536m diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java index 02f82b26..a22be146 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java @@ -1469,6 +1469,46 @@ public class SimilarAsinTaskService { return result; } + /** + * Task 9:按 id keyset 批量分页拉取 chunk,保持低内存读取。 + * 每轮只取 pageSize 条(id > lastId 升序),全部取回后按 chunkIndex 升序合并。 + * taskId 为 null 安全返回空;pageSize <= 0 回退默认 500; + * 中途查询失败抛项目约定异常,不返回半截结果。 + * 使用 QueryWrapper(列名直写)避免对 lambda 元数据缓存的依赖。 + */ + static List loadChunksKeyset(TaskChunkMapper mapper, Long taskId, String moduleType, int pageSize) { + if (taskId == null || mapper == null) { + return List.of(); + } + int batch = pageSize > 0 ? pageSize : 500; + List all = new ArrayList<>(); + long lastId = 0L; + while (true) { + List page; + try { + page = mapper.selectList(new com.baomidou.mybatisplus.core.conditions.query.QueryWrapper() + .eq("task_id", taskId) + .eq("module_type", moduleType) + .gt("id", lastId) + .orderByAsc("id") + .last("limit " + batch)); + } catch (Exception ex) { + throw new BusinessException("chunk keyset 分页查询失败: " + ex.getMessage(), ex); + } + if (page == null || page.isEmpty()) { + break; + } + all.addAll(page); + lastId = page.getLast().getId(); + if (page.size() < batch) { + break; + } + } + all.sort(java.util.Comparator.comparing(TaskChunkEntity::getChunkIndex, + java.util.Comparator.nullsLast(java.util.Comparator.naturalOrder()))); + return all; + } + private void applyCozeToPersistedChunks(FileTaskEntity task, Runnable progressHook) { if (task == null || task.getId() == null) { return; diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceChunkKeysetTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceChunkKeysetTest.java new file mode 100644 index 00000000..9d7cf142 --- /dev/null +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceChunkKeysetTest.java @@ -0,0 +1,204 @@ +package com.nanri.aiimage.modules.similarasin.service; + +import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper; +import com.nanri.aiimage.common.exception.BusinessException; +import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper; +import com.nanri.aiimage.modules.task.model.entity.TaskChunkEntity; +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 java.util.ArrayList; +import java.util.List; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +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.Mockito.doAnswer; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** + * Task 9:chunk 查询从单行分页改为批量 keyset 分页,保持低内存读取。 + * loadChunksKeyset 按 id 递增分批拉取(每批 pageSize),最后按 chunkIndex 升序合并, + * 避免超大任务一次 selectList 全量载入 chunk 元数据。 + * mock 分页由 wrapper 中 gt("id", lastId) 的 keyset 值驱动,保证重复调用可复现(幂等)。 + */ +@ExtendWith(MockitoExtension.class) +class SimilarAsinTaskServiceChunkKeysetTest { + + private static final String MODULE = "similar-asin"; + + @Mock private TaskChunkMapper taskChunkMapper; + + private static TaskChunkEntity chunk(long id, int chunkIndex) { + TaskChunkEntity chunk = new TaskChunkEntity(); + chunk.setId(id); + chunk.setTaskId(7004L); + chunk.setModuleType(MODULE); + chunk.setChunkIndex(chunkIndex); + return chunk; + } + + private static List chunks(long... idsAndIndexes) { + List result = new ArrayList<>(); + for (int i = 0; i < idsAndIndexes.length; i += 2) { + result.add(chunk(idsAndIndexes[i], (int) idsAndIndexes[i + 1])); + } + return result; + } + + /** 从 wrapper 的 SQL 片段解析 keyset:匹配 "id > #{ew.paramNameValuePairs.键}" 后从参数表取值。 */ + private static long keysetOf(QueryWrapper wrapper) { + java.util.regex.Matcher m = java.util.regex.Pattern + .compile("id\\s*>\\s*#\\{ew\\.paramNameValuePairs\\.(\\w+)\\}", java.util.regex.Pattern.CASE_INSENSITIVE) + .matcher(wrapper.getSqlSegment()); + if (m.find()) { + Object value = wrapper.getParamNameValuePairs().get(m.group(1)); + if (value instanceof Number number) { + return number.longValue(); + } + } + return 0L; + } + + /** 按 keyset 驱动分页:每次 selectList 返回 id > keyset 的下一批,天然支持重复调用。 */ + private void stubKeysetPages(List all, int pageSize) { + int batch = pageSize > 0 ? pageSize : 500; + doAnswer(invocation -> { + long lastId = keysetOf(invocation.getArgument(0)); + return all.stream().filter(c -> c.getId() > lastId).limit(batch).toList(); + }).when(taskChunkMapper).selectList(any()); + } + + private List keysetIdsFromRounds(int rounds) { + ArgumentCaptor> captor = + ArgumentCaptor.forClass(QueryWrapper.class); + verify(taskChunkMapper, times(rounds)).selectList(captor.capture()); + List keysets = new ArrayList<>(); + for (QueryWrapper wrapper : captor.getAllValues()) { + keysets.add(keysetOf(wrapper)); + } + return keysets; + } + + @Test + void test_task_009_chunk_normal_default_path() { + // 正常输入:chunk 数小于 pageSize,一轮拉完,结果全且按 chunkIndex 有序 + List all = chunks(1, 1, 2, 2, 3, 3); + stubKeysetPages(all, 500); + List result = SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, 7004L, MODULE, 500); + assertEquals(3, result.size()); + assertEquals(List.of(1, 2, 3), result.stream().map(TaskChunkEntity::getChunkIndex).toList()); + assertEquals(List.of(0L), keysetIdsFromRounds(1), "首轮 keyset 为 0"); + } + + @Test + void test_task_009_chunk_normal_multiple_items() { + // 超过 pageSize:多轮拉取,keyset 逐轮推进,全部合并且按 chunkIndex 升序、无重复 + List all = chunks(1, 1, 2, 2, 3, 3, 4, 4, 5, 5, 6, 6, 7, 7); + stubKeysetPages(all, 3); + List result = SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, 7004L, MODULE, 3); + assertEquals(7, result.size()); + assertEquals(List.of(1, 2, 3, 4, 5, 6, 7), result.stream().map(TaskChunkEntity::getChunkIndex).toList()); + long distinctIds = result.stream().map(TaskChunkEntity::getId).distinct().count(); + assertEquals(7, distinctIds, "keyset 分页不能产生重复 chunk"); + assertEquals(List.of(0L, 3L, 6L), keysetIdsFromRounds(3), "keyset 逐轮推进,不足一批即止"); + } + + @Test + void test_task_009_chunk_normal_repeated_operation_is_idempotent() { + // 重复执行同一输入:每轮都从 keyset=0 开始,结果一致 + List all = chunks(1, 1, 2, 2, 3, 3); + stubKeysetPages(all, 500); + List first = SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, 7004L, MODULE, 500); + List second = SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, 7004L, MODULE, 500); + assertEquals(first.size(), second.size()); + for (int i = 0; i < first.size(); i++) { + assertEquals(first.get(i).getId(), second.get(i).getId()); + } + assertEquals(List.of(0L, 0L), keysetIdsFromRounds(2), "重复执行每轮都从 keyset=0 开始"); + } + + @Test + void test_task_009_chunk_boundary_empty_input() { + // 空集合:返回空列表,不创建无效资源,且只查一轮 + stubKeysetPages(List.of(), 500); + List result = SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, 7004L, MODULE, 500); + assertNotNull(result); + assertEquals(0, result.size()); + verify(taskChunkMapper, times(1)).selectList(any()); + } + + @Test + void test_task_009_chunk_boundary_single_item() { + // 单 chunk:一轮返回后 keyset 推进即拉空,不依赖批量路径 + List all = chunks(42, 9); + stubKeysetPages(all, 1); + List result = SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, 7004L, MODULE, 1); + assertEquals(1, result.size()); + assertEquals(9, result.get(0).getChunkIndex()); + assertEquals(42L, result.get(0).getId()); + } + + @Test + void test_task_009_chunk_boundary_limit_and_overflow() { + // chunk 数恰好等于 pageSize 的倍数:最后一轮仍返回非空才继续,全部取回 + List all = chunks(1, 1, 2, 2, 3, 3, 4, 4, 5, 5, 6, 6); + stubKeysetPages(all, 3); + List result = SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, 7004L, MODULE, 3); + assertEquals(6, result.size()); + // pageSize 为 0/负数:回退默认 500,不抛异常 + stubKeysetPages(all, 0); + List fallback = SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, 7004L, MODULE, 0); + assertEquals(6, fallback.size()); + stubKeysetPages(all, -5); + List negative = SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, 7004L, MODULE, -5); + assertEquals(6, negative.size()); + } + + @Test + void test_task_009_chunk_invalid_input_rejected() { + // taskId 为 null:安全返回空列表,不发起查询 + List result = SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, null, MODULE, 500); + assertNotNull(result); + assertEquals(0, result.size()); + verify(taskChunkMapper, times(0)).selectList(any()); + // mapper 查询抛异常:转项目约定异常,消息可识别 + when(taskChunkMapper.selectList(any())).thenThrow(new IllegalStateException("db down")); + BusinessException ex = assertThrows(BusinessException.class, + () -> SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, 7004L, MODULE, 500)); + assertTrue(ex.getMessage() != null && ex.getMessage().contains("chunk"), + "异常消息必须可识别,实际: " + ex.getMessage()); + } + + @Test + void test_task_009_chunk_dependency_failure_releases_resources() { + // 第二轮查询失败:抛异常不返回半截结果;恢复后重试可完整返回 + List all = chunks(1, 1, 2, 2, 3, 3, 4, 4, 5, 5); + doAnswer(invocation -> { + long lastId = keysetOf(invocation.getArgument(0)); + if (lastId == 0L) { + return all.subList(0, 3); + } + if (lastId == 3L) { + throw new IllegalStateException("db down mid-page"); + } + return all.subList(3, 5); + }).when(taskChunkMapper).selectList(any()); + BusinessException ex = assertThrows(BusinessException.class, + () -> SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, 7004L, MODULE, 3)); + assertTrue(ex.getMessage() != null && ex.getMessage().contains("chunk"), + "异常消息必须可识别,实际: " + ex.getMessage()); + // 恢复后重试成功:5 个 chunk 全部取回 + stubKeysetPages(all, 3); + List recovered = SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, 7004L, MODULE, 3); + assertEquals(5, recovered.size()); + assertEquals(List.of(1, 2, 3, 4, 5), recovered.stream().map(TaskChunkEntity::getChunkIndex).toList()); + } +}