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

16 KiB
Raw Permalink Blame History

Reco Engine推荐引擎编排Plan

对应规范:spec_kit/Personalized Reco/modules/reco-engine/spec.md

规则来源(必须严格对齐):

  • 设计说明文档/個性化推薦算法規則.mdPipeline、回退梯度、场景差异、Hard Filter 关键规则)
  • 设计说明文档/句子文案打分規則.mdrisk_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.pydefaults.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 与 ORMdb-designcontent_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 阶段落地)

  • RecommendedItempydantic model 或 dataclass建议 pydantic与现有 RecoMeta 风格一致)。
  • RecoEngineResultitems + 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
    • RecoConstraintsRecommendedItemRecoEngineResultHardFilterResult
  • 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_finalmeta

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=10Feed/multiplier=30Push/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=1C.need_suitability[parenting_pressure]=1C.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_penaltyV1 可先只对 Push 强制Feed/Widget 保持默认)。

7.2 回退层级对个性化强度的约束

尽管 repository 已在候选拉取阶段约束 personalization_power但为保证“防御式一致性”引擎应再做一次 clamp

  • fallback_level>=1personalization_power = min(personalization_power, 0.5)
  • fallback_level>=2personalization_power = 0
  • fallback_level>=3personalization_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
  • Feeddedup + MMRmmr_lambda=0.7top_n_for_mmr=200 默认)
  • Push/Widgetdedup + 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_emptyraw==0
    • hard_filter_allraw>0 且 after_hard==0
    • freqcap_allraw>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(...),输出结构一致