task-9: chunk 查询改为批量 keyset 分页,保持低内存读取
loadChunksKeyset 按 id 递增分批拉取(每批 pageSize),最后按 chunkIndex 升序合并,避免超大任务一次 selectList 全量载入 chunk 元数据。surefire fork 堆显式限制 1536m 防止与主 JVM 叠加挤爆内存。
This commit is contained in:
@@ -183,7 +183,7 @@
|
||||
<groupId>org.apache.maven.plugins</groupId>
|
||||
<artifactId>maven-surefire-plugin</artifactId>
|
||||
<configuration>
|
||||
<argLine>-XX:+EnableDynamicAgentLoading -Xshare:off</argLine>
|
||||
<argLine>-XX:+EnableDynamicAgentLoading -Xshare:off -Xmx1536m</argLine>
|
||||
</configuration>
|
||||
</plugin>
|
||||
</plugins>
|
||||
|
||||
+40
@@ -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<TaskChunkEntity> loadChunksKeyset(TaskChunkMapper mapper, Long taskId, String moduleType, int pageSize) {
|
||||
if (taskId == null || mapper == null) {
|
||||
return List.of();
|
||||
}
|
||||
int batch = pageSize > 0 ? pageSize : 500;
|
||||
List<TaskChunkEntity> all = new ArrayList<>();
|
||||
long lastId = 0L;
|
||||
while (true) {
|
||||
List<TaskChunkEntity> page;
|
||||
try {
|
||||
page = mapper.selectList(new com.baomidou.mybatisplus.core.conditions.query.QueryWrapper<TaskChunkEntity>()
|
||||
.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;
|
||||
|
||||
+204
@@ -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<TaskChunkEntity> chunks(long... idsAndIndexes) {
|
||||
List<TaskChunkEntity> 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<TaskChunkEntity> 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<TaskChunkEntity> 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<Long> keysetIdsFromRounds(int rounds) {
|
||||
ArgumentCaptor<QueryWrapper<TaskChunkEntity>> captor =
|
||||
ArgumentCaptor.forClass(QueryWrapper.class);
|
||||
verify(taskChunkMapper, times(rounds)).selectList(captor.capture());
|
||||
List<Long> keysets = new ArrayList<>();
|
||||
for (QueryWrapper<TaskChunkEntity> wrapper : captor.getAllValues()) {
|
||||
keysets.add(keysetOf(wrapper));
|
||||
}
|
||||
return keysets;
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_009_chunk_normal_default_path() {
|
||||
// 正常输入:chunk 数小于 pageSize,一轮拉完,结果全且按 chunkIndex 有序
|
||||
List<TaskChunkEntity> all = chunks(1, 1, 2, 2, 3, 3);
|
||||
stubKeysetPages(all, 500);
|
||||
List<TaskChunkEntity> 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<TaskChunkEntity> all = chunks(1, 1, 2, 2, 3, 3, 4, 4, 5, 5, 6, 6, 7, 7);
|
||||
stubKeysetPages(all, 3);
|
||||
List<TaskChunkEntity> 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<TaskChunkEntity> all = chunks(1, 1, 2, 2, 3, 3);
|
||||
stubKeysetPages(all, 500);
|
||||
List<TaskChunkEntity> first = SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, 7004L, MODULE, 500);
|
||||
List<TaskChunkEntity> 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<TaskChunkEntity> 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<TaskChunkEntity> all = chunks(42, 9);
|
||||
stubKeysetPages(all, 1);
|
||||
List<TaskChunkEntity> 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<TaskChunkEntity> all = chunks(1, 1, 2, 2, 3, 3, 4, 4, 5, 5, 6, 6);
|
||||
stubKeysetPages(all, 3);
|
||||
List<TaskChunkEntity> result = SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, 7004L, MODULE, 3);
|
||||
assertEquals(6, result.size());
|
||||
// pageSize 为 0/负数:回退默认 500,不抛异常
|
||||
stubKeysetPages(all, 0);
|
||||
List<TaskChunkEntity> fallback = SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, 7004L, MODULE, 0);
|
||||
assertEquals(6, fallback.size());
|
||||
stubKeysetPages(all, -5);
|
||||
List<TaskChunkEntity> negative = SimilarAsinTaskService.loadChunksKeyset(taskChunkMapper, 7004L, MODULE, -5);
|
||||
assertEquals(6, negative.size());
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_009_chunk_invalid_input_rejected() {
|
||||
// taskId 为 null:安全返回空列表,不发起查询
|
||||
List<TaskChunkEntity> 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<TaskChunkEntity> 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<TaskChunkEntity> 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());
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user