refactor: extract Python package hotspot/ from .dspy files

- hotspot/parsers.py: 11 builtin parsers + 4 extractors (json_path/css/regex/rss)
- hotspot/engine.py: fetch + extract orchestration
- hotspot/analysis.py: 5-dim scoring + 7-day half-life status classifier
- hotspot/db.py: DB helpers (save_items, write_fetch_log, get_stats)
- hotspot/__init__.py: public API (run_fetch, run_analysis, get_stats)
- .dspy files reduced to 1-line wrappers
This commit is contained in:
yumoqing 2026-08-01 12:13:00 +08:00
parent 42dff9e77d
commit 4a8f2c22be
9 changed files with 741 additions and 912 deletions

View File

@ -1,92 +1,35 @@
# 热点雷达 (Hotspot Radar)
Sage 框架模块 — 指定来源的热点聚合分析平台,支持动态添加来源
Sage 框架模块 — 指定来源的热点聚合分析平台。
## 功能
- 11 个内置平台解析器(微博/知乎/百度/头条/B站/抖音/36氪/GitHub/HN/V2EX/RSSHub
- 4 种可配置提取器json_path / css / regex / rss
- 每条来源自带 parser_config同类型不同来源可用不同提取规则
- 五维分析 + 7 天半衰期自动状态分类
- **多类型来源**API / RSS / Browser (HTML抓取) / Crawler (递归爬虫)
- **调度管理**cron 定时 + 间隔模式retry 策略
- **抓取日志**:每次执行一条审计记录(耗时/成功数/错误/HTTP状态码
- **热点分类**5级状态 — emerging → rising → hot → cooling → expired
- **五维分析**:时效 / 热度 / 内容 / 传播 / 受众,每维度 0-100 分
- **预警规则**:热度阈值 / 加速度 / 情感触发email/webhook/sms 通知
## 目录结构
## 目录
```
hotspot/
├── json/ # Sage 表定义(放入 Sage 的 json/ 目录)
│ ├── hotspot_source.json
│ ├── hotspot_schedule.json
│ ├── hotspot_fetch_log.json
│ ├── hotspot_item.json
│ ├── hotspot_analysis.json
│ └── hotspot_alert.json
├── wwwroot/hotspot/ # Web UI + 后端端点(放入 Sage 的 wwwroot/hotspot/
│ ├── index.ui # 仪表盘7 统计卡片 + 6 标签页)
│ ├── stats.dspy # 统计 API
│ ├── fetch_now.dspy # 抓取引擎3种模式
│ └── analyze.dspy # 分析引擎5维 + 状态分类)
├── hotspot/ # Python 包(放入 Sage 根目录)
│ ├── __init__.py
│ ├── engine.py # 抓取引擎
│ ├── analysis.py # 五维分析引擎
│ ├── parsers.py # 11 内置解析器 + 通用提取器
│ └── db.py # 数据库操作
├── json/ # 表定义(放入 Sage json/ 目录)
├── wwwroot/hotspot/ # UI + 端点(放入 Sage wwwroot/hotspot/
├── ddl/mysql.sql # 建表 SQL
├── install/ # 安装片段
│ ├── load_path_append.txt # 追加到 load_path.py
│ └── menu_append.json # 追加到 menu.ui
└── README.md
└── install/ # 安装片段
```
## 安装
### 1. 复制文件到 Sage 部署目录
```bash
SAGE_ROOT=/path/to/sage
cp -r hotspot $SAGE_ROOT/ # Python 包
cp json/hotspot_*.json $SAGE_ROOT/json/
cp -r wwwroot/hotspot $SAGE_ROOT/wwwroot/
```
### 2. 建表
```bash
mysql -u user -p sage < ddl/mysql.sql
```
### 3. 注册权限
`load_path.py``"""` 结束符前追加 `install/load_path_append.txt` 的内容,然后运行:
```bash
cd $SAGE_ROOT && python load_path.py
```
### 4. 添加菜单
`wwwroot/menu.ui``items` 数组中追加 `install/menu_append.json` 的内容。
### 5. 重启 Sage
访问 `/hotspot` 即可看到仪表盘。
## 使用流程
1. **来源管理** 标签页 → 添加来源(选择 type: api/rss/browser/crawler
2. **调度配置** 标签页 → 为来源配置抓取频率
3. 仪表盘点击「立即抓取全部来源」
4. **抓取日志** 标签页 → 检查执行结果
5. 点击「执行分析」→ 自动计算五维评分和状态分类
6. **热点条目** 标签页 → 按状态筛选浏览
7. **预警规则** 标签页 → 设置自动告警
## 来源类型说明
| type | 适用场景 | 配置要点 |
|------|---------|---------|
| `rss` | RSS/Atom feed | url 指向 feed 地址 |
| `api` | REST API | url + api_key + headers + fetch_method |
| `browser` | 网页热点列表 | url + browser_wait_selector |
| `crawler` | 递归抓取 | url + crawler_depth + crawler_rules |
## 依赖
- Sage >= 0.0.1
- aiohttp (Python)
- MySQL / MariaDB
建表、注册权限、加菜单 → 见 README 原文。

14
hotspot/__init__.py Normal file
View File

@ -0,0 +1,14 @@
"""
热点雷达 (Hotspot Radar) Sage 框架模块
用法:
from hotspot.engine import run_fetch
from hotspot.analysis import run_analysis
from hotspot.db import get_stats
"""
from .engine import run_fetch
from .analysis import run_analysis
from .db import get_stats
from .parsers import BUILTIN as builtin_parsers
__all__ = ['run_fetch', 'run_analysis', 'get_stats', 'builtin_parsers']

159
hotspot/analysis.py Normal file
View File

@ -0,0 +1,159 @@
"""
五维分析引擎 + 自动状态分类
维度: time(时效) / heat(热度) / content(内容) / propagation(传播) / audience(受众)
状态: emerging rising hot cooling expired
"""
import json, math
from datetime import datetime, timedelta
from appPublic.uniqueID import getID
from sqlor.dbpools import DBPools
def now_str():
return datetime.now().strftime('%Y-%m-%d %H:%M:%S')
DIM_NAMES = {
'time': '时效维度',
'heat': '热度指标',
'content': '内容属性',
'propagation': '传播路径',
'audience': '受众画像',
}
async def analyze_item(item):
"""分析单条热点,返回 (new_status, dimensions_dict, detail)"""
now = datetime.now()
try:
first_seen = datetime.strptime(str(item.get('first_seen', '')), '%Y-%m-%d %H:%M:%S')
except Exception:
first_seen = now
hours_alive = max(0.1, (now - first_seen).total_seconds() / 3600)
heat = float(item.get('heat_score', 0))
engagement = int(item.get('engagement_count', 0))
comments = int(item.get('comment_count', 0))
shares = int(item.get('share_count', 0))
# ---- 五维评分 (0-100) ----
dims = {}
# 时效: 越新越高
dims['time'] = 95 if hours_alive < 1 else 85 if hours_alive < 6 else \
70 if hours_alive < 24 else 50 if hours_alive < 72 else \
30 if hours_alive < 168 else 10
# 热度: 绝对值
dims['heat'] = 95 if heat > 1e4 else 85 if heat > 5e3 else \
70 if heat > 1e3 else 55 if heat > 500 else \
35 if heat > 100 else 15
# 内容: 标题长度 + 摘要丰富度 + 标签
title_len = len(item.get('title', ''))
summary_len = len(item.get('summary', ''))
dims['content'] = min(100, (20 if title_len > 15 else 10) +
(30 if summary_len > 100 else 15) +
(25 if item.get('tags') else 0) +
(25 if item.get('category') else 0))
# 传播: 互动总量
total_eng = engagement + comments * 2 + shares * 3
dims['propagation'] = 95 if total_eng > 1e4 else 80 if total_eng > 5e3 else \
60 if total_eng > 1e3 else 35 if total_eng > 100 else 10
# 受众: 互动率
rate = total_eng / heat if heat > 0 else 0
dims['audience'] = 90 if rate > 0.5 else 70 if rate > 0.2 else \
45 if rate > 0.05 else 25 if total_eng > 0 else 5
# ---- 状态分类 ----
decay = math.exp(-hours_alive / 168) # 7天半衰期
adjusted_heat = heat * decay
velocity = round(heat / hours_alive, 2)
if hours_alive > 336: # >14天
new_status = 'expired'
elif hours_alive > 168: # 7-14天
new_status = 'cooling'
elif adjusted_heat > 5000:
new_status = 'hot'
elif adjusted_heat > 500:
new_status = 'rising' if velocity > 50 else 'emerging'
elif adjusted_heat > 100:
new_status = 'rising' if velocity > 100 else 'emerging'
else:
new_status = 'emerging'
detail = {
'hours_alive': round(hours_alive, 1),
'adjusted_heat': round(adjusted_heat, 2),
'heat_velocity': velocity,
'total_engagement': total_eng,
'decay_factor': round(decay, 4),
}
return new_status, dims, detail
async def save_analysis(item_id, dims, detail, now=None):
"""保存/更新五维分析记录"""
if now is None:
now = now_str()
db = DBPools()
async with db.sqlorContext('sage') as sor:
for dim, score in dims.items():
analysis_data = json.dumps({
'dimension': dim,
'dimension_cn': DIM_NAMES.get(dim, dim),
'score': score,
'detail': detail,
}, ensure_ascii=False)
old = await sor.sqlExe(
"select id from hotspot_analysis where item_id=${iid}$ and dimension=${dim}$",
{'iid': item_id, 'dim': dim})
if old:
await sor.U('hotspot_analysis', {
'id': old[0]['id'], 'score': score,
'analysis_data': analysis_data, 'analyzed_at': now,
})
else:
await sor.C('hotspot_analysis', {
'id': getID(), 'item_id': item_id,
'dimension': dim, 'score': score,
'analysis_data': analysis_data, 'analyzed_at': now,
})
async def run_analysis(limit=500):
"""批量分析最近的热点 (默认500条)"""
db = DBPools()
now = now_str()
results = {'analyzed': 0, 'status_changes': 0}
async with db.sqlorContext('sage') as sor:
items = await sor.sqlExe(
"select * from hotspot_item order by first_seen desc limit ${n}$",
{'n': limit})
for item in items:
current_status = item.get('status', 'emerging')
new_status, dims, detail = await analyze_item(item)
await save_analysis(item['id'], dims, detail, now)
if new_status != current_status:
results['status_changes'] += 1
async with db.sqlorContext('sage') as sor:
await sor.U('hotspot_item', {
'id': item['id'],
'heat_score': round(detail['adjusted_heat'], 2),
'heat_velocity': detail['heat_velocity'],
'status': new_status,
'last_updated': now,
})
results['analyzed'] += 1
return results

107
hotspot/db.py Normal file
View File

@ -0,0 +1,107 @@
"""
数据库操作: 保存热点条目写抓取日志统计查询
"""
from datetime import datetime
from appPublic.uniqueID import getID
from sqlor.dbpools import DBPools
def now_str():
return datetime.now().strftime('%Y-%m-%d %H:%M:%S')
async def get_enabled_sources():
"""获取所有启用的来源"""
db = DBPools()
async with db.sqlorContext('sage') as sor:
sources = await sor.R('hotspot_source', {})
return [s for s in sources if s.get('enabled') == '1']
async def save_items(items, source_id, fetch_time=None):
"""批量保存热点条目 (去重)"""
if fetch_time is None:
fetch_time = now_str()
new_count = 0
dup_count = 0
db = DBPools()
async with db.sqlorContext('sage') as sor:
for item in items:
url = item.get('url', '')
if url:
existing = await sor.sqlExe(
"select id from hotspot_item where url=${u}$ limit 1", {'u': url})
if existing:
dup_count += 1
continue
await sor.C('hotspot_item', {
'id': getID(), 'source_id': source_id,
'title': item.get('title', '')[:500],
'url': (url or '')[:2000],
'summary': item.get('summary', '')[:1000],
'heat_score': item.get('heat_score', 100),
'heat_velocity': 0,
'publish_time': fetch_time,
'first_seen': fetch_time,
'last_updated': fetch_time,
'status': 'emerging',
'category': item.get('category', ''),
'tags': item.get('tags', ''),
'engagement_count': item.get('engagement_count', 0),
'comment_count': item.get('comment_count', 0),
'share_count': item.get('share_count', 0),
'sentiment': 'neutral',
})
new_count += 1
return new_count, dup_count
async def update_source_last_fetch(source_id):
db = DBPools()
async with db.sqlorContext('sage') as sor:
await sor.U('hotspot_source', {'id': source_id, 'last_fetch': now_str()})
async def write_fetch_log(source_id, start_time, elapsed_ms,
status, items_total, items_new, items_duplicate,
error_msg='', response_code=0, response_size_bytes=0):
db = DBPools()
async with db.sqlorContext('sage') as sor:
await sor.C('hotspot_fetch_log', {
'id': getID(),
'source_id': source_id,
'start_time': start_time,
'end_time': now_str(),
'duration_ms': elapsed_ms,
'status': status,
'items_total': items_total,
'items_new': items_new,
'items_duplicate': items_duplicate,
'error_msg': error_msg[:1000] if error_msg else '',
'response_code': response_code,
'response_size_bytes': response_size_bytes,
})
async def get_stats():
"""全维度统计"""
from datetime import timedelta
now = datetime.now()
day_ago = (now - timedelta(hours=24)).strftime('%Y-%m-%d %H:%M:%S')
db = DBPools()
async with db.sqlorContext('sage') as sor:
items = await sor.R('hotspot_item', {'order': 'heat_score desc'})
sources = await sor.R('hotspot_source', {})
logs = await sor.sqlExe(
"select * from hotspot_fetch_log where start_time >= ${t}$",
{'t': day_ago})
status_count = {'emerging': 0, 'rising': 0, 'hot': 0, 'cooling': 0, 'expired': 0}
for item in items:
s = item.get('status', 'emerging')
status_count[s] = status_count.get(s, 0) + 1
return {
'total': len(items),
**status_count,
'sources': sum(1 for s in sources if s.get('enabled') == '1'),
'failures': sum(1 for log in logs if log.get('status') == 'failed'),
}

141
hotspot/engine.py Normal file
View File

@ -0,0 +1,141 @@
"""
抓取引擎 调度所有来源fetch + extract + save
"""
import json, time
import aiohttp
from datetime import datetime
from .parsers import (BUILTIN, extract_json_path, extract_css,
extract_regex, extract_rss, fetch_hn)
from .db import (save_items, update_source_last_fetch,
write_fetch_log, get_enabled_sources)
def now_str():
return datetime.now().strftime('%Y-%m-%d %H:%M:%S')
async def fetch_and_extract(session, src):
"""根据 source 的 parser_config 抓取并提取条目"""
url = src.get('url', '')
headers = {}
if src.get('headers'):
try:
headers = json.loads(src.get('headers', '{}'))
except Exception:
pass
parser_config = {}
if src.get('parser_config'):
try:
parser_config = json.loads(src.get('parser_config', '{}'))
except Exception:
pass
parser_type = parser_config.get('parser', 'json_path')
builtin_name = parser_config.get('builtin_name', '')
# === 内置解析器 ===
if parser_type == 'builtin' and builtin_name in BUILTIN:
bp = BUILTIN[builtin_name]
fetch_url = bp['url'] or url
fetch_headers = {**headers, **bp.get('headers', {})}
if builtin_name == 'hackernews':
items = await fetch_hn(session)
return items, 200, 0
if builtin_name == 'rsshub':
try:
async with session.get(url, headers=fetch_headers,
timeout=aiohttp.ClientTimeout(total=30)) as resp:
text = await resp.text()
items = extract_rss(text)
return items, resp.status, len(text.encode())
except Exception:
return [], 0, 0
# 标准内置: GET → JSON → parser_fn
try:
async with session.get(fetch_url, headers=fetch_headers,
timeout=aiohttp.ClientTimeout(total=30)) as resp:
if resp.status != 200:
return [], resp.status, 0
text = await resp.text()
size = len(text.encode())
try:
data = json.loads(text)
except json.JSONDecodeError:
return [], resp.status, size
items = bp['parser'](data)
return items, resp.status, size
except Exception:
return [], 0, 0
# === 通用提取器 ===
try:
async with session.get(url, headers=headers,
timeout=aiohttp.ClientTimeout(total=30)) as resp:
if resp.status != 200:
return [], resp.status, 0
text = await resp.text()
size = len(text.encode())
if parser_type == 'json_path':
try:
data = json.loads(text)
except json.JSONDecodeError:
return [], 0, size
items = extract_json_path(data, parser_config)
return items, resp.status, size
elif parser_type == 'css':
items = extract_css(text, parser_config)
return items, resp.status, size
elif parser_type == 'regex':
items = extract_regex(text, parser_config)
return items, resp.status, size
elif parser_type == 'rss':
items = extract_rss(text)
return items, resp.status, size
return [], resp.status, size
except Exception:
return [], 0, 0
async def run_fetch():
"""遍历所有启用的来源,执行抓取"""
sources = await get_enabled_sources()
now = now_str()
results = {'total': len(sources), 'new_items': 0, 'errors': 0, 'details': []}
async with aiohttp.ClientSession() as session:
for src in sources:
t0 = time.time()
try:
items, code, size = await fetch_and_extract(session, src)
except Exception as e:
await write_fetch_log(src['id'], now, int((time.time() - t0) * 1000),
'failed', 0, 0, 0, str(e)[:1000])
results['errors'] += 1
continue
elapsed = int((time.time() - t0) * 1000)
status = 'success' if code in (0, 200) else ('partial' if items else 'failed')
new_count, dup_count = await save_items(items, src['id'], now)
await update_source_last_fetch(src['id'])
await write_fetch_log(src['id'], now, elapsed, status,
len(items), new_count, dup_count,
response_code=code, response_size_bytes=size)
results['new_items'] += new_count
results['details'].append({
'source': src.get('name', src['id']),
'items': len(items), 'new': new_count, 'dup': dup_count,
'elapsed_ms': elapsed, 'code': code,
})
return results

297
hotspot/parsers.py Normal file
View File

@ -0,0 +1,297 @@
"""
内置平台解析器 每个返回 (fetch_url, headers_override, extract_fn)
添加新平台: 在模块底部调用 _reg() 注册即可
"""
import json
import re
import xml.etree.ElementTree as ET
import aiohttp
BUILTIN = {}
def _reg(name, url, parser_fn, headers=None):
BUILTIN[name] = {'url': url, 'parser': parser_fn, 'headers': headers or {}}
# ---- 国内平台 ----
_reg('weibo_hot',
'https://weibo.com/ajax/side/hotSearch',
lambda data: [
{'title': item.get('word', ''),
'url': f'https://s.weibo.com/weibo?q={item.get("word","")}',
'heat_score': item.get('raw_hot', item.get('num', 0))}
for item in data.get('data', {}).get('realtime', [])[:50]
if isinstance(item, dict) and item.get('word')
],
{'User-Agent': 'Mozilla/5.0', 'X-Requested-With': 'XMLHttpRequest'})
_reg('zhihu_hot',
'https://www.zhihu.com/api/v3/feed/topstory/hot-lists/total?limit=50',
lambda data: [
{'title': item.get('target', {}).get('title', ''),
'url': f'https://www.zhihu.com/question/{item.get("target",{}).get("id","")}',
'heat_score': int(str(item.get('detail_text', '0'))
.replace('','0000').replace('亿','00000000')
.replace('热度','').strip() or 0),
'summary': (item.get('target', {}).get('excerpt', '') or '')[:500]}
for item in data.get('data', []) if isinstance(item, dict)
],
{'User-Agent': 'Mozilla/5.0'})
_reg('baidu_hot',
'https://top.baidu.com/board?tab=realtime',
lambda data: [
{'title': item.get('word', ''),
'url': item.get('url', ''),
'heat_score': int(item.get('hotScore', 0)),
'summary': item.get('desc', '')}
for item in (data.get('data', {}).get('cards', [{}])[0].get('content', []))
if isinstance(item, dict)
],
{'User-Agent': 'Mozilla/5.0'})
_reg('toutiao_hot',
'https://www.toutiao.com/hot-event/hot-board/?origin=toutiao_pc',
lambda data: [
{'title': item.get('Title', ''),
'url': item.get('Url', ''),
'heat_score': int(float(item.get('HotValue', 0)))}
for item in (data.get('data', []) if isinstance(data.get('data'), list) else [])[:50]
if isinstance(item, dict) and item.get('Title')
],
{'User-Agent': 'Mozilla/5.0', 'Referer': 'https://www.toutiao.com/'})
_reg('bilibili_hot',
'https://api.bilibili.com/x/web-interface/popular?ps=50',
lambda data: [
{'title': item.get('title', ''),
'url': f'https://www.bilibili.com/video/{item.get("bvid","")}',
'heat_score': item.get('stat', {}).get('view', 0),
'summary': item.get('desc', '')[:500],
'engagement_count': item.get('stat', {}).get('like', 0),
'comment_count': item.get('stat', {}).get('reply', 0),
'share_count': item.get('stat', {}).get('share', 0),
'category': item.get('tname', '')}
for item in data.get('data', {}).get('list', [])
if isinstance(item, dict) and item.get('title')
],
{'User-Agent': 'Mozilla/5.0', 'Referer': 'https://www.bilibili.com/'})
_reg('douyin_hot',
'https://dy.feigua.cn/api/hot/board',
lambda data: [
{'title': item.get('title', item.get('word', '')),
'url': item.get('url', ''),
'heat_score': int(item.get('hot_value', item.get('heat', 0)))}
for item in (data.get('data', []) if isinstance(data.get('data'), list) else [])[:50]
if isinstance(item, dict)
],
{'User-Agent': 'Mozilla/5.0'})
_reg('36kr_hot',
'https://www.36kr.com/api/search/list/info-flow/hot?pageSize=30',
lambda data: [
{'title': item.get('title', ''),
'url': f'https://www.36kr.com/p/{item.get("id","")}',
'heat_score': item.get('total_count', 0),
'summary': (item.get('summary', '') or '')[:500]}
for item in data.get('data', {}).get('items', [])
if isinstance(item, dict) and item.get('title')
],
{'User-Agent': 'Mozilla/5.0'})
# ---- 海外平台 ----
_reg('github_trending',
'https://api.github.com/search/repositories?q=stars:>1&sort=stars&order=desc&per_page=20',
lambda data: [
{'title': item.get('full_name', ''),
'url': item.get('html_url', ''),
'heat_score': item.get('stargazers_count', 0),
'summary': (item.get('description', '') or '')[:500],
'category': item.get('language', '')}
for item in data.get('items', []) if isinstance(item, dict)
],
{'Accept': 'application/vnd.github.v3+json', 'User-Agent': 'HotspotRadar/1.0'})
_reg('hackernews',
'https://hacker-news.firebaseio.com/v0/topstories.json',
None, # 特殊: 两步抓取, 见 engine.py
{'User-Agent': 'Mozilla/5.0'})
_reg('v2ex_hot',
'https://www.v2ex.com/api/topics/hot.json',
lambda data: [
{'title': item.get('title', ''),
'url': item.get('url', ''),
'heat_score': item.get('replies', 0) * 10,
'summary': (item.get('content', '') or '')[:500],
'comment_count': item.get('replies', 0),
'category': item.get('node', {}).get('title', '')
if isinstance(item.get('node'), dict) else ''}
for item in data if isinstance(item, dict) and item.get('title')
],
{'User-Agent': 'Mozilla/5.0'})
_reg('rsshub', None, None, {}) # url 来自 source 配置
# ---- 通用提取器 ----
def extract_json_path(data, config):
"""简易 JSONPath: $.data.list[*] → {title, url, heat_score, ...}"""
path = config.get('item_path', '$[*]')
fm = config.get('field_map', {})
if path.startswith('$.'):
path = path[2:]
elif path.startswith('$'):
path = path[1:]
parts = []
for p in path.split('.'):
p = p.strip()
if not p: continue
if p.endswith('[*]'):
parts.append(p[:-3]); parts.append('*')
else:
parts.append(p)
items = data
for part in parts:
if part == '*':
continue
if isinstance(items, dict):
items = items.get(part)
elif isinstance(items, list):
items = [i.get(part) for i in items if isinstance(i, dict)]
items = [x for sub in items for x in (sub if isinstance(sub, list) else [sub])]
else:
return []
if not isinstance(items, list):
items = [items] if items else []
result = []
for item in items:
if not isinstance(item, dict):
continue
entry = {
'title': _field(item, fm.get('title', 'title')),
'url': _field(item, fm.get('url', 'url')),
'heat_score': float(_field(item, fm.get('heat_score', '0')) or 0),
'summary': str(_field(item, fm.get('summary', '')) or '')[:500],
'category': str(_field(item, fm.get('category', '')) or ''),
'tags': str(_field(item, fm.get('tags', '')) or ''),
'engagement_count': int(_field(item, fm.get('engagement_count', '0')) or 0),
'comment_count': int(_field(item, fm.get('comment_count', '0')) or 0),
'share_count': int(_field(item, fm.get('share_count', '0')) or 0),
}
if entry['title']:
result.append(entry)
return result
def _field(obj, path):
if not path or not isinstance(obj, dict):
return obj if not isinstance(obj, dict) else ''
val = obj
for p in path.split('.'):
if isinstance(val, dict):
val = val.get(p, '')
else:
return ''
return val
def extract_css(html, config):
"""CSS 选择器提取 (轻量正则版)"""
results = []
for m in re.finditer(
r'<a[^>]+href=["\']([^"\']+)["\'][^>]*>([^<]*)</a>',
html, re.IGNORECASE | re.DOTALL
):
url, title = m.group(1), m.group(2).strip()
if len(title) >= 3:
results.append({
'title': title[:200], 'url': url[:2000],
'heat_score': 60, 'summary': '',
'category': '', 'tags': '',
'engagement_count': 0, 'comment_count': 0, 'share_count': 0,
})
return results[:100]
def extract_regex(text, config):
pattern = config.get('item_regex', '')
fm = config.get('field_map', {})
if not pattern: return []
results = []
for m in re.finditer(pattern, text, re.DOTALL):
gd = m.groupdict()
entry = {
'title': gd.get(fm.get('title', 'title'), '')[:200],
'url': gd.get(fm.get('url', 'url'), ''),
'heat_score': float(gd.get(fm.get('heat_score', '0'), 0) or 0),
'summary': str(gd.get(fm.get('summary', ''), ''))[:500],
'category': '', 'tags': '',
'engagement_count': 0, 'comment_count': 0, 'share_count': 0,
}
if entry['title']:
results.append(entry)
return results
def extract_rss(xml_text):
items = []
try:
root = ET.fromstring(xml_text)
ns_a = 'http://www.w3.org/2005/Atom'
for item in root.findall('.//item'):
items.append({
'title': item.findtext('title', '')[:200],
'url': item.findtext('link', ''),
'summary': (item.findtext('description', '') or '')[:500],
'heat_score': 80, 'category': item.findtext('category', ''),
'tags': '', 'engagement_count': 0, 'comment_count': 0, 'share_count': 0,
})
if not items:
for e in root.findall(f'.//{{{ns_a}}}entry'):
link = e.find(f'{{{ns_a}}}link')
items.append({
'title': e.findtext(f'{{{ns_a}}}title', '')[:200],
'url': link.get('href', '') if link is not None else '',
'summary': (e.findtext(f'{{{ns_a}}}summary', '') or '')[:500],
'heat_score': 70, 'category': '',
'tags': '', 'engagement_count': 0, 'comment_count': 0, 'share_count': 0,
})
except ET.ParseError:
pass
return items
async def fetch_hn(session):
"""HackerNews 两步抓取"""
items = []
try:
async with session.get(
'https://hacker-news.firebaseio.com/v0/topstories.json',
timeout=aiohttp.ClientTimeout(total=15)
) as resp:
ids = await resp.json()
for sid in ids[:30]:
try:
async with session.get(
f'https://hacker-news.firebaseio.com/v0/item/{sid}.json',
timeout=aiohttp.ClientTimeout(total=10)
) as r2:
s = await r2.json()
if s and s.get('title'):
items.append({
'title': s.get('title', '')[:200],
'url': s.get('url', f'https://news.ycombinator.com/item?id={sid}'),
'heat_score': s.get('score', 0),
'engagement_count': s.get('score', 0),
'comment_count': s.get('descendants', 0),
'summary': '', 'category': s.get('type', ''),
'tags': '', 'share_count': 0,
})
except Exception:
continue
except Exception:
pass
return items

View File

@ -1,205 +1,2 @@
# 五维 + 状态自动分析
# emerging → rising → hot → cooling → expired
import json
from datetime import datetime, timedelta
from sqlor.dbpools import DBPools
from appPublic.uniqueID import getID
def now_str():
return datetime.now().strftime('%Y-%m-%d %H:%M:%S')
async def main(request):
db = DBPools()
now = datetime.now()
results = {'analyzed': 0, 'status_changes': 0, 'dimensions': {}}
async with db.sqlorContext('sage') as sor:
items = await sor.sqlExe(
"select * from hotspot_item order by first_seen desc limit 500",
{}
)
for item in items:
try:
first_seen = datetime.strptime(str(item.get('first_seen', '')), '%Y-%m-%d %H:%M:%S')
except:
first_seen = now
hours_alive = max(0.1, (now - first_seen).total_seconds() / 3600)
heat = float(item.get('heat_score', 0))
velocity = float(item.get('heat_velocity', 0))
engagement = int(item.get('engagement_count', 0))
comments = int(item.get('comment_count', 0))
shares = int(item.get('share_count', 0))
current_status = item.get('status', 'emerging')
item_id = item['id']
# === 五维分析 ===
dimensions = {}
# 1. 时效维度 (0-100) — 越新越高
if hours_alive < 1:
dim_time = 95
elif hours_alive < 6:
dim_time = 85
elif hours_alive < 24:
dim_time = 70
elif hours_alive < 72:
dim_time = 50
elif hours_alive < 168:
dim_time = 30
else:
dim_time = 10
dimensions['time'] = dim_time
# 2. 热度维度 (0-100)
if heat > 10000:
dim_heat = 95
elif heat > 5000:
dim_heat = 85
elif heat > 1000:
dim_heat = 70
elif heat > 500:
dim_heat = 55
elif heat > 100:
dim_heat = 35
else:
dim_heat = 15
dimensions['heat'] = dim_heat
# 3. 内容维度 (0-100) — 基于标题长度+摘要丰富度
title_len = len(item.get('title', ''))
summary_len = len(item.get('summary', ''))
has_tags = bool(item.get('tags'))
has_category = bool(item.get('category'))
dim_content = min(100,
(20 if title_len > 15 else 10) +
(30 if summary_len > 100 else 15) +
(25 if has_tags else 0) +
(25 if has_category else 0)
)
dimensions['content'] = dim_content
# 4. 传播维度 (0-100) — 互动量
total_engagement = engagement + comments * 2 + shares * 3
if total_engagement > 10000:
dim_propagation = 95
elif total_engagement > 5000:
dim_propagation = 80
elif total_engagement > 1000:
dim_propagation = 60
elif total_engagement > 100:
dim_propagation = 35
else:
dim_propagation = 10
dimensions['propagation'] = dim_propagation
# 5. 受众维度 (0-100) — 基于互动率
if heat > 0:
engagement_rate = total_engagement / heat
else:
engagement_rate = 0
if engagement_rate > 0.5:
dim_audience = 90
elif engagement_rate > 0.2:
dim_audience = 70
elif engagement_rate > 0.05:
dim_audience = 45
elif total_engagement > 0:
dim_audience = 25
else:
dim_audience = 5
dimensions['audience'] = dim_audience
# === 状态分类 (基于热度 + 时间) ===
# 热度衰减: heat * e^(-hours/168) ~ 7天半衰期
import math
decay = math.exp(-hours_alive / 168)
adjusted_heat = heat * decay
# 热度加速度 (简化:基于当前热度/time)
new_velocity = round(heat / max(hours_alive, 0.1), 2)
if hours_alive > 336: # 超过14天
new_status = 'expired'
elif hours_alive > 168: # 7-14天
new_status = 'cooling'
elif adjusted_heat > 5000:
new_status = 'hot'
elif adjusted_heat > 500:
new_status = 'rising' if new_velocity > 50 else 'emerging'
elif adjusted_heat > 100:
new_status = 'rising' if new_velocity > 100 else 'emerging'
else:
new_status = 'emerging'
# 保存
async with db.sqlorContext('sage') as sor:
# 更新条目
await sor.U('hotspot_item', {
'id': item_id,
'heat_score': round(adjusted_heat, 2),
'heat_velocity': new_velocity,
'status': new_status,
'last_updated': now_str(),
})
results['analyzed'] += 1
if new_status != current_status:
results['status_changes'] += 1
# 保存五维分析
for dim, score in dimensions.items():
dim_names = {
'time': '时效维度',
'heat': '热度指标',
'content': '内容属性',
'propagation': '传播路径',
'audience': '受众画像',
}
analysis_data = json.dumps({
'dimension': dim,
'dimension_cn': dim_names.get(dim, dim),
'score': score,
'detail': {
'hours_alive': round(hours_alive, 1),
'adjusted_heat': round(adjusted_heat, 2),
'heat_velocity': new_velocity,
'total_engagement': total_engagement,
'decay_factor': round(decay, 4),
}
}, ensure_ascii=False)
# Upsert: delete old analysis for this item+dimension, insert new
old = await sor.sqlExe(
"select id from hotspot_analysis where item_id=${iid}$ and dimension=${dim}$",
{'iid': item_id, 'dim': dim}
)
if old:
await sor.U('hotspot_analysis', {
'id': old[0]['id'],
'score': score,
'analysis_data': analysis_data,
'analyzed_at': now_str(),
})
else:
await sor.C('hotspot_analysis', {
'id': getID(),
'item_id': item_id,
'dimension': dim,
'score': score,
'analysis_data': analysis_data,
'analyzed_at': now_str(),
})
results['dimensions'] = {
'time': dim_names['time'],
'heat': dim_names['heat'],
'content': dim_names['content'],
'propagation': dim_names['propagation'],
'audience': dim_names['audience'],
}
return results
from hotspot.analysis import run_analysis
return await run_analysis()

View File

@ -1,570 +1,2 @@
"""
抓取引擎 — 可配置提取规则 + 内置平台解析器
每条 hotspot_source 可带 parser_config JSON:
方式 A — 内置解析器: {"parser": "builtin", "builtin_name": "weibo"}
方式 B — JSONPath: {"parser": "json_path", "item_path": "$.data.list[*]", "field_map": {...}}
方式 C — CSS选择器: {"parser": "css", "item_selector": ".hot-item", "field_map": {...}}
方式 D — Regex: {"parser": "regex", "item_regex": "...", "field_map": {...}}
方式 E — RSS: {"parser": "rss"}
field_map: {"title":"...", "url":"...", "heat_score":"...", "summary":"...", "category":"...",
"tags":"...", "engagement_count":"...", "comment_count":"...", "share_count":"..."}
"""
import json, re, time
import xml.etree.ElementTree as ET
from datetime import datetime
from appPublic.uniqueID import getID
from sqlor.dbpools import DBPools
def now_str():
return datetime.now().strftime('%Y-%m-%d %H:%M:%S')
# ============================================================
# 内置平台解析器 — 每个返回 (url, headers, extractor函数)
# ============================================================
BUILTIN_PARSERS = {}
def _reg(name, url, parser_fn, headers=None):
BUILTIN_PARSERS[name] = {'url': url, 'parser': parser_fn, 'headers': headers or {}}
# -- 微博热搜 --
_reg('weibo_hot',
'https://weibo.com/ajax/side/hotSearch',
lambda data: [
{'title': item.get('word', ''), 'url': f'https://s.weibo.com/weibo?q={item.get("word","")}',
'heat_score': item.get('raw_hot', item.get('num', 0)),
'summary': item.get('word_scheme', item.get('word', ''))}
for item in data.get('data', {}).get('realtime', [])[:50] if isinstance(item, dict)
],
{'User-Agent': 'Mozilla/5.0', 'X-Requested-With': 'XMLHttpRequest'})
# -- 知乎热榜 --
_reg('zhihu_hot',
'https://www.zhihu.com/api/v3/feed/topstory/hot-lists/total?limit=50',
lambda data: [
{'title': item.get('target', {}).get('title', ''),
'url': f'https://www.zhihu.com/question/{item.get("target",{}).get("id","")}',
'heat_score': int(item.get('detail_text', '0').replace('万', '0000').replace('亿', '00000000').replace('热度', '') or 0),
'summary': item.get('target', {}).get('excerpt', '')[:500]}
for item in data.get('data', []) if isinstance(item, dict)
],
{'User-Agent': 'Mozilla/5.0'})
# -- 百度热搜 --
_reg('baidu_hot',
'https://top.baidu.com/board?tab=realtime',
lambda data: [
{'title': item.get('word', ''), 'url': item.get('url', ''),
'heat_score': int(item.get('hotScore', 0)),
'summary': item.get('desc', '')}
for item in data.get('data', {}).get('cards', [{}])[0].get('content', []) if isinstance(item, dict)
],
{'User-Agent': 'Mozilla/5.0'})
# -- 头条热榜 --
_reg('toutiao_hot',
'https://www.toutiao.com/hot-event/hot-board/?origin=toutiao_pc',
lambda data: [
{'title': item.get('Title', ''), 'url': item.get('Url', ''),
'heat_score': int(float(item.get('HotValue', 0))),
'summary': item.get('Desc', '')}
for item in (data.get('data', []) if isinstance(data.get('data'), list) else [])[:50]
],
{'User-Agent': 'Mozilla/5.0', 'Referer': 'https://www.toutiao.com/'})
# -- B站热门 --
_reg('bilibili_hot',
'https://api.bilibili.com/x/web-interface/popular?ps=50',
lambda data: [
{'title': item.get('title', ''), 'url': f'https://www.bilibili.com/video/{item.get("bvid","")}',
'heat_score': item.get('stat', {}).get('view', 0),
'summary': item.get('desc', '')[:500],
'engagement_count': item.get('stat', {}).get('like', 0),
'comment_count': item.get('stat', {}).get('reply', 0),
'share_count': item.get('stat', {}).get('share', 0),
'category': item.get('tname', '')}
for item in data.get('data', {}).get('list', []) if isinstance(item, dict)
],
{'User-Agent': 'Mozilla/5.0', 'Referer': 'https://www.bilibili.com/'})
# -- 抖音热点 (通过第三方聚合接口) --
_reg('douyin_hot',
'https://dy.feigua.cn/api/hot/board',
lambda data: [
{'title': item.get('title', item.get('word', '')),
'url': item.get('url', ''),
'heat_score': int(item.get('hot_value', item.get('heat', 0))),
'summary': item.get('desc', '')[:500]}
for item in (data.get('data', []) if isinstance(data.get('data'), list) else [])[:50]
],
{'User-Agent': 'Mozilla/5.0'})
# -- 36氪热榜 --
_reg('36kr_hot',
'https://www.36kr.com/api/search/list/info-flow/hot?pageSize=30',
lambda data: [
{'title': item.get('title', ''), 'url': f'https://www.36kr.com/p/{item.get("id","")}',
'heat_score': item.get('total_count', 0),
'summary': item.get('summary', '')[:500],
'engagement_count': item.get('total_count', 0)}
for item in data.get('data', {}).get('items', []) if isinstance(item, dict)
],
{'User-Agent': 'Mozilla/5.0'})
# -- GitHub Trending (HTML parsing via API proxy) --
_reg('github_trending',
'https://api.github.com/search/repositories?q=stars:>1&sort=stars&order=desc&per_page=20',
lambda data: [
{'title': item.get('full_name', ''), 'url': item.get('html_url', ''),
'heat_score': item.get('stargazers_count', 0),
'summary': item.get('description', '')[:500],
'category': item.get('language', ''),
'engagement_count': item.get('forks_count', 0)}
for item in data.get('items', []) if isinstance(item, dict)
],
{'Accept': 'application/vnd.github.v3+json', 'User-Agent': 'HotspotRadar/1.0'})
# -- HackerNews --
_reg('hackernews',
'https://hacker-news.firebaseio.com/v0/topstories.json',
None, # special: needs 2-step fetch
{'User-Agent': 'Mozilla/5.0'})
# -- V2EX 热帖 --
_reg('v2ex_hot',
'https://www.v2ex.com/api/topics/hot.json',
lambda data: [
{'title': item.get('title', ''), 'url': item.get('url', ''),
'heat_score': item.get('replies', 0) * 10,
'summary': (item.get('content', '') or '')[:500],
'comment_count': item.get('replies', 0),
'category': item.get('node', {}).get('title', '') if isinstance(item.get('node'), dict) else ''}
for item in data if isinstance(item, dict)
],
{'User-Agent': 'Mozilla/5.0'})
# -- RSSHub (标准格式) --
_reg('rsshub',
None, # url comes from source config
None, # parsed as RSS/Atom feed
{})
# ============================================================
# 可配置提取器
# ============================================================
def extract_json_path(data, config):
"""JSONPath 提取 (简易实现: $.a.b[*] 格式)"""
item_path = config.get('item_path', '$[*]')
field_map = config.get('field_map', {})
# 解析路径: $.data.list[*] → ['data', 'list', '*']
path_parts = []
if item_path.startswith('$.'):
item_path = item_path[2:]
elif item_path.startswith('$'):
item_path = item_path[1:]
# 按 . 分割,处理 [*] 通配
for part in item_path.split('.'):
part = part.strip()
if not part:
continue
if part.endswith('[*]'):
path_parts.append(part[:-3])
path_parts.append('*')
else:
path_parts.append(part)
# 导航到 items 数组
items = data
for part in path_parts:
if part == '*':
if isinstance(items, list):
pass # wildcard means iterate
continue
if isinstance(items, dict):
items = items.get(part)
elif isinstance(items, list):
# 对每个元素取 part
items = [i.get(part) for i in items if isinstance(i, dict)]
items = [x for sub in items for x in (sub if isinstance(sub, list) else [sub])]
else:
return []
if not isinstance(items, list):
items = [items] if items else []
result = []
for item in items:
if not isinstance(item, dict):
continue
entry = {
'title': _extract_field(item, field_map.get('title', 'title')),
'url': _extract_field(item, field_map.get('url', 'url')),
'heat_score': float(_extract_field(item, field_map.get('heat_score', '0')) or 0),
'summary': str(_extract_field(item, field_map.get('summary', '')) or '')[:500],
'category': str(_extract_field(item, field_map.get('category', '')) or ''),
'tags': str(_extract_field(item, field_map.get('tags', '')) or ''),
'engagement_count': int(_extract_field(item, field_map.get('engagement_count', '0')) or 0),
'comment_count': int(_extract_field(item, field_map.get('comment_count', '0')) or 0),
'share_count': int(_extract_field(item, field_map.get('share_count', '0')) or 0),
}
if entry['title']:
result.append(entry)
return result
def _extract_field(obj, path):
"""从嵌套字典中按路径取值: 'target.title' 或 'stat.view'"""
if not path or not isinstance(obj, dict):
return obj if not isinstance(obj, dict) else ''
parts = path.split('.')
val = obj
for p in parts:
if isinstance(val, dict):
val = val.get(p, '')
else:
return ''
return val
def extract_css(html, config):
"""
CSS 选择器提取 (简易实现: 正则 + 文本匹配)
完整版需要 BeautifulSoup/pyquery — 这里做轻量正则提取
"""
item_sel = config.get('item_selector', '')
field_map = config.get('field_map', {})
results = []
# 尝试按常见的 HTML 模式分割
# 匹配 <a> 标签中的 title 和 href
link_pattern = re.compile(
r'<a[^>]+href=["\']([^"\']+)["\'][^>]*>([^<]*)</a>',
re.IGNORECASE | re.DOTALL
)
for m in link_pattern.finditer(html):
url, title = m.group(1), m.group(2).strip()
if len(title) < 3:
continue
results.append({
'title': title[:200],
'url': url[:2000],
'heat_score': 60,
'summary': '',
'category': '',
'tags': '',
'engagement_count': 0,
'comment_count': 0,
'share_count': 0,
})
return results[:100]
def extract_regex(text, config):
"""正则提取"""
pattern = config.get('item_regex', '')
field_map = config.get('field_map', {})
if not pattern:
return []
results = []
for m in re.finditer(pattern, text, re.DOTALL):
entry = {
'title': m.groupdict().get(field_map.get('title', 'title'), '')[:200],
'url': m.groupdict().get(field_map.get('url', 'url'), ''),
'heat_score': float(m.groupdict().get(field_map.get('heat_score', '0'), 0) or 0),
'summary': str(m.groupdict().get(field_map.get('summary', ''), ''))[:500],
'category': '',
'tags': '',
'engagement_count': 0,
'comment_count': 0,
'share_count': 0,
}
if entry['title']:
results.append(entry)
return results
def extract_rss(xml_text):
"""RSS/Atom 提取"""
items = []
try:
root = ET.fromstring(xml_text)
ns_atom = 'http://www.w3.org/2005/Atom'
for item in root.findall('.//item'):
items.append({
'title': item.findtext('title', '')[:200],
'url': item.findtext('link', ''),
'summary': (item.findtext('description', '') or '')[:500],
'heat_score': 80,
'category': item.findtext('category', ''),
'tags': '',
'engagement_count': 0, 'comment_count': 0, 'share_count': 0,
})
if not items:
for entry in root.findall(f'.//{{{ns_atom}}}entry'):
link_el = entry.find(f'{{{ns_atom}}}link')
items.append({
'title': entry.findtext(f'{{{ns_atom}}}title', '')[:200],
'url': link_el.get('href', '') if link_el is not None else '',
'summary': (entry.findtext(f'{{{ns_atom}}}summary', '') or '')[:500],
'heat_score': 70,
'category': '',
'tags': '',
'engagement_count': 0, 'comment_count': 0, 'share_count': 0,
})
except ET.ParseError:
pass
return items
# ============================================================
# 主逻辑: fetch + extract
# ============================================================
async def fetch_and_extract(session, src):
"""根据 source 配置,抓取并提取 items"""
url = src.get('url', '')
headers = {}
if src.get('headers'):
try:
headers = json.loads(src.get('headers', '{}'))
except:
pass
# 解析 parser_config
parser_config = {}
if src.get('parser_config'):
try:
parser_config = json.loads(src.get('parser_config', '{}'))
except:
pass
parser_type = parser_config.get('parser', 'json_path')
builtin_name = parser_config.get('builtin_name', '')
import aiohttp
# === 内置解析器路径 ===
if parser_type == 'builtin' and builtin_name in BUILTIN_PARSERS:
bp = BUILTIN_PARSERS[builtin_name]
fetch_url = bp['url'] or url
fetch_headers = {**headers, **bp.get('headers', {})}
if builtin_name == 'hackernews':
# 两步: 先拿 top IDs, 再批量拿详情
try:
async with aiohttp.ClientSession() as hn_session:
async with hn_session.get(
'https://hacker-news.firebaseio.com/v0/topstories.json',
timeout=aiohttp.ClientTimeout(total=15)
) as resp:
ids = await resp.json()
items = []
for story_id in ids[:30]:
try:
async with hn_session.get(
f'https://hacker-news.firebaseio.com/v0/item/{story_id}.json',
timeout=aiohttp.ClientTimeout(total=10)
) as r2:
story = await r2.json()
if story and story.get('title'):
items.append({
'title': story.get('title', '')[:200],
'url': story.get('url', f'https://news.ycombinator.com/item?id={story_id}'),
'heat_score': story.get('score', 0),
'engagement_count': story.get('score', 0),
'comment_count': story.get('descendants', 0),
'summary': '',
'category': story.get('type', ''),
'tags': '',
'share_count': 0,
})
except:
continue
return items, 200, 0
except Exception as e:
return [], 0, 0
elif builtin_name == 'rsshub':
# 来源的 url 就是 RSSHub 地址
try:
async with session.get(
url, headers=fetch_headers,
timeout=aiohttp.ClientTimeout(total=30)
) as resp:
text = await resp.text()
size = len(text.encode())
items = extract_rss(text)
return items, resp.status, size
except Exception as e:
return [], 0, 0
else:
try:
async with session.get(
fetch_url, headers=fetch_headers,
timeout=aiohttp.ClientTimeout(total=30)
) as resp:
if resp.status != 200:
return [], resp.status, 0
text = await resp.text()
size = len(text.encode())
try:
data = json.loads(text)
except json.JSONDecodeError:
return [], resp.status, size
items = bp['parser'](data)
return items, resp.status, size
except Exception as e:
return [], 0, 0
# === 通用提取器路径 ===
try:
async with session.get(
url, headers=headers,
timeout=aiohttp.ClientTimeout(total=30)
) as resp:
if resp.status != 200:
return [], resp.status, 0
text = await resp.text()
size = len(text.encode())
if parser_type == 'json_path':
try:
data = json.loads(text)
except json.JSONDecodeError:
return [], 0, size
items = extract_json_path(data, parser_config)
return items, resp.status, size
elif parser_type == 'css':
items = extract_css(text, parser_config)
return items, resp.status, size
elif parser_type == 'regex':
items = extract_regex(text, parser_config)
return items, resp.status, size
elif parser_type == 'rss':
items = extract_rss(text)
return items, resp.status, size
else:
return [], resp.status, size
except Exception as e:
return [], 0, 0
# ============================================================
# 入库
# ============================================================
async def save_items(db, src_id, items, fetch_time):
new_count = 0
dup_count = 0
async with db.sqlorContext('sage') as sor:
for item in items:
url = item.get('url', '')
if url:
existing = await sor.sqlExe(
"select id from hotspot_item where url=${u}$ limit 1", {'u': url})
if existing:
dup_count += 1
continue
await sor.C('hotspot_item', {
'id': getID(),
'source_id': src_id,
'title': item.get('title', '')[:500],
'url': url[:2000] if url else '',
'summary': item.get('summary', '')[:1000],
'heat_score': item.get('heat_score', 100),
'heat_velocity': 0,
'publish_time': fetch_time,
'first_seen': fetch_time,
'last_updated': fetch_time,
'status': 'emerging',
'category': item.get('category', ''),
'tags': item.get('tags', ''),
'engagement_count': item.get('engagement_count', 0),
'comment_count': item.get('comment_count', 0),
'share_count': item.get('share_count', 0),
'sentiment': 'neutral',
})
new_count += 1
return new_count, dup_count
# ============================================================
# 入口
# ============================================================
async def main(request):
db = DBPools()
now = now_str()
results = {'total': 0, 'new_items': 0, 'errors': 0, 'details': []}
async with db.sqlorContext('sage') as sor:
sources = await sor.R('hotspot_source', {})
import aiohttp
async with aiohttp.ClientSession() as session:
for src in sources:
if src.get('enabled') != '1':
continue
results['total'] += 1
src_id = src['id']
t0 = time.time()
try:
items, code, size = await fetch_and_extract(session, src)
except Exception as e:
async with db.sqlorContext('sage') as sor:
await sor.C('hotspot_fetch_log', {
'id': getID(), 'source_id': src_id,
'start_time': now, 'end_time': now_str(),
'duration_ms': int((time.time() - t0) * 1000),
'status': 'failed', 'items_total': 0, 'items_new': 0,
'items_duplicate': 0,
'error_msg': str(e)[:1000],
'response_code': 0, 'response_size_bytes': 0,
})
results['errors'] += 1
continue
elapsed = int((time.time() - t0) * 1000)
status = 'success' if code in (0, 200) else ('partial' if items else 'failed')
new_count, dup_count = await save_items(db, src_id, items, now)
async with db.sqlorContext('sage') as sor:
await sor.U('hotspot_source', {'id': src_id, 'last_fetch': now})
await sor.C('hotspot_fetch_log', {
'id': getID(), 'source_id': src_id,
'start_time': now, 'end_time': now_str(),
'duration_ms': elapsed, 'status': status,
'items_total': len(items), 'items_new': new_count,
'items_duplicate': dup_count, 'error_msg': '',
'response_code': code, 'response_size_bytes': size,
})
results['new_items'] += new_count
results['details'].append({
'source': src.get('name', src_id),
'items': len(items), 'new': new_count, 'dup': dup_count,
'elapsed_ms': elapsed, 'code': code,
})
return results
from hotspot.engine import run_fetch
return await run_fetch()

View File

@ -1,63 +1,2 @@
# 热点统计:全维度
from datetime import datetime, timedelta
from sqlor.dbpools import DBPools
async def main(request):
db = DBPools()
now = datetime.now()
day_ago = (now - timedelta(hours=24)).strftime('%Y-%m-%d %H:%M:%S')
async with db.sqlorContext('sage') as sor:
items = await sor.R('hotspot_item', {'order': 'heat_score desc'})
sources = await sor.R('hotspot_source', {})
logs = await sor.sqlExe(
"select * from hotspot_fetch_log where start_time >= ${t}$",
{'t': day_ago}
)
status_count = {'emerging': 0, 'rising': 0, 'hot': 0, 'cooling': 0, 'expired': 0}
for item in items:
s = item.get('status', 'emerging')
status_count[s] = status_count.get(s, 0) + 1
sources_active = sum(1 for s in sources if s.get('enabled') == '1')
failures = sum(1 for log in logs if log.get('status') == 'failed')
# Top sources by item count
src_count = {}
for item in items:
sid = item.get('source_id', '')
src_count[sid] = src_count.get(sid, 0) + 1
top_src = sorted(src_count.items(), key=lambda x: x[1], reverse=True)[:5]
src_names = {s['id']: s.get('name', s['id']) for s in sources}
# Avg heat by category
cat_heat = {}
for item in items:
cat = item.get('category', '未分类') or '未分类'
h = float(item.get('heat_score', 0))
if cat not in cat_heat:
cat_heat[cat] = {'sum': 0, 'cnt': 0}
cat_heat[cat]['sum'] += h
cat_heat[cat]['cnt'] += 1
top_categories = sorted(
[{'name': k, 'avg': round(v['sum']/v['cnt'], 1), 'cnt': v['cnt']}
for k, v in cat_heat.items()],
key=lambda x: x['cnt'], reverse=True
)[:10]
return {
'total': len(items),
'emerging': status_count['emerging'],
'rising': status_count['rising'],
'hot': status_count['hot'],
'cooling': status_count['cooling'],
'expired': status_count['expired'],
'sources': sources_active,
'failures': failures,
'top_sources': [{'name': src_names.get(k, k), 'count': v} for k, v in top_src],
'top_categories': top_categories,
'log_count_24h': len(logs),
}
from hotspot.db import get_stats
return await get_stats()