task-58: collectdata 结果明细读取失败降级(RustFS 超时跳过该 chunk,不中断其它 chunk)

chunk 级引用读取失败(RustFS 超时/不可用)时降级为跳过该 chunk 的引用行并记录 warn,
结果文件继续生成,缺失行由任务状态可观测;chunk 明细 JSON 损坏仍抛异常(数据问题不降级),
旧格式单行读取失败保持抛异常。失败 chunk 在单次调用内缓存为空哨兵,避免重复 resolve。
task-052 依赖失败测试同步更新为降级语义。
This commit is contained in:
2026-08-30 17:31:21 +08:00
parent 44abcdd4ef
commit 567cb1ab67
3 changed files with 234 additions and 7 deletions
@@ -5,6 +5,7 @@ import com.fasterxml.jackson.databind.ObjectMapper;
import com.nanri.aiimage.modules.collectdata.model.vo.CollectDataResultRowVo;
import com.nanri.aiimage.modules.task.model.entity.TaskResultItemEntity;
import com.nanri.aiimage.modules.task.service.TransientPayloadStorageService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import java.util.ArrayList;
@@ -19,6 +20,7 @@ import java.util.Map;
* 空列表,越界/损坏行安全跳过,读取失败抛可识别异常。
*/
@Component
@Slf4j
public class CollectDataResultDetailReader {
private static final String ERROR_DETAIL = "read collect data result detail failed";
@@ -54,10 +56,24 @@ public class CollectDataResultDetailReader {
if (ref != null) {
List<CollectDataResultRowVo> details = detailCache.get(ref.pointer());
if (details == null) {
try {
String detailJson = transientPayloadStorageService.resolvePayload(ref.pointer(), ERROR_DETAIL);
// 空内容用空列表占位:与缓存 miss 区分,避免同一 chunk 重复 resolve。
details = (detailJson == null || detailJson.isBlank())
? null
? List.of()
: objectMapper.readValue(detailJson, listType);
} catch (Exception readEx) {
// RustFS 超时/不可用等读取失败:跳过该 chunk 的引用行,
// 不中断其它 chunk 的行读取(结果文件降级生成,缺失行
// 由任务状态可观测);chunk 明细 JSON 损坏则正常上报,
// 属于数据问题而非依赖降级。
if (readEx instanceof com.fasterxml.jackson.core.JsonProcessingException) {
throw readEx;
}
log.warn("[collect-data] skip chunk detail read failed pointer={} err={}",
ref.pointer(), readEx.getMessage());
details = List.of();
}
detailCache.put(ref.pointer(), details);
}
if (details != null && ref.offset() >= 0 && ref.offset() < details.size()) {
@@ -26,7 +26,8 @@ import static org.mockito.Mockito.when;
* CollectDataResultDetailReader 从 biz_task_result_item 解析行:引用格式
* 按 pointer 缓存整 chunk 明细(同一 chunk 对象只 resolve/解析一次),
* 再按 offset 取行;旧格式逐行兜底。空输入返回空列表,越界/损坏行安全
* 跳过,读取失败抛可识别异常且可恢复。
* 跳过,chunk 级读取失败降级跳过该 chunk(不中断其它 chunk),
* 旧格式读取失败仍抛可识别异常。
*/
class CollectDataResultDetailReaderTest {
@@ -163,14 +164,16 @@ class CollectDataResultDetailReaderTest {
@Test
void test_task_052_chunk_dependency_failure_releases_resources() {
// 依赖失败:对象存储读取抛错时抛可识别异常(不泄漏内部状态);
// 依赖失败task-58 起降级语义):chunk 级读取失败 → 跳过该 chunk
// 的引用行(降级为空结果,不抛异常、不中断其它 chunk),
// 依赖恢复后同一实例再次调用成功,无资源残留。
doThrow(new IllegalStateException("object storage down"))
.doReturn("[{\"asin\":\"B000000001\",\"brand\":\"back\"}]")
.when(transientPayloadStorageService)
.resolvePayload("rustfs:detail/down", "read collect data result detail failed");
List<TaskResultItemEntity> items = List.of(item(1L, refJson(0, 0, "rustfs:detail/down")));
assertThrows(RuntimeException.class, () -> reader.readRows(items), "读取失败抛可识别异常");
List<CollectDataResultRowVo> degraded = reader.readRows(items);
assertEquals(0, degraded.size(), "读取失败 chunk 降级为空结果");
List<CollectDataResultRowVo> rows = reader.readRows(items);
assertEquals(1, rows.size(), "恢复后正常读取");
@@ -0,0 +1,208 @@
package com.nanri.aiimage.modules.collectdata.util;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.nanri.aiimage.modules.collectdata.model.vo.CollectDataResultRowVo;
import com.nanri.aiimage.modules.task.model.entity.TaskResultItemEntity;
import com.nanri.aiimage.modules.task.service.TransientPayloadStorageService;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.util.ArrayList;
import java.util.List;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
/**
* Task 58:外部品牌服务不可用、RustFS 超时和重复 chunk 的降级测试。
* 新增降级语义:chunk 级引用读取失败(RustFS 超时/不可用)时跳过该 chunk
* 的引用行并记录 warn,不中断其它 chunk 的行读取;旧格式单行读取失败仍抛
* 可识别异常(不可降级,避免静默丢数据)。重复 chunk(多 item 引用同一
* RustFS 对象)只 resolve 一次,失败降级对同一 chunk 只发生一次。
*/
class CollectDataRustfsDegradeTest {
private TransientPayloadStorageService transientPayloadStorageService;
private CollectDataResultDetailReader reader;
@BeforeEach
void setUp() {
transientPayloadStorageService = mock(TransientPayloadStorageService.class);
reader = new CollectDataResultDetailReader(
new CollectDataResultDetailCodec(new ObjectMapper()), new ObjectMapper(),
transientPayloadStorageService);
}
@Test
void test_task_058_chunk_brand_rustfs_normal_default_path() {
// 正常路径:无失败时行为与既有一致,多 chunk 全部取回,无降级误触发。
List<TaskResultItemEntity> items = List.of(
item(1L, refJson(0, 0, "rustfs:detail/a")),
item(2L, refJson(0, 1, "rustfs:detail/a")),
item(3L, refJson(1, 0, "rustfs:detail/b")));
when(transientPayloadStorageService.resolvePayload("rustfs:detail/a", "read collect data result detail failed"))
.thenReturn(rowsJson(0, 2));
when(transientPayloadStorageService.resolvePayload("rustfs:detail/b", "read collect data result detail failed"))
.thenReturn(rowsJson(2, 1));
List<CollectDataResultRowVo> rows = reader.readRows(items);
assertEquals(3, rows.size(), "3 行全部取回");
assertEquals("B000000001", rows.get(0).getAsin(), "chunk a 首行");
assertEquals("B000000003", rows.get(2).getAsin(), "chunk b 行");
}
@Test
void test_task_058_chunk_brand_rustfs_normal_multiple_items() {
// 批量场景(外部服务不可用降级):chunk a 读取失败 → 其引用行降级跳过,
// chunk b/c 正常取回,顺序稳定不丢失。
when(transientPayloadStorageService.resolvePayload("rustfs:detail/down", "read collect data result detail failed"))
.thenThrow(new IllegalStateException("rustfs read timeout"));
when(transientPayloadStorageService.resolvePayload("rustfs:detail/ok1", "read collect data result detail failed"))
.thenReturn(rowsJson(0, 2));
when(transientPayloadStorageService.resolvePayload("rustfs:detail/ok2", "read collect data result detail failed"))
.thenReturn(rowsJson(2, 2));
List<TaskResultItemEntity> items = List.of(
item(1L, refJson(0, 0, "rustfs:detail/down")),
item(2L, refJson(0, 1, "rustfs:detail/down")),
item(3L, refJson(1, 0, "rustfs:detail/ok1")),
item(4L, refJson(1, 1, "rustfs:detail/ok1")),
item(5L, refJson(2, 0, "rustfs:detail/ok2")),
item(6L, refJson(2, 1, "rustfs:detail/ok2")));
List<CollectDataResultRowVo> rows = reader.readRows(items);
assertEquals(4, rows.size(), "失败 chunk 降级跳过 2 行,其余 4 行取回");
assertEquals("B000000001", rows.get(0).getAsin(), "ok1 首行");
assertEquals("B000000003", rows.get(2).getAsin(), "ok2 行顺序稳定");
}
@Test
void test_task_058_chunk_brand_rustfs_normal_repeated_operation_is_idempotent() {
// 幂等(重复 chunk):多个 item 引用同一失败对象 → 降级只发生一次、
// 结果稳定;重复读取结果一致,无状态残留。
when(transientPayloadStorageService.resolvePayload("rustfs:detail/x", "read collect data result detail failed"))
.thenThrow(new IllegalStateException("rustfs down"))
.thenReturn(rowsJson(0, 2));
List<TaskResultItemEntity> items = List.of(
item(1L, refJson(0, 0, "rustfs:detail/x")),
item(2L, refJson(0, 1, "rustfs:detail/x")));
List<CollectDataResultRowVo> first = reader.readRows(items);
assertEquals(0, first.size(), "失败时整 chunk 降级为空");
verify(transientPayloadStorageService, times(1)).resolvePayload(anyString(), anyString());
// 恢复后再次读取:本次调用重新 resolve,chunk 行全部取回。
List<CollectDataResultRowVo> second = reader.readRows(items);
assertEquals(2, second.size(), "恢复后同一实例再次读取成功");
assertEquals("B000000002", second.get(1).getAsin(), "两行均恢复");
}
@Test
void test_task_058_chunk_brand_rustfs_boundary_empty_input() {
// 空输入:空列表安全跳过,不发起读取,不触发降级。
List<CollectDataResultRowVo> rows = reader.readRows(List.of());
assertEquals(0, rows.size(), "空输入返回空列表");
verify(transientPayloadStorageService, times(0)).resolvePayload(anyString(), anyString());
}
@Test
void test_task_058_chunk_brand_rustfs_boundary_single_item() {
// 单元素:单个 chunk 读取失败 → 返回空结果不抛异常(降级路径不依赖批量)。
when(transientPayloadStorageService.resolvePayload("rustfs:detail/solo", "read collect data result detail failed"))
.thenThrow(new IllegalStateException("rustfs timeout"));
List<CollectDataResultRowVo> rows = reader.readRows(List.of(item(1L, refJson(0, 0, "rustfs:detail/solo"))));
assertEquals(0, rows.size(), "单 chunk 失败降级为空结果");
}
@Test
void test_task_058_chunk_brand_rustfs_boundary_limit_and_overflow() {
// 上限/超限:全部 chunk 失败(10 个不同对象)→ 空结果,各 resolve 一次,
// 无无界累积、无重复调用。
List<TaskResultItemEntity> items = new ArrayList<>();
for (int i = 0; i < 10; i++) {
items.add(item((long) i + 1, refJson(i, 0, "rustfs:detail/fail" + i)));
}
for (int i = 0; i < 10; i++) {
when(transientPayloadStorageService.resolvePayload("rustfs:detail/fail" + i, "read collect data result detail failed"))
.thenThrow(new IllegalStateException("rustfs timeout"));
}
List<CollectDataResultRowVo> rows = reader.readRows(items);
assertEquals(0, rows.size(), "全部失败降级为空");
verify(transientPayloadStorageService, times(10)).resolvePayload(anyString(), anyString());
}
@Test
void test_task_058_chunk_brand_rustfs_invalid_input_rejected() {
// 非法参数:旧格式单行读取失败不可降级(抛可识别异常,避免静默丢数据);
// chunk 明细 JSON 损坏同样拒绝。
when(transientPayloadStorageService.resolvePayload("rustfs:detail/broken", "read collect data result detail failed"))
.thenReturn("[{\"asin\":\"B000000001\"");
assertThrows(RuntimeException.class,
() -> reader.readRows(List.of(item(1L, refJson(0, 0, "rustfs:detail/broken")))),
"损坏 chunk JSON 仍抛异常(数据损坏不降级)");
when(transientPayloadStorageService.resolvePayload("{\"asin\":\"B000000001\"", "read collect data result item failed"))
.thenThrow(new IllegalStateException("rustfs down"));
assertThrows(RuntimeException.class,
() -> reader.readRows(List.of(item(2L, "{\"asin\":\"B000000001\""))),
"旧格式行读取失败抛可识别异常(不可降级)");
}
@Test
void test_task_058_chunk_brand_rustfs_dependency_failure_releases_resources() {
// 依赖失败(品牌服务不可用语义):部分 chunk 失败降级、其余正常;
// 依赖恢复后同一实例再次调用全部取回,无资源残留。
when(transientPayloadStorageService.resolvePayload("rustfs:detail/partial", "read collect data result detail failed"))
.thenThrow(new IllegalStateException("brand service down"))
.thenReturn(rowsJson(0, 2), rowsJson(0, 2));
when(transientPayloadStorageService.resolvePayload("rustfs:detail/keep", "read collect data result detail failed"))
.thenReturn(rowsJson(2, 1));
List<TaskResultItemEntity> items = List.of(
item(1L, refJson(0, 0, "rustfs:detail/partial")),
item(2L, refJson(0, 1, "rustfs:detail/partial")),
item(3L, refJson(1, 0, "rustfs:detail/keep")));
List<CollectDataResultRowVo> degraded = reader.readRows(items);
assertEquals(1, degraded.size(), "失败 chunk 降级,keep 正常");
assertTrue(degraded.get(0).getAsin().contains("B000000003"), "正常 chunk 行不丢失");
verify(transientPayloadStorageService, times(1)).resolvePayload("rustfs:detail/partial", "read collect data result detail failed");
List<CollectDataResultRowVo> recovered = reader.readRows(items);
assertEquals(3, recovered.size(), "恢复后全部取回,无残留");
}
private static TaskResultItemEntity item(Long id, String payloadJson) {
TaskResultItemEntity entity = new TaskResultItemEntity();
entity.setId(id);
entity.setPayloadJson(payloadJson);
return entity;
}
private static String refJson(int chunk, int offset, String pointer) {
return "{\"chunk\":" + chunk + ",\"offset\":" + offset + ",\"payload\":\"" + pointer + "\"}";
}
private static String rowsJson(int start, int count) {
StringBuilder sb = new StringBuilder("[");
for (int i = 0; i < count; i++) {
if (i > 0) {
sb.append(",");
}
sb.append("{\"asin\":\"B").append(String.format("%09d", start + i + 1))
.append("\",\"brand\":\"brand").append(i).append("\"}");
}
return sb.append("]").toString();
}
}