修复采集数据结果明细 module_type 写读大小写不一致:写入端误用小写 collectdata 导致 task-50 之后所有采集任务明细行生成结果文件时读不到

线上现象:任务 26334 采集完成(去重 42/无效品牌 62/品牌拒绝 6/保留 85)
但结果文件「采集数据结果」sheet 只有表头 0 行,历史卡片「保留 0」。

根因链:
1. CollectDataResultItemBatchWriter(task-50 新增写入端)常量误写为小写
   "collectdata",与 CollectDataService.MODULE_TYPE(大写 COLLECT_DATA)不一致;
2. loadFinalRows 按大写查 biz_task_result_item 命中 0 行 → 组装空表 →
   stats.finalRowCount 被 rows.size()=0 覆盖,UI 显示保留 0;
3. 生产库 5386 行小写 collectdata(8/31 起约 15 个任务)全部读不到,
   大写 COLLECT_DATA 仅存 1067 行历史数据。

修复:
- 写入常量改回大写 COLLECT_DATA(与读取端一致),现有行查询 in(大写,小写)
  兼容存量,旧任务重试幂等命中不重复插入;
- loadFinalRows / deleteTask / deleteHistory 查询改 in(MODULE_TYPE, LEGACY_MODULE_TYPE),
  历史小写行可读出并正确清理(RustFS 对象不泄漏);
- 新增 2 条回归测试(写入 module_type 必须与读取端常量一致、小写存量行重试命中),
  collectdata 全包 168 测试全绿。
