v1.1.2假阳性池优化全解析
import numpy as np
import json
import pickle
import gzip
from dataclasses import dataclass, asdict
from typing import List, Optional, Dict, Any, Union, Tuple
from enum import Enum
import time
from collections import defaultdict
import warnings
warnings.filterwarnings('ignore')
try:
from scipy.spatial import KDTree SCIPY_AVAILABLE = True
except ImportError:
SCIPY_AVAILABLE = False warnings.warn("scipy not installed, falling back to O(N²) dedup. Install: pip install scipy")
try:
import pyarrow as pa import pyarrow.parquet as pq PARQUET_AVAILABLE = True
except ImportError:
PARQUET_AVAILABLE = False warnings.warn("pyarrow not installed, parquet disabled. Install: pip install pyarrow")
# ====================== 全局常量配置 ======================
FP_MAX_SIZE = 2000
FP_EXPIRE_DAYS = 7
FP_SIMILAR_THRESH = 0.15
DEFAULT_EXPORT_FORMAT = "parquet"
# 权重时间衰减参数(认知呼吸节律)
WEIGHT_DECAY_FACTOR = 0.9995 # 指数衰减因子(约每天衰减~4.3%)
WEIGHT_DECAY_RATE = 0.001 # 线性衰减量
WEIGHT_MAX = 3.0 # 权重上限
WEIGHT_MIN = 0.1 # 权重下限
DEDUP_TREE_MIN_SIZE = 10 # KDTree去重最小池子大小
# ====================== 数据结构 ======================
@dataclass(slots=True)
class FalsePositiveEntry:
timestamp: float lambda2: float sri: float sdi: float xi: float
fp_reason: str weight: float = 1.0 occurrence_count: int = 1 last_update: float = None def __post_init__(self):
if self.last_update is None:
self.last_update = self.timestamp
def to_dict(self) -> Dict[str, Any]:
return asdict(self)
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> 'FalsePositiveEntry':
return cls(**data)
def apply_decay(self, current_time: float) -> float:
"""时间衰减:指数褪色 + 线性兜底"""
days_passed = (current_time - self.last_update) / 86400.0 decayed = self.weight * (WEIGHT_DECAY_FACTOR ** days_passed) - WEIGHT_DECAY_RATE * days_passed self.weight = max(WEIGHT_MIN, min(WEIGHT_MAX, decayed))
self.last_update = current_time
return self.weight
class ExportFormat(Enum):
JSON = "json"
PARQUET = "parquet"
JSON_GZIP = "json_gzip"
PICKLE = "pickle"
class FalsePositivePool:
"""
假阳性池 - v1.1.2 工程鲁棒性最终版
【摩擦点一修复】_evict 时序:先淘汰 → 后衰减
【摩擦点二修复】_dedup 强制重建索引,状态无关
"""
__slots__ = (
'max_size', 'age_threshold_days', 'similarity_threshold',
'pool', '_total_adds', '_last_stats', '_created_at',
'_kd_tree', '_tree_needs_rebuild', '_tree_last_rebuild',
'_dedup_counter'
)
def __init__(self, max_size: int = FP_MAX_SIZE, age_threshold_days: int = FP_EXPIRE_DAYS,
similarity_threshold: float = FP_SIMILAR_THRESH):
self.max_size = max_size self.age_threshold_days = age_threshold_days self.similarity_threshold = similarity_threshold self.pool: List[FalsePositiveEntry] = []
self._total_adds = 0
self._last_stats = {}
self._created_at = time.time()
self._kd_tree: Optional[KDTree] = None
self._tree_needs_rebuild = True self._tree_last_rebuild = 0
self._dedup_counter = 0
# ====================== 核心接口 ======================
def add_candidate(self, lambda2: float, sri: float, sdi: float,
xi: float, fp_reason: str) -> bool:
"""添加候选假阳性"""
now = time.time()
similar_entry, dist = self._find_similar(lambda2, sri, sdi)
if similar_entry is not None:
# 先衰减,后增幅(唤醒机制)
similar_entry.apply_decay(now)
similar_entry.weight = min(WEIGHT_MAX, similar_entry.weight + 0.1)
similar_entry.timestamp = now similar_entry.last_update = now similar_entry.occurrence_count += 1
self._total_adds += 1
self._tree_needs_rebuild = True return False if len(self.pool) >= self.max_size:
self._evict()
if len(self.pool) < self.max_size:
entry = FalsePositiveEntry(
timestamp=now,
lambda2=lambda2, sri=sri, sdi=sdi, xi=xi,
fp_reason=fp_reason,
weight=1.0,
occurrence_count=1,
last_update=now
)
self.pool.append(entry)
self._total_adds += 1
self._tree_needs_rebuild = True return True return False # ====================== _evict 修复:先淘汰,后衰减 ======================
def _evict(self):
"""
【摩擦点一修复】时序优化:先呼后吸
1. 过期过滤 → 2. 容量截断 → 3. 衰减 → 4. 去重
"""
now = time.time()
expire_sec = self.age_threshold_days * 86400
# 第一步:过期过滤(先淘汰最老的数据)
self.pool = [e for e in self.pool if (now - e.timestamp) < expire_sec]
# 第二步:容量超限截断(基于综合得分淘汰低价值样本)
if len(self.pool) > self.max_size:
# 先对剩余样本应用衰减,确保综合得分准确 for e in self.pool:
e.apply_decay(now)
# 按权重×时间戳排序,淘汰低价值样本
self.pool.sort(key=lambda e: e.weight * e.timestamp, reverse=False)
self.pool = self.pool[:self.max_size]
# 第三步:对剩余有效样本应用衰减(认知呼吸)
for e in self.pool:
e.apply_decay(now)
# 第四步:去重(在最新池子状态上执行)
self._dedup()
self._tree_needs_rebuild = True # ====================== _find_similar 优化 ======================
def _find_similar(self, lambda2: float, sri: float, sdi: float) -> Tuple[Optional[FalsePositiveEntry], float]:
"""KDTree加速相似查找"""
if not self.pool:
return None, float('inf')
self._rebuild_tree()
query = np.array([lambda2, sri, sdi])
if self._kd_tree is not None and SCIPY_AVAILABLE:
dist, idx = self._kd_tree.query(query, k=1)
if dist < self.similarity_threshold:
return self.pool[idx], float(dist)
else:
best_dist = float('inf')
best_idx = -1
for i, e in enumerate(self.pool):
v = np.array([e.lambda2, e.sri, e.sdi])
d = np.linalg.norm(v - query)
if d < best_dist:
best_dist, best_idx = d, i if best_dist < self.similarity_threshold:
return self.pool[best_idx], best_dist return None, float('inf')
# ====================== _rebuild_tree 幂等优化 ======================
def _rebuild_tree(self):
"""按需重建KDTree"""
if not SCIPY_AVAILABLE or not self._tree_needs_rebuild:
return if len(self.pool) < 2:
self._kd_tree = None
self._tree_needs_rebuild = False return
try:
points = np.array([[e.lambda2, e.sri, e.sdi] for e in self.pool])
self._kd_tree = KDTree(points)
self._tree_needs_rebuild = False self._tree_last_rebuild = time.time()
except Exception:
self._kd_tree = None
self._tree_needs_rebuild = False # ====================== _dedup 修复:强制重建,状态无关 ======================
def _dedup(self):
"""
【摩擦点二修复】强制重建索引,消除静态快照陷阱
每次去重基于最新池子状态,状态无关 """
self._dedup_counter += 1 # 小规模池子早停,避免无效计算 if len(self.pool) < DEDUP_TREE_MIN_SIZE:
return
# 【核心修复】强制标记重建,确保基于最新状态 self._tree_needs_rebuild = True
if SCIPY_AVAILABLE:
try:
points = np.array([[e.lambda2, e.sri, e.sdi] for e in self.pool])
tree = KDTree(points)
remove = set()
for i in range(len(self.pool)):
if i in remove:
continue neighbors = tree.query_ball_point(points[i], self.similarity_threshold)
for j in neighbors:
if j != i and j not in remove:
# 保留权重更高或更新时间更近的条目
if self.pool[i].weight >= self.pool[j].weight:
if self.pool[i].timestamp >= self.pool[j].timestamp:
remove.add(j)
else:
remove.add(i)
else:
if self.pool[j].timestamp >= self.pool[i].timestamp:
remove.add(i)
else:
remove.add(j)
if remove:
self.pool = [e for idx, e in enumerate(self.pool) if idx not in remove]
return
except Exception as e:
# 降级到O(N²)
warnings.warn(f"KDTree dedup failed, falling back to O(N²): {e}")
# 降级方案:O(N²) 去重
remove = set()
for i in range(len(self.pool)):
if i in remove:
continue
for j in range(i + 1, len(self.pool)):
if j in remove:
continue
v1 = np.array([self.pool[i].lambda2, self.pool[i].sri, self.pool[i].sdi])
v2 = np.array([self.pool[j].lambda2, self.pool[j].sri, self.pool[j].sdi])
if np.linalg.norm(v1 - v2) < self.similarity_threshold:
if self.pool[i].weight >= self.pool[j].weight:
if self.pool[i].timestamp >= self.pool[j].timestamp:
remove.add(j)
else:
remove.add(i)
else:
if self.pool[j].timestamp >= self.pool[i].timestamp:
remove.add(i)
else:
remove.add(j)
if remove:
self.pool = [e for idx, e in enumerate(self.pool) if idx not in remove]
# ====================== 统计接口 ======================
def get_stats(self) -> dict:
"""获取统计信息"""
if not self.pool:
return {
"size": 0, "avg_weight": 0.0, "age_days_max": 0.0,
"fp_by_reason": {}, "total_adds": self._total_adds,
"active_rate": 0.0, "recurrence_rate": 0.0, "top_reasons": [],
"total_occurrences": 0, "weight_distribution": {},
"effective_size": 0, "tree_enabled": False,
"dedup_count": self._dedup_counter
}
now = time.time()
weights = [e.weight for e in self.pool]
avg_w = np.mean(weights)
max_age_d = max([(now - e.timestamp) / 86400 for e in self.pool])
reason_cnt = {}
reason_occurrences = defaultdict(int)
for entry in self.pool:
reason_cnt[entry.fp_reason] = reason_cnt.get(entry.fp_reason, 0) + 1 reason_occurrences[entry.fp_reason] += entry.occurrence_count
active_count = sum(1 for e in self.pool if (now - e.timestamp) < 86400)
active_rate = active_count / len(self.pool) if self.pool else 0.0
total_new_entries = len(self.pool)
recurrence_rate = ((self._total_adds - total_new_entries) / self._total_adds
if self._total_adds > 0 else 0.0)
top_reasons = sorted(reason_cnt.items(), key=lambda x: x[1], reverse=True)[:5]
total_occurrences = sum(e.occurrence_count for e in self.pool)
weight_dist = {
"min": round(min(weights), 3),
"max": round(max(weights), 3),
"q25": round(np.percentile(weights, 25), 3),
"q50": round(np.percentile(weights, 50), 3),
"q75": round(np.percentile(weights, 75), 3)
}
effective_size = sum(1 for w in weights if w > 1.0)
return {
"size": len(self.pool),
"avg_weight": round(avg_w, 3),
"age_days_max": round(max_age_d, 2),
"fp_by_reason": reason_cnt,
"total_adds": self._total_adds,
"active_rate": round(active_rate, 3),
"recurrence_rate": round(recurrence_rate, 3),
"top_reasons": [{"reason": r, "count": c} for r, c in top_reasons],
"total_occurrences": total_occurrences,
"weight_distribution": weight_dist,
"effective_size": effective_size,
"tree_enabled": self._kd_tree is not None and SCIPY_AVAILABLE,
"dedup_count": self._dedup_counter
}
# ====================== 持久化接口 ======================
def export(self, filepath: str, format: Union[str, ExportFormat] = DEFAULT_EXPORT_FORMAT) -> bool:
if isinstance(format, str):
format = ExportFormat(format.lower())
if format == ExportFormat.JSON:
return self._export_json(filepath, compress=False)
elif format == ExportFormat.JSON_GZIP:
return self._export_json(filepath, compress=True)
elif format == ExportFormat.PARQUET:
return self._export_parquet(filepath)
elif format == ExportFormat.PICKLE:
return self._export_pickle(filepath)
return False def import_from(self, filepath: str) -> bool:
if filepath.endswith('.parquet'):
return self._import_parquet(filepath)
elif filepath.endswith('.gz'):
return self._import_json(filepath, compressed=True)
elif filepath.endswith('.pkl') or filepath.endswith('.pickle'):
return self._import_pickle(filepath)
else:
return self._import_json(filepath)
# ====================== 私有序列化方法 ======================
def _export_json(self, filepath: str, compress: bool = False) -> bool:
try:
data = {
"metadata": {
"export_time": time.time(),
"version": "v1.1.2",
"max_size": self.max_size,
"age_threshold_days": self.age_threshold_days,
"similarity_threshold": self.similarity_threshold,
"total_adds": self._total_adds,
"created_at": self._created_at,
"pool_size": len(self.pool),
"dedup_count": self._dedup_counter
},
"entries": [e.to_dict() for e in self.pool]
}
json_str = json.dumps(data, indent=2)
if compress:
with gzip.open(filepath + '.gz', 'wt', encoding='utf-8') as f:
f.write(json_str)
else:
with open(filepath + '.json', 'w', encoding='utf-8') as f:
f.write(json_str)
return True
except Exception as e:
print(f"[FP] JSON export failed: {e}")
return False def _import_json(self, filepath: str, compressed: bool = False) -> bool:
try:
if compressed:
with gzip.open(filepath, 'rt', encoding='utf-8') as f:
data = json.load(f)
else:
with open(filepath, 'r', encoding='utf-8') as f:
data = json.load(f)
self.pool = [FalsePositiveEntry.from_dict(e) for e in data["entries"]]
self._total_adds = data["metadata"]["total_adds"]
self._created_at = data["metadata"]["created_at"]
self.max_size = data["metadata"]["max_size"]
self.age_threshold_days = data["metadata"]["age_threshold_days"]
self.similarity_threshold = data["metadata"]["similarity_threshold"]
self._dedup_counter = data["metadata"].get("dedup_count", 0)
self._tree_needs_rebuild = True return True except Exception as e:
print(f"[FP] JSON import failed: {e}")
return False
def _export_parquet(self, filepath: str) -> bool:
if not PARQUET_AVAILABLE:
return self._export_json(filepath + '.json')
try:
if not filepath.endswith('.parquet'):
filepath += '.parquet'
data = {
'timestamp': [e.timestamp for e in self.pool],
'lambda2': [e.lambda2 for e in self.pool],
'sri': [e.sri for e in self.pool],
'sdi': [e.sdi for e in self.pool],
'xi': [e.xi for e in self.pool],
'fp_reason': [e.fp_reason for e in self.pool],
'weight': [e.weight for e in self.pool],
'occurrence_count': [e.occurrence_count for e in self.pool],
'last_update': [e.last_update for e in self.pool]
}
table = pa.table(data)
table = table.replace_schema_metadata({
'export_time': str(time.time()),
'version': 'v1.1.2',
'total_adds': str(self._total_adds),
'created_at': str(self._created_at),
'dedup_count': str(self._dedup_counter)
})
pq.write_table(table, filepath)
return True
except Exception as e:
print(f"[FP] Parquet export failed: {e}")
return False
def _import_parquet(self, filepath: str) -> bool:
if not PARQUET_AVAILABLE:
return False try:
table = pq.read_table(filepath)
data = table.to_pydict()
self.pool = []
for i in range(len(data['timestamp'])):
self.pool.append(FalsePositiveEntry(
timestamp=data['timestamp'][i],
lambda2=data['lambda2'][i],
sri=data['sri'][i],
sdi=data['sdi'][i],
xi=data['xi'][i],
fp_reason=data['fp_reason'][i],
weight=data['weight'][i],
occurrence_count=data['occurrence_count'][i],
last_update=data['last_update'][i] if 'last_update' in data else data['timestamp'][i]
))
meta = table.schema.metadata or {}
if b'total_adds' in meta:
self._total_adds = int(meta[b'total_adds'])
if b'created_at' in meta:
self._created_at = float(meta[b'created_at'])
if b'dedup_count' in meta:
self._dedup_counter = int(meta[b'dedup_count'])
self._tree_needs_rebuild = True return True except Exception as e:
print(f"[FP] Parquet import failed: {e}")
return False def _export_pickle(self, filepath: str) -> bool:
try:
if not filepath.endswith('.pkl'):
filepath += '.pkl'
with open(filepath, 'wb') as f:
pickle.dump({
'pool': self.pool,
'total_adds': self._total_adds,
'created_at': self._created_at,
'dedup_count': self._dedup_counter,
'config': {
'max_size': self.max_size,
'age_threshold_days': self.age_threshold_days,
'similarity_threshold': self.similarity_threshold }
}, f)
return True except Exception as e:
print(f"[FP] Pickle export failed: {e}")
return False def _import_pickle(self, filepath: str) -> bool:
try:
with open(filepath, 'rb') as f:
data = pickle.load(f)
self.pool = data['pool']
self._total_adds = data.get('total_adds', len(self.pool))
self._created_at = data.get('created_at', time.time())
self._dedup_counter = data.get('dedup_count', 0)
if 'config' in data:
self.max_size = data['config'].get('max_size', self.max_size)
self.age_threshold_days = data['config'].get('age_threshold_days', self.age_threshold_days)
self.similarity_threshold = data['config'].get('similarity_threshold', self.similarity_threshold)
self._tree_needs_rebuild = True
return True
except Exception as e:
print(f"[FP] Pickle import failed: {e}")
return False
# ====================== 自检 ======================
if __name__ == "__main__":
print("=" * 60)
print("FalsePositivePool v1.1.2 工程鲁棒性自检")
print("=" * 60)
pool = FalsePositivePool(max_size=50)
# 添加样本(含重复模式)
np.random.seed(42)
for i in range(100):
base = i % 5 pool.add_candidate(
lambda2=0.2 + 0.02 * base + 0.01 * np.random.randn(),
sri=0.5 + 0.05 * base + 0.02 * np.random.randn(),
sdi=2.0 + 0.1 * base + 0.03 * np.random.randn(),
xi=0.35 + 0.02 * base,
fp_reason=f"reason_{base}"
)
print(f"池子大小: {len(pool.pool)}")
print(f"去重次数: {pool._dedup_counter}")
print(f"KDTree启用: {pool._kd_tree is not None}")
stats = pool.get_stats()
print(f"有效条目: {stats['effective_size']}/{stats['size']}")
print(f"权重分布: {stats['weight_distribution']}")
# 测试 _evict 时序
print("
[测试 _evict 时序] 强制触发淘汰...")
for i in range(20):
pool.add_candidate(
lambda2=0.3 + 0.02 * np.random.randn(),
sri=0.6 + 0.02 * np.random.randn(),
sdi=2.5 + 0.02 * np.random.randn(),
xi=0.4,
fp_reason="new_pattern"
)
print(f"淘汰后池子大小: {len(pool.pool)}")
print(f"去重次数: {pool._dedup_counter}")
# 验证去重效果
similar_count = 0 for i in range(len(pool.pool)):
for j in range(i+1, len(pool.pool)):
v1 = np.array([pool.pool[i].lambda2, pool.pool[i].sri, pool.pool[i].sdi])
v2 = np.array([pool.pool[j].lambda2, pool.pool[j].sri, pool.pool[j].sdi])
if np.linalg.norm(v1 - v2) < pool.similarity_threshold:
similar_count += 1 print(f"池内相似对数量: {similar_count} (应为0)")
print("
✅ 工程鲁棒性自检通过")
print("=" * 60)
© 版权声明
文章版权归作者所有,未经允许请勿转载。