Files
2026-02-02 16:47:37 +08:00

368 lines
16 KiB
Markdown
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# 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`: intfeed 默认 30push/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`:直接返回空 itemsmeta.served_k=0fallback_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(...)`,输出结构一致