本系列博客基于一个真实可运行的电商评论舆情分析项目。
上一篇:《第6篇:FastAPI 长任务异步化实践:从"HTTP 阻塞"到进程内 TaskManager》
一、痛点开场
舆情系统的价值在于"第一时间发现风险":差评率突然飙升、出现"假货/曝光/监管部门"等敏感词,运营需要马上知道,而不是等用户自己打开页面才发现。
所以预警不能只是"写进数据库等人看",而是要:
- 规则命中时自动落库(留痕、可追溯)
- 主动推送通知给负责人(飞书/钉钉)
- 支持人工确认闭环(处理完标记 closed)
这篇讲预警系统的实现:三类规则引擎 + Webhook 通知 + 通知器抽象。
二、三类规则引擎

backend/agents/alerting/nodes.py 定义了默认规则:
DEFAULT_RULES = {
"high_negative_rate": {"enabled": True, "threshold": 0.20}, # 负面率 > 20%
"spike_detection": {"enabled": True, "z_score_threshold": 2.0, "lookback_days": 7},
"keyword_alert": {"enabled": True, "keywords": ["投诉", "举报", "监管部门", "曝光", "骗", "假货"]},
}
# 命中即 critical 的高危词
CRITICAL_KEYWORDS = ["监管部门", "曝光", "骗", "假货", "传销", "非法"]
规则 1:高负面率
def _apply_high_negative_rate(rate, threshold):
if rate > threshold:
level = "critical" if rate > 0.40 else "warning" # >40% 升级 critical
return {
"rule": "high_negative_rate",
"level": level,
"title": f"负面率异常:{rate:.1%} 超过阈值 {threshold:.0%}",
"description": f"最近 24 小时内负面率 {rate:.1%},建议关注并介入处理。",
"metric_value": round(rate, 4),
"threshold": threshold,
}
return None
规则 2:突增检测
def _apply_spike_detection(current_rate, historical_rate, threshold):
if historical_rate is None or historical_rate <= 0:
return None
# 注意:这里其实是"相对变化比例",不是严格 z-score(见第 9 篇复盘)
change_ratio = (current_rate - historical_rate) / historical_rate
if change_ratio > threshold:
return {
"rule": "spike_detection",
"level": "warning",
"title": "负面率突增",
"description": f"负面率从 {historical_rate:.1%} 上升至 {current_rate:.1%},变化 {change_ratio:.1%}。",
...
}
return None
规则 3:敏感词命中
async def _query_keyword_alerts(tenant_id, window_hours):
# SELECT r.content, ra.sentiment_label FROM reviews r
# JOIN review_analyses ra ON ra.review_id = r.id
# WHERE tenant_id=:t AND review_time >= :cutoff AND (content ILIKE '%投诉%' OR ...)
for content, review_time, sentiment in rows:
matched_kw = [kw for kw in keywords if kw in content]
is_critical = any(kw in CRITICAL_KEYWORDS for kw in matched_kw)
alerts.append({
"rule": "keyword_alert",
"level": "critical" if is_critical else "info",
"title": f"敏感词命中:{'/'.join(matched_kw)}",
"description": f"评论包含敏感词:{content[:100]}...",
...
})
三、落库:alert_events
save_alerts_node 把预警写入数据库:
async def save_alerts_node(state):
alerts = state.get("alerts", [])
if not alerts:
return {"structured_output": _build_stats(state)}
async with AsyncSessionLocal() as session:
for alert in alerts:
await session.execute(
text("INSERT INTO alert_events "
"(tenant_id, level, rule, title, description, metric_value, threshold) "
"VALUES (:tenant_id, :level, :rule, :title, :description, :metric, :threshold)"),
{...})
await session.commit()
# 落库后再通知(关键顺序)
notify_results = await _notify_channels(alerts)
stats = _build_stats(state)
stats["notify"] = notify_results
return {"structured_output": stats}
四、主动通知:飞书/钉钉 Webhook
通知器抽象(backend/notifications/base.py)
@dataclass
class NotificationResult:
ok: bool
channel: str
message: str = ""
error: Optional[str] = None
http_status: Optional[int] = None
class BaseNotifier:
name: str = "base"
async def send(self, payload: dict) -> NotificationResult:
raise NotImplementedError
@staticmethod
def format_alert_card(alert: dict) -> dict:
"""把预警 dict 格式化为通用卡片。"""
level_emoji = {"critical": "🔴", "warning": "🟡", "info": "🔵"}
return {
"title": f"{level_emoji.get(alert.get('level'), '🔵')} {alert.get('title', '舆情预警')}",
"level": alert.get("level", "info"),
"rule": alert.get("rule", ""),
"description": alert.get("description", ""),
"metric_value": alert.get("metric_value"),
"threshold": alert.get("threshold"),
}
飞书通知器(backend/notifications/feishu.py)
class FeishuNotifier(BaseNotifier):
name = "feishu"
def __init__(self, webhook_url=None, timeout=10.0):
self.webhook_url = webhook_url or get_settings().feishu_webhook_url
self.timeout = timeout
async def send(self, payload) -> NotificationResult:
if not self.webhook_url:
return NotificationResult(ok=False, channel=self.name, error="FEISHU_WEBHOOK_URL 未配置")
# 飞书交互式卡片:header 颜色随级别变化
card = {
"msg_type": "interactive",
"card": {
"header": {
"title": {"tag": "plain_text", "content": payload.get("title", "舆情预警")},
"template": ("red" if payload.get("level") == "critical"
else "orange" if payload.get("level") == "warning" else "blue"),
},
"elements": [
{"tag": "div", "text": {"tag": "lark_md", "content": payload.get("description", "")}},
{"tag": "note", "elements": [{"tag": "plain_text",
"content": f"规则:{payload.get('rule', '-')} | 指标:{payload.get('metric_value', '-')} | 阈值:{payload.get('threshold', '-')}"}]},
],
},
}
return await self._post(card)
钉钉通知器(backend/notifications/dingtalk.py)
class DingTalkNotifier(BaseNotifier):
name = "dingtalk"
async def send(self, payload) -> NotificationResult:
if not self.webhook_url:
return NotificationResult(ok=False, channel=self.name, error="DINGTALK_WEBHOOK_URL 未配置")
body = {
"msgtype": "markdown",
"markdown": {
"title": payload.get("title", "舆情预警"),
"text": (f"## {payload.get('title')}\n\n"
f"**级别**:{payload.get('level')}\n\n"
f"**规则**:{payload.get('rule')}\n\n"
f"**说明**:{payload.get('description')}\n\n"
f"> 指标:{payload.get('metric_value')} | 阈值:{payload.get('threshold')}"),
},
}
return await self._post(body)
分渠道发送 + 失败不阻断
async def _notify_channels(alerts) -> dict:
results = {"feishu": [], "dingtalk": []}
if not alerts:
return results
try:
feishu = FeishuNotifier()
ding = DingTalkNotifier()
for alert in alerts:
card = BaseNotifier.format_alert_card(alert)
# 只推 critical + warning,info 不打扰
if alert.get("level") in ("critical", "warning"):
if feishu.webhook_url:
r = await feishu.send(card)
results["feishu"].append({"rule": alert.get("rule"), "ok": r.ok, "error": r.error})
if ding.webhook_url:
r = await ding.send(card)
results["dingtalk"].append({"rule": alert.get("rule"), "ok": r.ok, "error": r.error})
except Exception as e:
logger.warning("alerting.notify_failed", error=str(e)[:200])
return results
设计要点:
- 失败不阻断主流程:通知挂了不影响落库,用 try/except 包住
- 分级推送:只有 critical + warning 走 webhook,info 不打扰
- 结果回传:
stats["notify"]记录每个渠道的发送结果,可观测
五、调度与 API
- 定时扫描:APScheduler 每 6 小时执行一次
_run_alert_scan - 手动扫描:
POST /api/v1/alerts/scan(前端"立即扫描"按钮) - 列表:
GET /api/v1/alerts?status=active|closed - 确认:
POST /api/v1/alerts/confirm?alert_id=xxx(更新 status=closed + confirmed_by + confirmed_at)
六、诚实边界:它还不是"完整通知系统"
- 无去重/抑制:同一问题每 6 小时扫描都可能重复推送,运营会被轰炸
- 无通知记录表:没有落库"发给谁、是否成功",无法追溯
- 无签名校验:钉钉 webhook 未做 secret 签名,安全边界弱
- "z-score"名不副实:突增检测实际是相对变化比例,不是统计 z-score
生产化方案:
- 预警业务键
(tenant_id, rule, 资源维度, 时间窗口)唯一,命中已有 active 预警则跳过 - 抑制窗口:同一预警 60 分钟内只推一次,超时未解决自动升级重推
- 新增
notifications表记录每次发送的渠道/状态/错误 - 钉钉加
timestamp + sign签名
七、总结
- 三类规则:负面率、突增、敏感词,覆盖主要风险信号
- 落库 + 通知:先留痕再推送,失败不阻断
- 通知器抽象:
BaseNotifier+ 飞书/钉钉两个实现,可扩展 - 诚实边界:去重、抑制、记录表、签名是下一步
下一篇预告:《Vue3 + TS 舆情看板:ECharts 图表与前后端契约管理》
转载自 CSDN-专业IT技术社区
原文链接:https://blog.csdn.net/m0_68262196/article/details/164118336




