用AI做舆情监控:从社交媒体爬取到情感分析的实时Pipeline
一、场景痛点与技术挑战
企业舆情监控是品牌安全的底线防线。
一条负面消息在30分钟内就能引爆全网。
传统人工巡检根本跟不上信息传播速度。
社交媒体日均产生数亿条内容。
人工筛选耗时且遗漏率高。
核心痛点有四个。
一是数据源碎片化。
微博、抖音、小红书、Twitter各成孤岛。
每个平台的API和页面结构完全不同。
二是实时性要求极高。
舆情窗口期仅30-60分钟。
延迟1小时的报告基本失去价值。
三是情感分析精度不足。
中文语境下讽刺、反话、隐喻难以识别。
简单词典匹配误判率超过40%。
四是数据量与算力矛盾。
日处理百万条文本需要弹性算力调度。
峰值流量是均值3-5倍。
技术挑战更棘手。
反爬机制越来越严格。
验证码、IP限制、登录墙层层叠加。
数据清洗噪音大。
广告、水军、机器人内容占比超过30%。
模型推理延迟与吞吐量平衡困难。
实时Pipeline需要毫秒级响应。
二、核心原理与架构设计
实时舆情Pipeline分五个核心模块。
采集模块负责多平台数据获取。
每个平台一个独立爬虫worker。
微博用API+页面混合采集。
小红书和抖音通过移动端接口。
Twitter用官方API v2流式端点。
所有worker将原始数据推入Kafka。
清洗模块负责数据去噪。
规则引擎过滤广告和水军内容。
正则匹配+模式识别去除营销模板。
机器人检测通过行为特征识别。
发帖频率、内容重复度、账号属性综合判断。
清洗后的数据写入Elasticsearch。
预处理模块负责文本标准化。
中文分词用jieba+自定义行业词典。
去停用词、统一缩写、修复错别字。
提取关键词和实体(品牌名、人名、地名)。
实体识别用BERT-based NER模型。
情感分析模块是核心决策层。
三级情感分类:正面、负面、中性。
细粒度情感:愤怒、失望、担忧、期待。
讽刺检测用专用子模型。
输出情感置信度而非硬标签。
置信度低于阈值的内容标记为待人工复核。
告警模块负责实时通知。
负面情感超过阈值触发告警。
同一话题在短时间内密集出现触发告警。
告警级别:关注、预警、紧急。
紧急告警同步推送钉钉和企业微信。
三、生产级代码实现
多平台数据采集框架
"""多平台舆情数据采集框架""" import asyncio import hashlib import json import logging import time from abc import ABC, abstractmethod from dataclasses import dataclass, field logger = logging.getLogger("sentiment_collector") @dataclass class RawPost: platform: str post_id: str content: str author: str timestamp: float metadata: dict = field(default_factory=dict) class BaseCollector(ABC): """爬虫基类:统一采集接口""" def __init__(self, platform: str, rate_limit: float = 1.0): self.platform = platform self.rate_limit = rate_limit # 请求间隔秒数 self._last_request = 0.0 async def _throttle(self): """请求限速""" elapsed = time.time() - self._last_request if elapsed < self.rate_limit: await asyncio.sleep(self.rate_limit - elapsed) self._last_request = time.time() @abstractmethod async def fetch_recent(self, keywords: list[str], since: float) -> list[RawPost]: """获取最近帖子""" ... class WeiboCollector(BaseCollector): """微博数据采集器""" def __init__(self, cookie: str, rate_limit: float = 2.0): super().__init__("weibo", rate_limit) self.cookie = cookie self.headers = { "User-Agent": "Mozilla/5.0...", "Cookie": cookie, } async def fetch_recent(self, keywords: list[str], since: float) -> list[RawPost]: posts = [] for kw in keywords: await self._throttle() # 微博搜索API模拟 url = f"https://m.weibo.cn/api/container/getContainer?" url += f"containerid=100103type%3D1%26q%3D{kw}" # 实际生产中用aiohttp请求 # resp = await aiohttp_get(url, headers=self.headers) # data = json.loads(resp) # 这里模拟数据返回 mock_data = [ { "id": hashlib.md5(f"{kw}_{i}".encode()).hexdigest()[:12], "text": f"关于{kw}的讨论内容{i}", "user": f"user_{i}", "created_at": time.time() - i * 60, } for i in range(5) ] for item in mock_data: if item["created_at"] >= since: posts.append(RawPost( platform=self.platform, post_id=item["id"], content=item["text"], author=item["user"], timestamp=item["created_at"], )) return posts class TwitterCollector(BaseCollector): """Twitter数据采集器(API v2)""" def __init__(self, bearer_token: str, rate_limit: float = 1.0): super().__init__("twitter", rate_limit) self.bearer = bearer_token async def fetch_recent(self, keywords: list[str], since: float) -> list[RawPost]: posts = [] query = " OR ".join(keywords) # Twitter API v2 filtered stream # 实际生产中用streaming endpoint mock_data = [ { "id": f"tw_{hashlib.md5(query.encode()).hexdigest()[:8]}", "text": f"Twitter discussion about {query}", "user": f"tw_user_{i}", "created_at": time.time() - i * 120, } for i in range(3) ] for item in mock_data: if item["created_at"] >= since: posts.append(RawPost( platform=self.platform, post_id=item["id"], content=item["text"], author=item["user"], timestamp=item["created_at"], )) return posts class CollectorScheduler: """采集调度器:协调多平台Worker""" def __init__(self, collectors: list[BaseCollector]): self.collectors = collectors async def collect_all(self, keywords: list[str], since: float) -> list[RawPost]: """并行采集所有平台数据""" tasks = [ c.fetch_recent(keywords, since) for c in self.collectors ] results = await asyncio.gather(*tasks, return_exceptions=True) all_posts = [] for result in results: if isinstance(result, Exception): logger.error(f"采集失败: {result}") continue all_posts.extend(result) logger.info(f"本轮采集: {len(all_posts)}条") return all_posts数据清洗与情感分析Pipeline
"""数据清洗与情感分析Pipeline""" import re import hashlib from dataclasses import dataclass from enum import Enum from collections import Counter class SentimentLevel(Enum): POSITIVE = "正面" NEGATIVE = "负面" NEUTRAL = "中性" UNCERTAIN = "待复核" class EmotionTag(Enum): ANGER = "愤怒" DISAPPOINT = "失望" WORRY = "担忧" EXPECT = "期待" SATISFY = "满意" @dataclass class AnalyzedPost: post_id: str content: str platform: str sentiment: SentimentLevel sentiment_score: float # -1.0 ~ 1.0 emotion: EmotionTag | None keywords: list[str] entities: list[str] is_spam: bool confidence: float # 0.0 ~ 1.0 class DataCleaner: """数据清洗:去除噪音内容""" # 广告关键词模式 AD_PATTERNS = [ r"加微信|加VX|扫码领取|限时优惠|点击购买", r"代购|批发价|工厂直供|最低价|免费送", ] # 水军行为特征阈值 SPAM_THRESHOLD = { "max_posts_per_hour": 20, "min_unique_ratio": 0.3, } _ad_regex = None _content_history: dict[str, list[str]] = {} def __init__(self): self._ad_regex = re.compile( "|".join(self.AD_PATTERNS), re.IGNORECASE ) def is_ad(self, content: str) -> bool: """广告内容检测""" return bool(self._ad_regex.search(content)) def is_water_army(self, author: str, content: str) -> bool: """水军检测:基于行为特征""" self._content_history.setdefault(author, []).append(content) history = self._content_history[author] if len(history) > self.SPAM_THRESHOLD["max_posts_per_hour"]: return True unique_ratio = len(set(history)) / max(len(history), 1) if unique_ratio < self.SPAM_THRESHOLD["min_unique_ratio"]: return True return False def clean(self, post) -> AnalyzedPost | None: """清洗单条数据""" if self.is_ad(post.content): return None # 广告内容直接丢弃 is_spam = self.is_water_army(post.author, post.content) if is_spam: return None # 水军内容丢弃 return AnalyzedPost( post_id=post.post_id, content=post.content, platform=post.platform, sentiment=SentimentLevel.NEUTRAL, sentiment_score=0.0, emotion=None, keywords=[], entities=[], is_spam=False, confidence=0.0, ) class SentimentAnalyzer: """情感分析引擎""" # 简化版情感词典(生产环境用BERT模型) NEGATIVE_WORDS = { "垃圾": -0.8, "糟糕": -0.7, "失望": -0.6, "气愤": -0.9, "投诉": -0.5, "退货": -0.4, "骗": -0.85, "坑": -0.75, "差": -0.5, "黑心": -0.9, "抵制": -0.7, "曝光": -0.6, } POSITIVE_WORDS = { "好": 0.5, "赞": 0.7, "满意": 0.6, "推荐": 0.65, "优秀": 0.8, "喜欢": 0.6, "棒": 0.7, "支持": 0.55, "靠谱": 0.65, } # 讽刺模式识别 IRONY_PATTERNS = [ r"真是.*的好", # "真是垃圾的好" r"所谓.*其实", # "所谓高端其实低端" r"还.*呢", # "还好呢" ] def __init__(self, confidence_threshold: float = 0.6): self.threshold = confidence_threshold self._irony_regex = re.compile( "|".join(self.IRONY_PATTERNS) ) def detect_irony(self, text: str) -> bool: """讽刺检测""" return bool(self._irony_regex.search(text)) def analyze(self, post: AnalyzedPost) -> AnalyzedPost: """分析情感""" text = post.content # 计算情感得分 score = 0.0 matched_words = 0 for word, weight in self.NEGATIVE_WORDS.items(): if word in text: score += weight matched_words += 1 for word, weight in self.POSITIVE_WORDS.items(): if word in text: score += weight matched_words += 1 # 讽刺检测翻转情感 if self.detect_irony(text) and score > 0: score = -score * 0.8 # 计算置信度 confidence = min(matched_words / 3.0, 1.0) # 情感分类 if confidence < self.threshold: post.sentiment = SentimentLevel.UNCERTAIN elif score > 0.2: post.sentiment = SentimentLevel.POSITIVE elif score < -0.2: post.sentiment = SentimentLevel.NEGATIVE else: post.sentiment = SentimentLevel.NEUTRAL post.sentiment_score = round(score, 3) post.confidence = round(confidence, 2) # 情感标签 if post.sentiment == SentimentLevel.NEGATIVE: if score < -0.7: post.emotion = EmotionTag.ANGER elif score < -0.5: post.emotion = EmotionTag.DISAPPOINT else: post.emotion = EmotionTag.WORRY return post class AlertEngine: """舆情告警引擎""" def __init__(self, negative_threshold: float = 0.3, cluster_threshold: int = 10): self.negative_threshold = negative_threshold self.cluster_threshold = cluster_threshold self._topic_counter: Counter = Counter() def evaluate(self, posts: list[AnalyzedPost]) -> list[dict]: """评估告警条件""" alerts = [] negative_count = 0 for post in posts: if post.sentiment == SentimentLevel.NEGATIVE: negative_count += 1 # 提取话题关键词用于聚类 for kw in post.keywords or [post.content[:10]]: self._topic_counter[kw] += 1 # 负面比例告警 if len(posts) > 0: ratio = negative_count / len(posts) if ratio > self.negative_threshold: level = "紧急" if ratio > 0.6 else "预警" if ratio > 0.4 else "关注" alerts.append({ "type": "negative_ratio", "level": level, "ratio": round(ratio, 3), "total": len(posts), "negative": negative_count, }) # 话题密集告警 hot_topics = [ (topic, count) for topic, count in self._topic_counter.items() if count >= self.cluster_threshold ] for topic, count in hot_topics: alerts.append({ "type": "topic_cluster", "level": "预警" if count > 20 else "关注", "topic": topic, "count": count, }) return alerts四、性能优化与工程实践
采集层性能优化。
每个平台Worker独立进程运行。
进程间通过Kafka Topic解耦。
Worker崩溃不影响其他平台采集。
请求限速通过令牌桶算法实现。
突发流量时排队而非粗暴拒绝。
Kafka分区策略按平台划分。
每个平台一个分区,避免数据交叉。
消费者组按分析模块划分。
清洗、情感分析、告警各一个消费组。
并行消费提升处理吞吐量。
清洗层优化。
广告过滤用预编译正则。
一次编译多次匹配,避免重复开销。
水军检测用滑动窗口。
只保留最近1小时的发帖历史。
超过窗口的历史数据自动清理。
情感分析模型优化。
生产环境用BERT-based模型。
词典匹配仅作为低置信度的快速通道。
BERT模型用ONNX Runtime加速推理。
batch_size=32的批量推理吞吐量提升5倍。
GPU推理延迟<10ms/条。
CPU推理用模型蒸馏版本延迟<50ms/条。
讽刺检测是精度瓶颈。
词典匹配的讽刺检测覆盖率仅60%。
生产环境需要专项训练讽刺数据集。
中文讽刺语料稀缺是最大障碍。
自建讽刺标注数据集至少5000条。
标注质量比数量更重要。
告警引擎优化。
负面比例告警用滑动窗口计算。
窗口大小5分钟,步长1分钟。
避免单条数据触发误告警。
话题聚类用TF-IDF+余弦相似度。
相似度>0.7的内容归入同一话题。
五、总结与技术提炼
多平台采集框架用抽象基类统一接口。
每个平台一个Worker,独立进程运行。
请求限速用令牌桶,采集调度器并行协调。数据清洗三层过滤:广告、水军、机器人。
广告用正则模式匹配,水军用行为特征检测。
机器人用发帖频率+内容重复度综合判断。情感分析引擎分层设计。
快速通道用词典匹配,高精度通道用BERT模型。
讽刺检测翻转正负情感,置信度低于阈值标记待复核。Kafka解耦采集与处理层。
按平台分区,按模块划分消费组。
Worker崩溃不影响Pipeline其他环节。告警引擎双维度触发。
负面比例超过阈值触发比例告警。
同一话题密集出现触发聚类告警。
滑动窗口避免单条误触发。模型推理用ONNX Runtime加速。
GPU批量推理吞吐量5倍提升。
讽刺数据集自建是精度突破的关键。
5000条高质量标注优于10000条低质标注。