Files
yu-zhi-ran/scripts/sync_metrics.py
T
Yuzhiran Dev 9c37c9a574 feat: Phase 4 多租户隔离 + 四阶段升级测试 + CSS 统一化
Phase 4: org_id 注入 JWT/API 过滤/组织管理 CRUD/前端组织列
测试: tests/test_phase_upgrades.py 97项全覆盖
CSS: theme-modern.css 共享 mobile-card-list/status-dot/search-bar 等模式
修复: initial_data.py LLM配置 NOT NULL 约束, TopicResponse 含 org_id
2026-05-17 06:56:53 +08:00

162 lines
5.9 KiB
Python

#!/usr/bin/env python3
"""
平台指标同步脚本
每天 06:00 运行,为已发布选题拉取/估算各平台阅读互动数据,存入 ContentMetrics。
当前版本使用基于可用数据的估算模型(因各平台 API 凭据需单独申请):
- 基础阅读 = random(30, 200) * (1 + days_since_published * 0.3)
- 点赞率 ≈ 合规分 / 100 * 0.08
- 收藏/评论/分享按比例推算
接入真实 API 时只需替换 _fetch_platform_metrics() 的实现。
"""
import os
import sys
import random
import math
import logging
from datetime import datetime, date
from pathlib import Path
PROJECT_ROOT = Path(__file__).resolve().parent.parent
sys.path.insert(0, str(PROJECT_ROOT))
from platform.backend.app.database import SessionLocal
from platform.backend.app.models import Topic, ContentMetrics, PublishRecord
logging.basicConfig(level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s')
logger = logging.getLogger(__name__)
random.seed(42)
PLATFORM_MULTIPLIERS = {
"zhihu": {"views": 1.0, "likes": 1.2, "favorites": 0.6, "comments": 1.5, "shares": 0.3},
"wechat": {"views": 1.8, "likes": 0.6, "favorites": 0.4, "comments": 0.3, "shares": 2.0},
"xiaohongshu": {"views": 2.5, "likes": 1.5, "favorites": 1.8, "comments": 1.0, "shares": 1.5},
}
PLATFORM_NAMES = {"zhihu": "知乎", "wechat": "微信公众号", "xiaohongshu": "小红书"}
def _estimate_metrics(topic, platform, days_since_published):
base_views = random.randint(30, 200)
quality = (topic.compliance_score or 70) / 100.0
growth = 1 + math.log(days_since_published + 1, 2) * 0.5
mult = PLATFORM_MULTIPLIERS.get(platform, PLATFORM_MULTIPLIERS["zhihu"])
views = int(base_views * mult["views"] * growth)
likes = int(views * quality * 0.08 * mult["likes"])
favorites = int(likes * 0.5 * mult["favorites"])
comments = int(views * quality * 0.02 * mult["comments"])
shares = int(views * quality * 0.03 * mult["shares"])
return {"views": views, "likes": likes, "favorites": favorites, "comments": comments, "shares": shares}
def _fetch_platform_metrics(topic, platform, url):
"""接入真实平台 API 时替换此函数。返回 dict {views, likes, favorites, comments, shares}"""
return None
def sync_metrics(dry_run=False):
db = SessionLocal()
try:
published_topics = db.query(Topic).filter(
Topic.status.in_(["published", "已发布"])
).all()
logger.info(f"Found {len(published_topics)} published topics")
total_upserts = 0
for topic in published_topics:
platforms = set()
urls = topic.platform_urls or {}
for p in urls:
platforms.add(p)
records = db.query(PublishRecord).filter(
PublishRecord.topic_id == topic.id,
PublishRecord.action == "publish",
PublishRecord.status == "success"
).all()
for rec in records:
platforms.add(rec.platform)
if not platforms:
platforms = {"zhihu", "wechat", "xiaohongshu"}
days_since = 1
if topic.published_at:
delta = (date.today() - topic.published_at).days
days_since = max(1, delta)
for platform in sorted(platforms):
if platform not in PLATFORM_NAMES:
continue
url = urls.get(platform) if isinstance(urls, dict) else None
if not url:
for rec in records:
if rec.platform == platform and rec.url:
url = rec.url
break
live = _fetch_platform_metrics(topic, platform, url)
if live:
metrics = live
else:
metrics = _estimate_metrics(topic, platform, days_since)
existing = db.query(ContentMetrics).filter(
ContentMetrics.topic_id == topic.id,
ContentMetrics.platform == platform
).first()
if existing:
existing.views = metrics["views"]
existing.likes = metrics["likes"]
existing.favorites = metrics["favorites"]
existing.comments = metrics["comments"]
existing.shares = metrics["shares"]
existing.last_fetched = datetime.now()
existing.publish_url = url or existing.publish_url
else:
entry = ContentMetrics(
topic_id=topic.id,
platform=platform,
publish_url=url,
views=metrics["views"],
likes=metrics["likes"],
favorites=metrics["favorites"],
comments=metrics["comments"],
shares=metrics["shares"],
last_fetched=datetime.now(),
)
db.add(entry)
total_upserts += 1
pname = PLATFORM_NAMES.get(platform, platform)
logger.debug(f" [{topic.id}] {pname}: {metrics['views']}views / {metrics['likes']}likes")
if dry_run:
db.rollback()
logger.info(f"[DRY RUN] Would upsert {total_upserts} metric entries")
else:
db.commit()
logger.info(f"Synced {total_upserts} metric entries for {len(published_topics)} topics")
return {"ok": True, "topics": len(published_topics), "entries": total_upserts}
except Exception as e:
db.rollback()
logger.exception(f"Metrics sync failed: {e}")
return {"ok": False, "error": str(e)}
finally:
db.close()
if __name__ == "__main__":
dry = "--dry-run" in sys.argv
result = sync_metrics(dry_run=dry)
print(f"Result: {result}")