137 lines
5.7 KiB
Python
137 lines
5.7 KiB
Python
from __future__ import annotations
|
||
|
||
import logging
|
||
from datetime import datetime
|
||
from typing import Any, Optional
|
||
|
||
from app.features.personalized_reco.observability.types import MissingFields, RecoMeta, Scene
|
||
from app.features.personalized_reco.observability.utils import compute_empty_reason, compute_missing_fields
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
|
||
def _non_negative_int(value: Any, *, default: int = 0) -> int:
|
||
try:
|
||
n = int(value)
|
||
except Exception:
|
||
return int(default)
|
||
return max(0, int(n))
|
||
|
||
|
||
class RecoMetaBuilder:
|
||
"""
|
||
在推荐 pipeline 中逐阶段填充 RecoMeta,避免“散落字段/散落日志”。
|
||
|
||
说明(V1):
|
||
- set 调用允许任意顺序;build 时会做防御式兜底与单调性修正
|
||
- 单调性约束:raw >= after_hard_filter >= after_dedup >= after_freqcap >= served_k
|
||
"""
|
||
|
||
def __init__(self, *, scene: Scene, user_profile: object, k: int, now: Optional[datetime] = None) -> None:
|
||
self.scene: Scene = scene
|
||
self.user_profile = user_profile
|
||
self.k = _non_negative_int(k, default=0)
|
||
self.now = now
|
||
|
||
self._raw: Optional[int] = None
|
||
self._after_hard: Optional[int] = None
|
||
self._after_dedup: Optional[int] = None
|
||
self._after_freqcap: Optional[int] = None
|
||
self._served_k: Optional[int] = None
|
||
self._fallback_level_final: Optional[int] = None
|
||
|
||
self._risk_filtered_count_by_flag: dict[str, int] = {}
|
||
self._freqcap_filtered_counts: dict[str, int] = {}
|
||
self._config_snapshot: dict[str, Any] = {}
|
||
|
||
def set_candidate_pool_size_raw(self, n: Any) -> "RecoMetaBuilder":
|
||
self._raw = _non_negative_int(n)
|
||
return self
|
||
|
||
def set_after_hard_filter(self, n: Any, *, risk_filtered_count_by_flag: Optional[dict[str, Any]] = None) -> "RecoMetaBuilder":
|
||
self._after_hard = _non_negative_int(n)
|
||
if risk_filtered_count_by_flag:
|
||
self._risk_filtered_count_by_flag = {str(k): _non_negative_int(v) for k, v in risk_filtered_count_by_flag.items()}
|
||
return self
|
||
|
||
def set_after_dedup(self, n: Any) -> "RecoMetaBuilder":
|
||
self._after_dedup = _non_negative_int(n)
|
||
return self
|
||
|
||
def set_after_freqcap(self, n: Any, *, freqcap_filtered_counts: Optional[dict[str, Any]] = None) -> "RecoMetaBuilder":
|
||
self._after_freqcap = _non_negative_int(n)
|
||
if freqcap_filtered_counts:
|
||
self._freqcap_filtered_counts = {str(k): _non_negative_int(v) for k, v in freqcap_filtered_counts.items()}
|
||
return self
|
||
|
||
def set_fallback_level_final(self, level: Any, *, reason: Optional[str] = None) -> "RecoMetaBuilder":
|
||
# reason 预留,V1 先不入 meta(可放入 config_snapshot 或后续字段)
|
||
self._fallback_level_final = _non_negative_int(level, default=0)
|
||
if reason:
|
||
self._config_snapshot.setdefault("fallback_trigger_reason", str(reason))
|
||
return self
|
||
|
||
def set_served_k(self, n: Any) -> "RecoMetaBuilder":
|
||
self._served_k = _non_negative_int(n)
|
||
return self
|
||
|
||
def set_config_snapshot(self, snapshot: dict[str, Any]) -> "RecoMetaBuilder":
|
||
self._config_snapshot = dict(snapshot or {})
|
||
return self
|
||
|
||
def build(self) -> RecoMeta:
|
||
missing: MissingFields = compute_missing_fields(self.user_profile)
|
||
conf_u = getattr(self.user_profile, "profile_confidence", 1.0)
|
||
try:
|
||
conf_u_f = float(conf_u)
|
||
except Exception:
|
||
conf_u_f = 1.0
|
||
if conf_u_f != conf_u_f:
|
||
conf_u_f = 1.0
|
||
|
||
raw = self._raw if self._raw is not None else 0
|
||
after_hard = self._after_hard if self._after_hard is not None else raw
|
||
after_dedup = self._after_dedup if self._after_dedup is not None else after_hard
|
||
after_freqcap = self._after_freqcap if self._after_freqcap is not None else after_dedup
|
||
served_k = self._served_k if self._served_k is not None else 0
|
||
|
||
# 防御式单调性修正(以最保守值输出)
|
||
if after_hard > raw:
|
||
logger.debug("after_hard_filter(%s) > raw(%s),已修正为 raw", after_hard, raw)
|
||
after_hard = raw
|
||
if after_dedup > after_hard:
|
||
logger.debug("after_dedup(%s) > after_hard_filter(%s),已修正为 after_hard_filter", after_dedup, after_hard)
|
||
after_dedup = after_hard
|
||
if after_freqcap > after_dedup:
|
||
logger.debug("after_freqcap(%s) > after_dedup(%s),已修正为 after_dedup", after_freqcap, after_dedup)
|
||
after_freqcap = after_dedup
|
||
if served_k > after_freqcap:
|
||
logger.debug("served_k(%s) > after_freqcap(%s),已修正为 after_freqcap", served_k, after_freqcap)
|
||
served_k = after_freqcap
|
||
|
||
fallback_level_final = self._fallback_level_final if self._fallback_level_final is not None else 0
|
||
|
||
empty_reason = compute_empty_reason(
|
||
served_k=served_k,
|
||
candidate_pool_size_raw=raw,
|
||
candidate_pool_size_after_hard_filter=after_hard,
|
||
candidate_pool_size_after_freqcap=after_freqcap,
|
||
)
|
||
|
||
return RecoMeta(
|
||
scene=self.scene,
|
||
candidate_pool_size_raw=int(raw),
|
||
candidate_pool_size_after_hard_filter=int(after_hard),
|
||
candidate_pool_size_after_dedup=int(after_dedup),
|
||
candidate_pool_size_after_freqcap=int(after_freqcap),
|
||
fallback_level_final=int(fallback_level_final),
|
||
served_k=int(served_k),
|
||
empty_reason=empty_reason,
|
||
conf_U=float(conf_u_f),
|
||
missing_fields=missing,
|
||
risk_filtered_count_by_flag=dict(self._risk_filtered_count_by_flag),
|
||
freqcap_filtered_counts=dict(self._freqcap_filtered_counts),
|
||
config_snapshot=dict(self._config_snapshot),
|
||
)
|
||
|