Browse Source

feat(indicator): NormalizeIndicatorPipeline 单一写者+三钩子(入库/编辑/食材快照)

iwt 2 days ago
parent
commit
6d3163bb4e

+ 3 - 0
cfc-backend/src/main/java/com/etotem/cfc/dto/ParsedReportPayload.java

@@ -104,6 +104,7 @@ public class ParsedReportPayload {
         private String description;
         private String category;
         private String level;
+        private String status;
 
         public String getBacteriaName() { return bacteriaName; }
         public void setBacteriaName(String bacteriaName) { this.bacteriaName = bacteriaName; }
@@ -121,6 +122,8 @@ public class ParsedReportPayload {
         public void setCategory(String category) { this.category = category; }
         public String getLevel() { return level; }
         public void setLevel(String level) { this.level = level; }
+        public String getStatus() { return status; }
+        public void setStatus(String status) { this.status = status; }
     }
 
     public static class FoodItem {

+ 13 - 0
cfc-backend/src/main/java/com/etotem/cfc/service/FoodRecommendService.java

@@ -38,6 +38,9 @@ public class FoodRecommendService {
     @Resource
     private HealthReportService healthReportService;
 
+    @Resource
+    private com.etotem.cfc.service.NormalizeIndicatorPipeline normalizeIndicatorPipeline;
+
     @Resource
     private FoodService foodService;
 
@@ -157,6 +160,16 @@ public class FoodRecommendService {
 
             log.info("用户{}食材{}已从最新报告中移除", userId, removed.getFoodId());
         }
+
+        // 归一化管道:全量快照最新食材推荐指数到 indicator_values(幂等覆盖)
+        try {
+            List<FoodRecommendIndex> finalIndices = foodRecommendIndexMapper.selectList(
+                    new LambdaQueryWrapper<FoodRecommendIndex>()
+                            .eq(FoodRecommendIndex::getUserId, userId));
+            normalizeIndicatorPipeline.snapshotFoodIndices(userId, finalIndices);
+        } catch (Exception e) {
+            log.warn("食材指数归一化快照失败 userId={}: {}", userId, e.getMessage());
+        }
     }
 
     public List<FoodRecommendIndex> getIndexByUser(Long userId) {

+ 12 - 0
cfc-backend/src/main/java/com/etotem/cfc/service/HealthReportService.java

@@ -112,6 +112,9 @@ public class HealthReportService {
     @Resource
     private FoodMapper foodMapper;
 
+    @Resource
+    private com.etotem.cfc.service.NormalizeIndicatorPipeline normalizeIndicatorPipeline;
+
     /**
      * 创建健康报告(含指标列表)
      */
@@ -1912,6 +1915,15 @@ public class HealthReportService {
         } catch (Exception e) {
             log.warn("编辑后刷新展示块失败 reportId={}: {}", reportId, e.getMessage());
         }
+
+        // 归一化管道:编辑后重写统一指标层(同事务;payload 为编辑后最新值)
+        try {
+            ParsedReportPayload.Payload editedPayload = reportBlockAssembler.fromMap(payload);
+            oldReport.setSubjectId(subjectId != null && subjectId > 0 ? subjectId : oldSubjectId);
+            normalizeIndicatorPipeline.processReport(oldReport, editedPayload);
+        } catch (Exception e) {
+            log.warn("归一化管道刷新失败 reportId={}: {}", reportId, e.getMessage());
+        }
     }
 
     private Integer toInteger(Object value) {

+ 130 - 0
cfc-backend/src/main/java/com/etotem/cfc/service/NormalizeIndicatorPipeline.java

@@ -0,0 +1,130 @@
+package com.etotem.cfc.service;
+
+import com.etotem.cfc.dto.ParsedReportPayload;
+import com.etotem.cfc.entity.FoodRecommendIndex;
+import com.etotem.cfc.entity.HealthReport;
+import com.etotem.cfc.entity.IndicatorValue;
+import com.etotem.cfc.mapper.IndicatorValueMapper;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.stereotype.Service;
+
+import javax.annotation.Resource;
+import java.math.BigDecimal;
+import java.sql.Date;
+import java.util.List;
+
+/**
+ * 归一化管道 —— 统一指标层「单一写者」。
+ * 报告入库/编辑流程照旧写 health_indicators(不动现有消费者),
+ * 本管道在同流程/同事务内将 指标+菌属+食材 幂等写入 indicator_values。
+ * 唯一键 (source_type, source_id, definition_id) 保证重解析只覆盖更新。
+ */
+@Service
+public class NormalizeIndicatorPipeline {
+
+    private static final Logger log = LoggerFactory.getLogger(NormalizeIndicatorPipeline.class);
+
+    @Resource
+    private IndicatorMappingService mappingService;
+
+    @Resource
+    private IndicatorValueMapper valueMapper;
+
+    /**
+     * 报告指标 + 菌属归一化写入。
+     * 调用点:insertReport()(subjectId 块后)/ updateReportFromPayload()(步骤 5 后)。
+     * 无 subject 的报告不入时序(无法归属观测对象)。
+     *
+     * @param report  已入库/已更新的报告(reportDate 须为真实日期或 null)
+     * @param payload 解析载荷(gutFlora 从 payload 读;indicators 从 payload 读,与 health_indicators 同源)
+     */
+    public void processReport(HealthReport report, ParsedReportPayload.Payload payload) {
+        if (report == null || report.getSubjectId() == null || report.getId() == null) {
+            return;
+        }
+        Long subjectId = report.getSubjectId();
+        Long familyId = report.getFamilyId();
+        Long reportId = report.getId();
+        String sourceType = report.getReportType();
+        Date reportDate = report.getReportDate() != null
+                ? new Date(report.getReportDate().getTime()) : null;
+
+        int written = 0;
+        if (payload != null && payload.getIndicators() != null) {
+            for (ParsedReportPayload.Indicator ind : payload.getIndicators()) {
+                if (ind.getIndicatorName() == null || ind.getIndicatorName().trim().isEmpty()) {
+                    continue;
+                }
+                Long defId = mappingService.resolve(sourceType, ind.getIndicatorName().trim());
+                if (defId == null) {
+                    continue;
+                }
+                IndicatorValue v = baseValue(sourceType, reportId, defId, subjectId, familyId, reportDate);
+                v.setValue(ind.getIndicatorValue());
+                v.setNumericValue(mappingService.parseNumericValue(ind.getIndicatorValue()));
+                v.setUnit(ind.getUnit());
+                v.setStatus(ind.getStatus());
+                valueMapper.upsertMapping(v);
+                written++;
+            }
+        }
+        if (payload != null && payload.getGutFlora() != null) {
+            for (ParsedReportPayload.Flora flora : payload.getGutFlora()) {
+                if (flora.getBacteriaName() == null || flora.getBacteriaName().trim().isEmpty()) {
+                    continue;
+                }
+                Long defId = mappingService.resolve(sourceType, flora.getBacteriaName().trim());
+                if (defId == null) {
+                    continue;
+                }
+                IndicatorValue v = baseValue(sourceType, reportId, defId, subjectId, familyId, reportDate);
+                v.setValue(flora.getBacteriaValue());
+                v.setNumericValue(mappingService.parseNumericValue(flora.getBacteriaValue()));
+                v.setUnit("%");
+                v.setStatus(flora.getStatus());
+                valueMapper.upsertMapping(v);
+                written++;
+            }
+        }
+        log.info("归一化管道 reportId={}, subjectId={}, type={},写入 {} 条观测", reportId, subjectId, sourceType, written);
+    }
+
+    /**
+     * 食材推荐指数全量快照(钩子 3:FoodRecommendService.recalculateForUser 尾部调用)。
+     * source_type=food_recommend,source_id=food_recommend_idx.id,幂等覆盖。
+     */
+    public void snapshotFoodIndices(Long userId, List<FoodRecommendIndex> indices) {
+        if (userId == null || indices == null || indices.isEmpty()) {
+            return;
+        }
+        for (FoodRecommendIndex idx : indices) {
+            if (idx.getFoodName() == null || idx.getFoodName().trim().isEmpty()) {
+                continue;
+            }
+            Long defId = mappingService.resolve("food_recommend", idx.getFoodName().trim());
+            if (defId == null) {
+                continue;
+            }
+            IndicatorValue v = baseValue("food_recommend", idx.getId(), defId, userId, null, null);
+            v.setValue(idx.getIndexScore() != null ? idx.getIndexScore().toString() : null);
+            v.setNumericValue(idx.getIndexScore() != null
+                    ? BigDecimal.valueOf(idx.getIndexScore()) : null);
+            v.setReportId(idx.getHealthReportId());
+            valueMapper.upsertMapping(v);
+        }
+        log.info("食材指数快照 userId={}, 共 {} 条", userId, indices.size());
+    }
+
+    private IndicatorValue baseValue(String sourceType, Long sourceId, Long defId,
+                                     Long subjectId, Long familyId, Date reportDate) {
+        IndicatorValue v = new IndicatorValue();
+        v.setSourceType(sourceType);
+        v.setSourceId(sourceId);
+        v.setDefinitionId(defId);
+        v.setSubjectId(subjectId);
+        v.setFamilyId(familyId);
+        v.setReportDate(reportDate);
+        return v;
+    }
+}

+ 10 - 0
cfc-backend/src/main/java/com/etotem/cfc/service/ReportCollectService.java

@@ -95,6 +95,9 @@ public class ReportCollectService {
     @Resource
     private ObjectMapper objectMapper;
 
+    @Resource
+    private com.etotem.cfc.service.NormalizeIndicatorPipeline normalizeIndicatorPipeline;
+
     @Value("${report.collect.opencode-base:http://127.0.0.1:4090}")
     private String opencodeBase;
 
@@ -648,6 +651,13 @@ public class ReportCollectService {
             log.warn("保存展示块失败 reportId={}", created.getId(), e);
         }
 
+        // 归一化管道:统一指标层幂等写入(同流程,subjectId 非空才有效)
+        try {
+            normalizeIndicatorPipeline.processReport(created, payload);
+        } catch (Exception e) {
+            log.warn("归一化管道写入失败 reportId={}: {}", created.getId(), e.getMessage());
+        }
+
         healthReportDraftService.markCollectCompleted(draft.getId(), created.getId());
 
         // 新手任务联动:完成首次报告上传任务(老用户无任务记录时静默跳过)

+ 86 - 0
cfc-backend/src/test/java/com/etotem/cfc/service/NormalizeIndicatorPipelineTest.java

@@ -0,0 +1,86 @@
+package com.etotem.cfc.service;
+
+import com.etotem.cfc.dto.ParsedReportPayload;
+import com.etotem.cfc.entity.HealthReport;
+import com.etotem.cfc.entity.IndicatorValue;
+import com.etotem.cfc.mapper.IndicatorValueMapper;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.Date;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.*;
+
+class NormalizeIndicatorPipelineTest {
+
+    private NormalizeIndicatorPipeline pipeline;
+    private IndicatorMappingService mappingService;
+    private IndicatorValueMapper valueMapper;
+
+    @BeforeEach
+    void setUp() {
+        pipeline = new NormalizeIndicatorPipeline();
+        mappingService = Mockito.mock(IndicatorMappingService.class);
+        valueMapper = Mockito.mock(IndicatorValueMapper.class);
+        ReflectionTestUtils.setField(pipeline, "mappingService", mappingService);
+        ReflectionTestUtils.setField(pipeline, "valueMapper", valueMapper);
+    }
+
+    @Test
+    void processReportWritesIndicatorsAndFlora() {
+        HealthReport report = new HealthReport();
+        report.setId(10L);
+        report.setFamilyId(2L);
+        report.setSubjectId(3L);
+        report.setReportType("gut_flora");
+        report.setReportDate(new Date());
+        when(mappingService.resolve(eq("gut_flora"), eq("双歧杆菌属"))).thenReturn(1L);
+
+        ParsedReportPayload.Payload payload = new ParsedReportPayload.Payload();
+        ParsedReportPayload.Indicator ind = new ParsedReportPayload.Indicator();
+        ind.setIndicatorName("双歧杆菌属");
+        ind.setIndicatorValue("12.8");
+        ind.setUnit("%");
+        ind.setStatus("偏高");
+        payload.setIndicators(Collections.singletonList(ind));
+
+        pipeline.processReport(report, payload);
+        verify(valueMapper).upsertMapping(any(IndicatorValue.class));
+    }
+
+    @Test
+    void processReportSkipsWhenNoSubject() {
+        HealthReport report = new HealthReport();
+        report.setId(10L);
+        report.setReportType("gut_flora");
+        report.setSubjectId(null);
+        ParsedReportPayload.Payload payload = new ParsedReportPayload.Payload();
+        payload.setIndicators(Collections.emptyList());
+
+        pipeline.processReport(report, payload);
+        verify(valueMapper, never()).upsertMapping(any());
+    }
+
+    @Test
+    void snapshotFoodIndicesWritesFoodValues() {
+        com.etotem.cfc.entity.FoodRecommendIndex idx = new com.etotem.cfc.entity.FoodRecommendIndex();
+        idx.setId(5L);
+        idx.setUserId(3L);
+        idx.setFoodId(9L);
+        idx.setFoodName("燕麦");
+        idx.setIndexScore(88);
+        idx.setHealthReportId(10L);
+        idx.setCalculatedAt(new Date());
+        when(mappingService.resolve(eq("food_recommend"), eq("燕麦"))).thenReturn(6L);
+
+        pipeline.snapshotFoodIndices(3L, java.util.Collections.singletonList(idx));
+        verify(valueMapper).upsertMapping(argThat(v -> v.getSourceType().equals("food_recommend")
+                && v.getSourceId().equals(5L)
+                && v.getDefinitionId().equals(6L)));
+    }
+}