This commit is contained in:
2026-09-02 13:47:52 +08:00
parent 5f4fcad2ef
commit 478233ba7a
3 changed files with 60 additions and 7 deletions
@@ -97,6 +97,13 @@ public class CollectDataService {
public static final String MODULE_TYPE = "COLLECT_DATA"; public static final String MODULE_TYPE = "COLLECT_DATA";
/**
* task-50 上线初期写入端(CollectDataResultItemBatchWriter)曾误用小写
* "collectdata",遗留 5386 行存量数据;读取端查询统一 in(MODULE_TYPE, LEGACY_MODULE_TYPE)
* 兼容,否则这些任务的明细行在生成结果文件时查不出来。
*/
private static final String LEGACY_MODULE_TYPE = "collectdata";
public TaskProgressLightBatchVo progressLight(List<Long> taskIds) { public TaskProgressLightBatchVo progressLight(List<Long> taskIds) {
return taskProgressLightAssembler.assemble(MODULE_TYPE, null, taskIds); return taskProgressLightAssembler.assemble(MODULE_TYPE, null, taskIds);
} }
@@ -937,7 +944,7 @@ public class CollectDataService {
private List<CollectDataResultRowVo> loadFinalRows(Long taskId) { private List<CollectDataResultRowVo> loadFinalRows(Long taskId) {
List<TaskResultItemEntity> rows = taskResultItemMapper.selectList(new LambdaQueryWrapper<TaskResultItemEntity>() List<TaskResultItemEntity> rows = taskResultItemMapper.selectList(new LambdaQueryWrapper<TaskResultItemEntity>()
.eq(TaskResultItemEntity::getTaskId, taskId) .eq(TaskResultItemEntity::getTaskId, taskId)
.eq(TaskResultItemEntity::getModuleType, MODULE_TYPE) .in(TaskResultItemEntity::getModuleType, MODULE_TYPE, LEGACY_MODULE_TYPE)
.orderByAsc(TaskResultItemEntity::getId)); .orderByAsc(TaskResultItemEntity::getId));
if (rows == null || rows.isEmpty()) { if (rows == null || rows.isEmpty()) {
return new ArrayList<>(); return new ArrayList<>();
@@ -1258,7 +1265,7 @@ public class CollectDataService {
.eq(TaskChunkEntity::getModuleType, MODULE_TYPE)); .eq(TaskChunkEntity::getModuleType, MODULE_TYPE));
taskResultItemMapper.delete(new LambdaQueryWrapper<TaskResultItemEntity>() taskResultItemMapper.delete(new LambdaQueryWrapper<TaskResultItemEntity>()
.eq(TaskResultItemEntity::getTaskId, task.getId()) .eq(TaskResultItemEntity::getTaskId, task.getId())
.eq(TaskResultItemEntity::getModuleType, MODULE_TYPE)); .in(TaskResultItemEntity::getModuleType, MODULE_TYPE, LEGACY_MODULE_TYPE));
deleteTransientTaskPayloads( deleteTransientTaskPayloads(
taskChunkMapper.selectList(new LambdaQueryWrapper<TaskChunkEntity>() taskChunkMapper.selectList(new LambdaQueryWrapper<TaskChunkEntity>()
.select(TaskChunkEntity::getPayloadJson) .select(TaskChunkEntity::getPayloadJson)
@@ -1267,7 +1274,7 @@ public class CollectDataService {
taskResultItemMapper.selectList(new LambdaQueryWrapper<TaskResultItemEntity>() taskResultItemMapper.selectList(new LambdaQueryWrapper<TaskResultItemEntity>()
.select(TaskResultItemEntity::getPayloadJson) .select(TaskResultItemEntity::getPayloadJson)
.eq(TaskResultItemEntity::getTaskId, task.getId()) .eq(TaskResultItemEntity::getTaskId, task.getId())
.eq(TaskResultItemEntity::getModuleType, MODULE_TYPE))); .in(TaskResultItemEntity::getModuleType, MODULE_TYPE, LEGACY_MODULE_TYPE)));
taskFileJobService.deleteTaskJobs(task.getId(), MODULE_TYPE); taskFileJobService.deleteTaskJobs(task.getId(), MODULE_TYPE);
fileTaskMapper.deleteById(task.getId()); fileTaskMapper.deleteById(task.getId());
} }
@@ -1280,11 +1287,11 @@ public class CollectDataService {
List<TaskResultItemEntity> resultItems = taskResultItemMapper.selectList(new LambdaQueryWrapper<TaskResultItemEntity>() List<TaskResultItemEntity> resultItems = taskResultItemMapper.selectList(new LambdaQueryWrapper<TaskResultItemEntity>()
.select(TaskResultItemEntity::getPayloadJson) .select(TaskResultItemEntity::getPayloadJson)
.eq(TaskResultItemEntity::getTaskId, row.getTaskId()) .eq(TaskResultItemEntity::getTaskId, row.getTaskId())
.eq(TaskResultItemEntity::getModuleType, MODULE_TYPE) .in(TaskResultItemEntity::getModuleType, MODULE_TYPE, LEGACY_MODULE_TYPE)
.eq(TaskResultItemEntity::getResultId, row.getId())); .eq(TaskResultItemEntity::getResultId, row.getId()));
taskResultItemMapper.delete(new LambdaQueryWrapper<TaskResultItemEntity>() taskResultItemMapper.delete(new LambdaQueryWrapper<TaskResultItemEntity>()
.eq(TaskResultItemEntity::getTaskId, row.getTaskId()) .eq(TaskResultItemEntity::getTaskId, row.getTaskId())
.eq(TaskResultItemEntity::getModuleType, MODULE_TYPE) .in(TaskResultItemEntity::getModuleType, MODULE_TYPE, LEGACY_MODULE_TYPE)
.eq(TaskResultItemEntity::getResultId, row.getId())); .eq(TaskResultItemEntity::getResultId, row.getId()));
// 与 deleteTask 一致:先删 DB 行再物理删对象,保证行删除与对象删除一致。 // 与 deleteTask 一致:先删 DB 行再物理删对象,保证行删除与对象删除一致。
deleteResultItemPayloads(resultItems); deleteResultItemPayloads(resultItems);
@@ -27,7 +27,12 @@ import java.util.Objects;
@Component @Component
public class CollectDataResultItemBatchWriter { public class CollectDataResultItemBatchWriter {
private static final String MODULE_TYPE = "collectdata"; /** 必须与 CollectDataService.MODULE_TYPE 一致(大写 COLLECT_DATA)。
* 曾因写端误用小写 "collectdata",导致 task-50 之后所有采集任务的明细行
* 在生成结果文件时查不出来(biz_task_result_item.module_type 大小写分叉)。 */
private static final String MODULE_TYPE = "COLLECT_DATA";
/** 历史存量小写值,查询现有行时兼容,避免旧任务重试重复插入(与 CollectDataService.LEGACY_MODULE_TYPE 对应)。 */
private static final String LEGACY_MODULE_TYPE = "collectdata";
private static final int DEFAULT_BATCH_SIZE = 100; private static final int DEFAULT_BATCH_SIZE = 100;
private final TaskResultItemMapper taskResultItemMapper; private final TaskResultItemMapper taskResultItemMapper;
@@ -63,10 +68,12 @@ public class CollectDataResultItemBatchWriter {
} }
String scopeHash = sha256(scopeKey); String scopeHash = sha256(scopeKey);
// 一次性取回本 scope 现有行,构建 item_key → 现有行 映射(hash 相等即跳过)。 // 一次性取回本 scope 现有行,构建 item_key → 现有行 映射(hash 相等即跳过)。
// 现有行查询兼容历史小写 "collectdata"(见 MODULE_TYPE 注释),避免旧任务重试
// 在同一唯一键 (task_id, module_type, scope_hash, item_key) 下重复插入。
List<TaskResultItemEntity> existingList = taskResultItemMapper.selectList( List<TaskResultItemEntity> existingList = taskResultItemMapper.selectList(
new LambdaQueryWrapper<TaskResultItemEntity>() new LambdaQueryWrapper<TaskResultItemEntity>()
.eq(TaskResultItemEntity::getTaskId, taskId) .eq(TaskResultItemEntity::getTaskId, taskId)
.eq(TaskResultItemEntity::getModuleType, MODULE_TYPE) .in(TaskResultItemEntity::getModuleType, MODULE_TYPE, LEGACY_MODULE_TYPE)
.eq(TaskResultItemEntity::getScopeHash, scopeHash)); .eq(TaskResultItemEntity::getScopeHash, scopeHash));
Map<String, TaskResultItemEntity> existingByKey = new HashMap<>(); Map<String, TaskResultItemEntity> existingByKey = new HashMap<>();
for (TaskResultItemEntity existing : existingList) { for (TaskResultItemEntity existing : existingList) {
@@ -72,6 +72,45 @@ class CollectDataResultItemBatchWriterTest {
assertEquals(1L, entities.get(0).getTaskId(), "task_id 保留"); assertEquals(1L, entities.get(0).getTaskId(), "task_id 保留");
} }
@Test
void test_module_type_matches_reader_constant() {
// 回归:写入端 module_type 必须与 CollectDataService.MODULE_TYPE(读取端)
// 写同一大小写值,否则生成结果文件时 loadFinalRows 查不出来(曾用小写
// "collectdata" 导致 task-50 之后所有采集任务明细行读不到)。
when(taskResultItemMapper.selectList(any())).thenReturn(List.of());
when(taskResultItemMapper.upsertBatch(anyList())).thenReturn(1);
writer.upsertAccepted(1L, 2L, "task:1", 0,
List.of(row("B000000001", "Nike")), "rustfs:detail/abc");
ArgumentCaptor<List> captor = ArgumentCaptor.forClass(List.class);
verify(taskResultItemMapper).upsertBatch(captor.capture());
TaskResultItemEntity entity = (TaskResultItemEntity) captor.getValue().get(0);
assertEquals("COLLECT_DATA", entity.getModuleType(), "写入 module_type 须为读取端大写常量");
}
@Test
void test_legacy_lowercase_rows_are_found_on_retry() {
// 回归:历史小写 "collectdata" 存量行在旧任务重提时也要被现有行查询
// 命中(按大小写兼容匹配),避免同一唯一键重复插入。
String refJson = codec.encodeRef(0, 0, "rustfs:detail/x");
TaskResultItemEntity legacy = new TaskResultItemEntity();
legacy.setId(100L);
legacy.setModuleType("collectdata");
legacy.setItemKey("asin:B000000001");
legacy.setPayloadJson(refJson);
legacy.setPayloadHash(sha256(refJson));
when(taskResultItemMapper.selectList(any())).thenReturn(List.of(legacy));
CollectDataResultItemBatchWriter.UpsertCounts counts =
writer.upsertAccepted(1L, 2L, "task:1", 0,
List.of(row("B000000001", "Nike")), "rustfs:detail/x");
assertEquals(0, counts.insertedOrUpdated(), "小写存量行命中 → 全部跳过");
assertEquals(1, counts.skipped(), "1 行跳过");
verify(taskResultItemMapper, never()).upsertBatch(anyList());
}
@Test @Test
void test_task_050_task_normal_multiple_items() { void test_task_050_task_normal_multiple_items() {
// 批量场景:75 行(8 批次),顺序稳定不丢失,每批数量正确。 // 批量场景:75 行(8 批次),顺序稳定不丢失,每批数量正确。