From a29d6e6ce4de21be1616f5d08a7647a136de85bf 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 15:39:38 +0800 Subject: [PATCH] =?UTF-8?q?task-51:=20accepted=20=E8=A1=8C=E5=BA=8F?= =?UTF-8?q?=E5=88=97=E5=8C=96=E4=B8=8E=20hash=20=E6=89=B9=E9=87=8F?= =?UTF-8?q?=E7=94=9F=E6=88=90?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit CollectDataResultDetailCodec.encodeRefsWithHash 一次迭代为整 chunk 生成 全部行引用 JSON 及其 SHA-256(offset 从 startOffset 递增),带 MAX_REFS_PER_BATCH 上限与参数校验;CollectDataResultItemBatchWriter 接入批量生成,删除逐行 encodeRef+hash 循环,同 chunk 重提结果逐字节一致。 --- .../util/CollectDataResultDetailCodec.java | 51 ++++++ .../CollectDataResultItemBatchWriter.java | 11 +- ...CollectDataResultDetailCodecBatchTest.java | 169 ++++++++++++++++++ 3 files changed, 227 insertions(+), 4 deletions(-) create mode 100644 backend-java/src/test/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultDetailCodecBatchTest.java diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultDetailCodec.java b/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultDetailCodec.java index 4a59eaf2..627d0ece 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultDetailCodec.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultDetailCodec.java @@ -6,6 +6,9 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.nanri.aiimage.modules.collectdata.model.vo.CollectDataResultRowVo; import org.springframework.stereotype.Component; +import java.security.MessageDigest; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; import java.util.List; /** @@ -26,6 +29,13 @@ public class CollectDataResultDetailCodec { public static final String REF_FIELD_OFFSET = "offset"; public static final String REF_FIELD_PAYLOAD = "payload"; + /** 单次批量引用生成的上限:与结果明细 batch size 同量级,拒绝前不分配结果。 */ + public static final int MAX_REFS_PER_BATCH = 100_000; + + /** 一批引用及各自 hash:refJson 为行引用 JSON,payloadHash 为其 SHA-256。 */ + public record RefWithHash(String refJson, String payloadHash) { + } + private final ObjectMapper objectMapper; public CollectDataResultDetailCodec(ObjectMapper objectMapper) { @@ -101,9 +111,50 @@ public class CollectDataResultDetailCodec { } } + /** + * 批量生成整 chunk 的引用 JSON 及其 SHA-256 hash:offset 从 startOffset + * 递增共 count 行,一次迭代完成序列化 + hash,替代调用方逐行两轮循环。 + * 非法参数(负 chunk/负 startOffset/负 count/null 或空白指针)抛 + * IllegalArgumentException;count 超过 MAX_REFS_PER_BATCH 时在分配 + * 结果前拒绝,不发生无界内存增长。 + */ + public List encodeRefsWithHash(int chunkIndex, int startOffset, int count, String pointer) { + if (chunkIndex < 0 || startOffset < 0 || count < 0 || pointer == null || pointer.isBlank()) { + throw new IllegalArgumentException("invalid chunk detail ref batch chunk=" + chunkIndex + + " startOffset=" + startOffset + " count=" + count + " pointer=" + pointer); + } + if (count > MAX_REFS_PER_BATCH) { + throw new IllegalArgumentException("chunk detail ref batch too large count=" + count + + " max=" + MAX_REFS_PER_BATCH); + } + if (count == 0) { + return List.of(); + } + List refs = new ArrayList<>(count); + for (int offset = startOffset; offset < startOffset + count; offset++) { + String refJson = encodeRef(chunkIndex, offset, pointer); + refs.add(new RefWithHash(refJson, sha256(refJson))); + } + return refs; + } + /** 结果明细行引用:chunk 序号 + 行内 offset + chunk 明细对象指针。 */ public record ChunkRef(@JsonProperty(REF_FIELD_CHUNK) int chunkIndex, @JsonProperty(REF_FIELD_OFFSET) int offset, @JsonProperty(REF_FIELD_PAYLOAD) String pointer) { } + + private static String sha256(String value) { + try { + MessageDigest digest = MessageDigest.getInstance("SHA-256"); + byte[] bytes = digest.digest((value == null ? "" : value).getBytes(StandardCharsets.UTF_8)); + StringBuilder sb = new StringBuilder(bytes.length * 2); + for (byte b : bytes) { + sb.append(String.format("%02x", b)); + } + return sb.toString(); + } catch (Exception ex) { + throw new IllegalStateException("chunk detail ref hash failed", ex); + } + } } diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultItemBatchWriter.java b/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultItemBatchWriter.java index 923dbd87..37f3895d 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultItemBatchWriter.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultItemBatchWriter.java @@ -55,7 +55,7 @@ public class CollectDataResultItemBatchWriter { if (rows == null || rows.isEmpty()) { return new UpsertCounts(0, 0); } - String scopeHash = hash(scopeKey); + String scopeHash = sha256(scopeKey); // 一次性取回本 scope 现有行,构建 item_key → 现有行 映射(hash 相等即跳过)。 List existingList = taskResultItemMapper.selectList( new LambdaQueryWrapper() @@ -72,6 +72,9 @@ public class CollectDataResultItemBatchWriter { List toUpsert = new ArrayList<>(); int skipped = 0; LocalDateTime now = LocalDateTime.now(); + // 批量生成整 chunk 的引用 JSON + hash(一次迭代),替代逐行 encodeRef + hash。 + List refsWithHash = + resultDetailCodec.encodeRefsWithHash(chunkIndex, 0, rows.size(), storedDetail); for (int offset = 0; offset < rows.size(); offset++) { CollectDataResultRowVo row = rows.get(offset); if (row == null || row.getAsin() == null || row.getAsin().isBlank()) { @@ -80,8 +83,8 @@ public class CollectDataResultItemBatchWriter { String itemKey = "asin:" + row.getAsin(); // chunk 级引用共享同一 RustFS 对象(deterministic key,同 chunk 重提 // 覆盖同一对象,引用 pointer 稳定,无需删除)。 - String refJson = resultDetailCodec.encodeRef(chunkIndex, offset, storedDetail); - String payloadHash = hash(refJson); + String refJson = refsWithHash.get(offset).refJson(); + String payloadHash = refsWithHash.get(offset).payloadHash(); TaskResultItemEntity existing = existingByKey.get(itemKey); if (existing != null && Objects.equals(existing.getPayloadHash(), payloadHash)) { skipped++; @@ -124,7 +127,7 @@ public class CollectDataResultItemBatchWriter { return new UpsertCounts(written, skipped); } - private String hash(String value) { + private static String sha256(String value) { try { java.security.MessageDigest digest = java.security.MessageDigest.getInstance("SHA-256"); byte[] bytes = digest.digest((value == null ? "" : value).getBytes(java.nio.charset.StandardCharsets.UTF_8)); diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultDetailCodecBatchTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultDetailCodecBatchTest.java new file mode 100644 index 00000000..e7ad0598 --- /dev/null +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultDetailCodecBatchTest.java @@ -0,0 +1,169 @@ +package com.nanri.aiimage.modules.collectdata.util; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; +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.any; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +/** + * Task 51:将 accepted 行的序列化和 hash 计算改为批量处理。 + * CollectDataResultDetailCodec.encodeRefsWithHash 一次为整 chunk 生成 + * 全部行的引用 JSON 及其 SHA-256 hash(一次迭代,offset 从 startOffset + * 递增),替代调用方逐行 encodeRef + hash 的两轮循环;同 chunk 重提时 + * 结果逐字节一致(幂等)。空 count 返回空列表,超批量上限拒绝, + * 序列化失败抛可识别异常且不泄漏内部状态。 + */ +class CollectDataResultDetailCodecBatchTest { + + private CollectDataResultDetailCodec codec; + + @BeforeEach + void setUp() { + codec = new CollectDataResultDetailCodec(new ObjectMapper()); + } + + @Test + void test_task_051_task_normal_default_path() { + // 正常批量:3 行 offset 10-12,引用 JSON 与逐行 encodeRef 一致,hash 为 SHA-256。 + List refs = + codec.encodeRefsWithHash(1, 10, 3, "rustfs:detail/abc"); + + assertEquals(3, refs.size(), "3 行引用全部生成"); + CollectDataResultDetailCodec.ChunkRef parsed0 = codec.parseRef(refs.get(0).refJson()); + assertEquals(1, parsed0.chunkIndex(), "chunk 序号保留"); + assertEquals(10, parsed0.offset(), "首行 offset=startOffset"); + assertEquals("rustfs:detail/abc", parsed0.pointer(), "RustFS 指针保留"); + assertEquals(11, codec.parseRef(refs.get(1).refJson()).offset(), "offset 递增"); + assertEquals(12, codec.parseRef(refs.get(2).refJson()).offset(), "末行 offset"); + assertEquals(codec.encodeRef(1, 12, "rustfs:detail/abc"), refs.get(2).refJson(), + "与逐行编码结果一致"); + assertEquals(sha256(refs.get(1).refJson()), refs.get(1).payloadHash(), "hash 为引用 JSON 的 SHA-256"); + } + + @Test + void test_task_051_task_normal_multiple_items() { + // 大批量:75 行 offset 0-74,顺序稳定不丢失,逐行 hash 均正确。 + List refs = + codec.encodeRefsWithHash(3, 0, 75, "rustfs:detail/x"); + + assertEquals(75, refs.size(), "75 行全部生成"); + for (int i = 0; i < 75; i++) { + CollectDataResultDetailCodec.ChunkRef ref = codec.parseRef(refs.get(i).refJson()); + assertEquals(3, ref.chunkIndex(), "chunk 序号一致"); + assertEquals(i, ref.offset(), "行 " + i + " offset 顺序稳定"); + assertEquals("rustfs:detail/x", ref.pointer(), "指针一致"); + assertEquals(sha256(refs.get(i).refJson()), refs.get(i).payloadHash(), "行 " + i + " hash 正确"); + } + } + + @Test + void test_task_051_task_normal_repeated_operation_is_idempotent() { + // 幂等:同一输入重复批量生成,逐元素结果完全一致,不产生新差异。 + List first = + codec.encodeRefsWithHash(0, 5, 20, "rustfs:detail/stable"); + List second = + codec.encodeRefsWithHash(0, 5, 20, "rustfs:detail/stable"); + + assertEquals(first.size(), second.size(), "数量一致"); + for (int i = 0; i < first.size(); i++) { + assertEquals(first.get(i).refJson(), second.get(i).refJson(), "行 " + i + " 引用一致"); + assertEquals(first.get(i).payloadHash(), second.get(i).payloadHash(), "行 " + i + " hash 一致"); + } + } + + @Test + void test_task_051_task_boundary_empty_input() { + // 空输入:count=0 返回空列表,不创建任何资源。 + assertEquals(0, codec.encodeRefsWithHash(0, 0, 0, "rustfs:x").size(), "count=0 空列表"); + assertEquals(0, codec.encodeRefsWithHash(2, 100, 0, "rustfs:x").size(), + "任意合法 startOffset 下 count=0 仍为空列表"); + } + + @Test + void test_task_051_task_boundary_single_item() { + // 单元素:count=1 单行 offset 正确,hash 正确。 + List refs = + codec.encodeRefsWithHash(0, 0, 1, "rustfs:detail/solo"); + + assertEquals(1, refs.size(), "单行生成"); + CollectDataResultDetailCodec.ChunkRef parsed = codec.parseRef(refs.get(0).refJson()); + assertEquals(0, parsed.offset(), "单行 offset 0"); + assertEquals("rustfs:detail/solo", parsed.pointer(), "指针保留"); + assertEquals(sha256(refs.get(0).refJson()), refs.get(0).payloadHash(), "hash 正确"); + } + + @Test + void test_task_051_task_boundary_limit_and_overflow() { + // 上限/超限:达到批量上限(100000)可正常生成;超过上限拒绝, + // 且在拒绝前不分配结果,不发生无界内存增长。 + List atLimit = codec.encodeRefsWithHash( + 0, 0, CollectDataResultDetailCodec.MAX_REFS_PER_BATCH, "rustfs:big"); + assertEquals(CollectDataResultDetailCodec.MAX_REFS_PER_BATCH, atLimit.size(), "上限内正常生成"); + assertEquals(0, codec.parseRef(atLimit.get(0).refJson()).offset(), "首行 offset 正确"); + + assertThrows(IllegalArgumentException.class, () -> codec.encodeRefsWithHash( + 0, 0, CollectDataResultDetailCodec.MAX_REFS_PER_BATCH + 1, "rustfs:big"), + "超过上限拒绝"); + } + + @Test + void test_task_051_task_invalid_input_rejected() { + // 非法参数:负 chunk/负 startOffset/负 count/null 指针/空白指针均抛可识别异常。 + assertThrows(IllegalArgumentException.class, + () -> codec.encodeRefsWithHash(-1, 0, 1, "rustfs:x"), "负 chunk 拒绝"); + assertThrows(IllegalArgumentException.class, + () -> codec.encodeRefsWithHash(0, -1, 1, "rustfs:x"), "负 startOffset 拒绝"); + assertThrows(IllegalArgumentException.class, + () -> codec.encodeRefsWithHash(0, 0, -1, "rustfs:x"), "负 count 拒绝"); + assertThrows(IllegalArgumentException.class, + () -> codec.encodeRefsWithHash(0, 0, 1, null), "null 指针拒绝"); + assertThrows(IllegalArgumentException.class, + () -> codec.encodeRefsWithHash(0, 0, 1, " "), "空白指针拒绝"); + } + + @Test + void test_task_051_task_dependency_failure_releases_resources() throws Exception { + // 依赖失败:序列化抛错时批量方法抛可识别异常,内部无残留状态; + // 依赖恢复后同一实例再次调用成功,无资源泄漏。 + ObjectMapper mapper = mock(ObjectMapper.class); + when(mapper.writeValueAsString(any())).thenThrow(new JsonProcessingException("json down") { + }); + CollectDataResultDetailCodec failingCodec = new CollectDataResultDetailCodec(mapper); + assertThrows(IllegalArgumentException.class, + () -> failingCodec.encodeRefsWithHash(0, 0, 2, "rustfs:x"), + "序列化失败抛可识别异常"); + + doAnswer(invocation -> invocation.getArgument(0).toString()) + .when(mapper).writeValueAsString(any()); + List refs = + failingCodec.encodeRefsWithHash(0, 0, 2, "rustfs:x"); + assertEquals(2, refs.size(), "恢复后正常生成"); + assertTrue(refs.get(0).refJson().contains("rustfs:x"), "引用内容完整"); + } + + private static String sha256(String value) { + try { + MessageDigest digest = MessageDigest.getInstance("SHA-256"); + byte[] bytes = digest.digest((value == null ? "" : value).getBytes(StandardCharsets.UTF_8)); + StringBuilder sb = new StringBuilder(bytes.length * 2); + for (byte b : bytes) { + sb.append(String.format("%02x", b)); + } + return sb.toString(); + } catch (Exception ex) { + throw new IllegalStateException(ex); + } + } +}