# Reco Engine(推荐引擎编排)|Plan > 对应规范:`spec_kit/Personalized Reco/modules/reco-engine/spec.md` > > 规则来源(必须严格对齐): > > - `设计说明文档/個性化推薦算法規則.md`(Pipeline、回退梯度、场景差异、Hard Filter 关键规则) > - `设计说明文档/句子文案打分規則.md`(risk_flags 命名与语义唯一准绳;需与 DB→DTO 归一化一致) > > 依赖模块(已实现): > > - `server/app/features/personalized_reco/content_repository/`(候选拉取) > - `server/app/features/personalized_reco/scoring/`(软打分) > - `server/app/features/personalized_reco/rerank_freqcap/`(去重/重排/频控) > - `server/app/features/personalized_reco/observability/`(统一 meta 构建与 empty_reason 口径) --- ## 1. 目标与交付物 ### 1.1 目标 - 实现推荐主编排器(Orchestrator),将候选拉取、硬过滤、软打分、重排/频控、回退梯度串成一个稳定 Pipeline。 - 任意输入(字段缺失、历史为空/很大、候选不足)均不报错,并返回结构稳定的 `items + meta`。 - 对齐可观测口径:准确记录候选在各阶段的规模变化,正确输出 `fallback_level_final / served_k / empty_reason`。 - 保持“无框架耦合”:同一引擎可被 FastAPI 与 Celery 调用。 ### 1.2 交付物 - `modules/reco-engine/plan.md`:本技术计划(本文件)。 - 代码实现(tasks 阶段落地)建议位置: - `server/app/features/personalized_reco/reco_engine/` - 包含: - 编排器:`orchestrator.py` - 硬过滤:`hard_filter.py` - 类型与配置:`types.py`、`defaults.py` - (可选)同步封装:`sync.py`(供 Celery 直接调用) - 单元测试(tasks 阶段落地)建议位置: - `server/tests/test_reco_engine.py` --- ## 2. 模块职责边界(V1 约定) ### 2.1 本模块负责 - **候选拉取编排**:调用 `ContentRepository.fetch_candidates(...)`,并按回退层级控制拉取策略与上限。 - **Hard Filter(硬过滤)**:按 risk_flags 与跨维度产品规则剔除高风险内容,并输出按 flag 聚合的统计。 - **Soft Scoring(软打分)编排**:调用 `scoring.score_content(...)`,并根据场景/回退层级/画像缺失控制配置开关(例如 Push 强制启用 `P_uncertainty`)。 - **Rerank/Freqcap(重排/频控)编排**:调用 `rerank_freqcap.rerank_and_freqcap(...)`,并将其 meta 写入统一 `RecoMeta`。 - **Fallback Ladder(回退梯度)**:实现 L0→L3 逐级回退与补齐策略(尤其 Feed 可配置是否继续回退补齐)。 - **统一输出结构**:`items: List[RecommendedItem]` + `meta: RecoMeta`(来自 `RecoMetaBuilder`)。 ### 2.2 本模块不负责 - 数据库 schema 与 ORM(由 `db-design` 与 `content_repository` 负责)。 - risk_flags 旧→新映射、suitability 默认值补齐(由 `content_repository.normalization` 负责)。 - 软打分的公式实现(由 `scoring.score_content` 负责)。 - 去重/频控/Feed MMR 具体算法实现(由 `rerank_freqcap.rerank_and_freqcap` 负责)。 - 打点上报/落库(由调用方:API/Worker 负责;本模块只生成可观测 `meta`)。 --- ## 3. 输入/输出与数据结构(V1) ### 3.1 编排器输入(对齐 spec,并补齐工程必需字段) 规范 `spec.md` 输入基础上,为满足 `ContentRepository` 的强约束,本模块额外引入 `locale`: - `scene`: `feed | push | widget` - `user_profile`: `UserProfileV1_2`(允许字段缺失/跳过) - `already_recommended_ids`: `List[str|int]` - `touched_or_viewed_ids`: `List[str|int]` - `k`: int(feed 默认 30;push/widget 默认 1) - `now`: 时间戳(`datetime`) - `locale`: `en | tc`(必填;不允许语言回退;若未传则由上层决定默认值) - (可选)`constraints`: - `exclude_content_ids`: `List[int]`(额外排除;会与 already/touched 合并) - `exclude_author_ids`: `List[str]` - `exclude_template_ids`: `List[str]` - `max_candidates_limit`: int(候选池上限;用于保护数据库与后续计算) - `recent_author_ids` / `recent_template_ids`(用于 Push/Widget 增强频控;不提供则由 `rerank_freqcap` 记录缺失并跳过该维度过滤) > 说明:`ContentRepository.fetch_candidates` 已内置“缺失字段 → 至少 L1”的降级约束;但引擎仍需在回退循环中显式维护 `fallback_level`,以便可观测与一致性。 ### 3.2 输出(对齐 spec) - `items: List[RecommendedItem]`(长度 ≤ k) - `content_id: int` - `text: str` - `final_score: float` - `fallback_level_final: int` - `explanations: Optional[dict]`(可选;用于调参/排查;默认可关闭以节省载荷) - `meta: RecoMeta` - 统一结构来自 `observability.RecoMetaBuilder.build()` ### 3.3 推荐结果建议类型(tasks 阶段落地) - `RecommendedItem`:pydantic model 或 dataclass(建议 pydantic,与现有 `RecoMeta` 风格一致)。 - `RecoEngineResult`:`items + meta` 的容器类型(便于 API/Worker 复用)。 --- ## 4. 总体架构与代码组织(建议) 建议新增目录:`server/app/features/personalized_reco/reco_engine/` - `orchestrator.py` - `async def recommend(...) -> RecoEngineResult` - `async def recommend_one(...)`(push/widget 便捷入口) - `hard_filter.py` - `def hard_filter(...) -> HardFilterResult`(返回 kept + 统计 + reasons) - `types.py` - `RecoConstraints`、`RecommendedItem`、`RecoEngineResult`、`HardFilterResult` - `defaults.py` - 场景默认参数(例如候选拉取上限、Feed 是否允许回退补齐等) - `utils.py` - 小工具:id 归一化、personalization_power clamp、解释字段构造等 --- ## 5. Pipeline 设计(Candidate → Hard Filter → Soft Scoring → Rerank/Freqcap → Serve) ### 5.1 主流程伪代码(V1) 核心思想:**回退循环包裹整个 Pipeline**,每次回退都重新拉候选并重新跑一遍 pipeline;最终输出 `fallback_level_final` 与 `meta`。 ```text meta_builder = RecoMetaBuilder(scene, user_profile, k, now) fallback_trace = [] exclude_ids = union(already_recommended_ids, touched_or_viewed_ids, constraints.exclude_content_ids) for level in [0, 1, 2, 3]: # 1) Candidate cands = await repo.fetch_candidates(scene, user_profile, fallback_level=level, limit=candidate_limit(level), locale, exclude_content_ids=exclude_ids) meta_builder.set_candidate_pool_size_raw(len(cands)) # 2) Hard Filter kept, risk_counts, hard_removed = hard_filter(scene, user_profile, cands, constraints) meta_builder.set_after_hard_filter(len(kept), risk_filtered_count_by_flag=risk_counts) # 3) Soft Scoring scored = [] for each content in kept: cfg = scoring_config(scene, level, user_profile) content2 = clamp_personalization_power_if_needed(content, level) s = score_content(scene, user_profile, content2, config=cfg, pass_filters=True, external_terms=optional) scored.append(ScoredCandidate.from(content2, final_score=s.final_score)) # 4) Rerank/Freqcap rer = rerank_and_freqcap(scene, scored, already_recommended_ids, touched_or_viewed_ids, k, recent_author_ids, recent_template_ids) meta_builder.set_after_dedup(rer.meta.candidate_pool_size_after_dedup) meta_builder.set_after_freqcap(rer.meta.candidate_pool_size_after_freqcap, freqcap_filtered_counts=rer.meta.freqcap_filtered_counts) served = rer.ranked_items[:k] meta_builder.set_served_k(len(served)) meta_builder.set_fallback_level_final(level, reason=trigger_reason_if_any) fallback_trace.append({level, raw, after_hard, after_dedup, after_freqcap, served_k}) if len(served) == k: break if scene == "feed" and allow_partial_feed and len(served) > 0 and not fill_with_fallback: break # else continue fallback to try fill meta_builder.set_config_snapshot({"fallback_trace": fallback_trace, ...}) return items=served_as_recommended_items, meta=meta_builder.build() ``` ### 5.2 候选拉取策略(与回退梯度一致) 依赖 `ContentRepository.fetch_candidates(...)`: - `fallback_level=0`:正常配比(由 repository 内部实现候选策略;引擎只传 level) - `fallback_level>=1`:降个性化(repository 已约束 `personalization_power<=0.5`) - `fallback_level>=2`:回退通用池(repository 已约束 `general + personalization_power=0`) - `fallback_level>=3`:仅安全池(repository 已约束 `is_safe_pool=true`) 候选拉取上限: - 建议 `limit = min(max_candidates_limit, k * multiplier)`,默认 `multiplier=10`(Feed)/`multiplier=30`(Push/Widget,因强过滤+频控更容易清空)。 - `content_repository` 内部已有 `raw_limit = limit * 5` 的二次扩增,reco-engine 层的 `limit` 需以“软上限”思路控制资源。 --- ## 6. Hard Filter(硬过滤)设计 ### 6.1 规则集合(V1 必做) 对每条候选 `Cᵢ`,若命中任一规则则过滤: - **全场景必挡**: - `block_health_medical`(注意:旧 flag 归一化已在 repository 做;引擎只消费归一化后的 `risk_flags`) - **与用户阶段相关**: - 若 `U.stage.unknown=1`:过滤含 `unsafe_for_stage_unknown` - 若 `U.stage.parenting=1`:过滤含 `unsafe_for_stage_parenting` - **与用户情绪相关**: - 若 `U.emotion_score <= 0.2`:过滤含 `unsafe_for_emotion_low` - **跨维度产品规则(示例,来自算法规则文档)**: - 若 `U.stage.unknown=1` 且 `C.need_suitability[parenting_pressure]=1` 且 `C.personalization_power=1`:过滤 > 说明:Hard Filter 只做“剔除”,不做分数惩罚;软风险(例如 `soft_health_sensitive`)应由 `scoring` 的外部项 `P_risk` 或未来扩展处理(V1 可先不实现软风险)。 ### 6.2 与 `UserProfileV1_2_Extended.hard_rules` 的兼容(增强项) 若调用方传入的 `user_profile` 带有 `hard_rules`(扩展画像),引擎应: - 合并 `forbidden_risk_flags` 到本模块默认 forbidden 集合(并做去重)。 - 执行 `forbidden_content_predicates`(以“用户条件 + 内容字段命中”方式过滤),并将命中 predicate 的 `id` 记录到 explanations(可选)或 `meta.config_snapshot`。 ### 6.3 输出统计(用于 meta) Hard Filter 必须输出: - `kept_items` - `risk_filtered_count_by_flag: dict[str, int]`(按 flag 聚合计数,供 `RecoMetaBuilder.set_after_hard_filter(..., risk_filtered_count_by_flag=...)`) - (可选)`filtered_by_rule_ids: dict[str, int]`(跨维度规则命中计数,可放 `config_snapshot`) --- ## 7. Soft Scoring 编排策略(V1) ### 7.1 配置选择 默认使用 `scoring.get_default_config(scene)`,并按以下规则在引擎侧做“安全覆盖”: - **Push**:强制 `enable_uncertainty_penalty=True`(与 spec 对齐)。 - **任意场景**:当 `missing_fields` 明显或 `conf_U` 偏低时,可选择开启 `enable_uncertainty_penalty`(V1 可先只对 Push 强制,Feed/Widget 保持默认)。 ### 7.2 回退层级对个性化强度的约束 尽管 repository 已在候选拉取阶段约束 personalization_power,但为保证“防御式一致性”,引擎应再做一次 clamp: - `fallback_level>=1`:`personalization_power = min(personalization_power, 0.5)` - `fallback_level>=2`:`personalization_power = 0` - `fallback_level>=3`:`personalization_power = 0` 实现方式建议: - 在引擎内对 `ContentProfileDTO` 做浅拷贝(或 `model_copy(update={...})`)后再传入 `score_content`。 ### 7.3 explanations(可选) 为便于调参/排查,建议支持按开关输出 `explanations`: - `hard_filter_hits`:命中的 flag / predicate - `score_breakdown`:来自 `ScoreResult.breakdown`(注意载荷大小,默认关闭) - `fallback_level_used` --- ## 8. Rerank/Freqcap 编排策略(V1) 依赖 `rerank_freqcap.rerank_and_freqcap(...)`: - **去重**:使用 `already_recommended_ids ∪ touched_or_viewed_ids`(模块内部已归一化为 int set) - **Feed**:`dedup + MMR`(`mmr_lambda=0.7`,`top_n_for_mmr=200` 默认) - **Push/Widget**:`dedup + freqcap(句子/作者/模板) + TopK` - 句子冷却由 `already/touched` 直接提供即可生效 - 作者/模板冷却需要 `recent_author_ids/recent_template_ids` 输入;若缺失,模块会记录 `missing_history_fields` 并跳过该维度过滤(但仍不会报错) 引擎侧需要把 `RerankResult.meta` 写入统一 `RecoMetaBuilder`: - `set_after_dedup(rer.meta.candidate_pool_size_after_dedup)` - `set_after_freqcap(rer.meta.candidate_pool_size_after_freqcap, freqcap_filtered_counts=rer.meta.freqcap_filtered_counts)` --- ## 9. Fallback Ladder(回退梯度)实现细节 ### 9.1 触发条件(对齐 spec) 任一满足即可进入下一层回退: - 候选池为空 / Hard Filter 清空 / 去重清空 / 频控清空 - `served_k < k` - Feed:允许“部分不足”,但需记录;是否继续回退补齐由配置控制 - Push/Widget:建议默认继续回退直到 `served_k==k` 或达到 L3 ### 9.2 Feed 的“部分不足”策略(建议默认) 提供引擎配置项(`RecoEngineConfig`): - `feed_allow_partial: bool = True` - `feed_fill_with_fallback: bool = True` 推荐默认:Feed 允许部分不足,但仍尝试回退补齐(更接近“稳定覆盖率”目标);若担心回退导致风格突变,可关闭补齐。 ### 9.3 回退过程可观测(建议) 由于 `RecoMeta` 为单结构,建议把每次回退的过程写入 `meta.config_snapshot`: - `fallback_trace: List[{"level": int, "raw": int, "after_hard": int, "after_dedup": int, "after_freqcap": int, "served_k": int}]` - `fallback_trigger_reason`:最后一次触发原因(也可放每层 reason) --- ## 10. 可观测 meta 构建与 empty_reason 口径 使用 `observability.RecoMetaBuilder` 统一生成 meta: - 初始化:`RecoMetaBuilder(scene=scene, user_profile=user_profile, k=k, now=now)` - 每阶段 set: - `set_candidate_pool_size_raw` - `set_after_hard_filter(..., risk_filtered_count_by_flag=...)` - `set_after_dedup` - `set_after_freqcap(..., freqcap_filtered_counts=...)` - `set_served_k` - `set_fallback_level_final(level, reason=...)` - `set_config_snapshot({"fallback_trace": ..., "engine_config": ...})` - 最终:`meta = builder.build()` empty_reason: - 由 `observability.compute_empty_reason(...)` 在 `build()` 内计算(无需引擎手动写入) - 关键在于引擎必须正确设置 `raw/after_hard/after_freqcap/served_k`,以便区分: - `pool_empty`:raw==0 - `hard_filter_all`:raw>0 且 after_hard==0 - `freqcap_all`:raw>0 且 after_freqcap==0(并且 after_hard>0) - `unknown`:其他异常情况 --- ## 11. 稳定性与错误处理(V1) ### 11.1 防御式输入处理 - `k<=0`:直接返回空 items,meta.served_k=0,fallback_level_final=0。 - `already_recommended_ids / touched_or_viewed_ids`:允许混合类型(str/int),统一按 int 解析(无效值忽略)。 - `locale`:由 `content_repository.types.normalize_locale` 约束;若不支持,建议在上层拦截;引擎内部需捕获异常并返回空结果(避免 500)。 ### 11.2 异常兜底 任何阶段发生异常: - 不抛出到调用方(除非调用方明确要求),而是返回: - `items=[]` - `meta`:尽可能填充已知字段,`config_snapshot` 记录错误信息(例如 `{"error": "...", "stage": "fetch_candidates"}`) - 目的:保证 API/Worker 稳定,不因单条数据问题导致任务/请求失败。 --- ## 12. 测试计划(对应验收标准) ### 12.1 单元测试覆盖 - **稳定性**: - 缺失字段组合(need/context/emotion 任意缺失)不报错 - 历史集合为空/很大(包含非数字 id)不报错 - **回退可观测**: - raw=0 → `empty_reason="pool_empty"` - raw>0 且 after_hard=0 → `empty_reason="hard_filter_all"` - raw>0 且 after_freqcap=0 且 after_hard>0 → `empty_reason="freqcap_all"` - fallback_trace 写入且 `fallback_level_final` 正确 - **去重生效**: - 输出不包含 already/touched 中的 id(覆盖 feed/push/widget) - **风险优先**: - `block_health_medical` 必挡(全场景) - unknown stage + `unsafe_for_stage_unknown` 必挡 - **跨调用复用**: - 同样输入(固定 now)重复调用结果稳定(允许 score 浮点微差) ### 12.2 集成测试建议(tasks 阶段可选) - 在 `integration-api-worker` 完成后: - FastAPI 与 Celery 调用同一 `recommend(...)`,输出结构一致