16 KiB
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 | widgetuser_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: inttext: strfinal_score: floatfallback_level_final: intexplanations: 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.pyasync def recommend(...) -> RecoEngineResultasync def recommend_one(...)(push/widget 便捷入口)
hard_filter.pydef hard_filter(...) -> HardFilterResult(返回 kept + 统计 + reasons)
types.pyRecoConstraints、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。
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_itemsrisk_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 = 0fallback_level>=3:personalization_power = 0
实现方式建议:
- 在引擎内对
ContentProfileDTO做浅拷贝(或model_copy(update={...}))后再传入score_content。
7.3 explanations(可选)
为便于调参/排查,建议支持按开关输出 explanations:
hard_filter_hits:命中的 flag / predicatescore_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 = Truefeed_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_rawset_after_hard_filter(..., risk_filtered_count_by_flag=...)set_after_dedupset_after_freqcap(..., freqcap_filtered_counts=...)set_served_kset_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==0hard_filter_all:raw>0 且 after_hard==0freqcap_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正确
- raw=0 →
- 去重生效:
- 输出不包含 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(...),输出结构一致
- FastAPI 与 Celery 调用同一