1013 lines
36 KiB
Python
1013 lines
36 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
懂球帝数据爬虫 - 主入口
|
|
用法:
|
|
python main.py standings # 采集所有联赛积分榜
|
|
python main.py schedule # 采集所有联赛赛程
|
|
python main.py menu # 采集比赛类型菜单
|
|
python main.py lottery # 采集香港/澳门六合彩开奖号码
|
|
python main.py match_data # 从la_league同步赛事比赛数据
|
|
python main.py csl_match # 采集中超赛程页面数据
|
|
python main.py nba_match # 采集NBA赛程页面数据
|
|
python main.py cba_match # 采集CBA赛程页面数据
|
|
python main.py epl_match # 采集英超赛程页面数据
|
|
python main.py bundesliga_match # 采集德甲赛程页面数据
|
|
python main.py laliga_match # 采集西甲赛程页面数据
|
|
python main.py seriea_match # 采集意甲赛程页面数据
|
|
python main.py ligue1_match # 采集法甲赛程页面数据
|
|
python main.py ucl_match # 采集欧冠赛程页面数据
|
|
python main.py uel_match # 采集欧联赛程页面数据
|
|
python main.py tennis_match # 采集网球赛程页面数据
|
|
python main.py esports_match # 采集电竞赛程页面数据
|
|
python main.py sports_match # 采集体坛赛程页面数据
|
|
python main.py truth_social # 采集Truth Social帖子并入库
|
|
python main.py crypto_news # 采集加密货币新闻资讯
|
|
python main.py lottery_news # 采集彩票资讯
|
|
python main.py taiwan_lottery_news # 采集台湾彩券资讯
|
|
python main.py hkjc_lottery_news # 采集香港赛马会资讯
|
|
python main.py fifa_worldcup_news # 采集 FIFA 世界杯资讯
|
|
python main.py dqd_worldcup # 采集懂球帝世界杯赛程/积分榜/球员榜
|
|
python main.py nba_news # 采集NBA新闻资讯
|
|
python main.py cba_news # 采集CBA新闻资讯
|
|
python main.py ai_comment_dispatch # Docker AI 评论调度
|
|
python main.py match_live_stream # 抓取赛事直播线路播放流
|
|
python main.py match_odds # 抓取第三方赛事赔率
|
|
python main.py article_content_fetch # 抓取资讯文章正文内容
|
|
python main.py all # 采集积分榜 + 赛程
|
|
python main.py single <season_id> # 采集单个联赛积分榜
|
|
python main.py cron # 启动定时调度
|
|
python main.py init-db # 初始化数据库表
|
|
python main.py test # 测试连通性
|
|
"""
|
|
import asyncio
|
|
import sys
|
|
import time
|
|
import io
|
|
import os
|
|
from pathlib import Path
|
|
from contextlib import redirect_stdout
|
|
|
|
from src.core.config import load_config
|
|
from src.core.error_collector import ErrorCollector
|
|
from src.core.logger import setup_logger
|
|
|
|
|
|
TASK_ALERT_POLICY = {
|
|
"error_report": "never",
|
|
"article_content": "saved_if_candidates",
|
|
"article_content_fetch": "saved_if_candidates",
|
|
"live_detail": "saved_if_candidates",
|
|
"match_finish": "saved_if_candidates",
|
|
"lottery_draw": "saved_if_candidates",
|
|
"lottery_draw_force": "saved_if_candidates",
|
|
"ai_comment_dispatch": "saved_if_candidates",
|
|
"match_live_stream": "saved_if_candidates",
|
|
"match_odds": "saved_if_candidates",
|
|
}
|
|
|
|
DEFAULT_ACTION_TIMEOUT = 600
|
|
ACTION_TIMEOUTS = {
|
|
"csl_match": 180,
|
|
"nba_match": 180,
|
|
"cba_match": 180,
|
|
"epl_match": 180,
|
|
"bundesliga_match": 180,
|
|
"laliga_match": 180,
|
|
"seriea_match": 180,
|
|
"ligue1_match": 180,
|
|
"ucl_match": 180,
|
|
"uel_match": 180,
|
|
"tennis_match": 180,
|
|
"esports_match": 180,
|
|
"sports_match": 180,
|
|
"live_detail": 180,
|
|
"match_finish": 180,
|
|
"lottery_draw": 180,
|
|
"lottery_draw_force": 180,
|
|
"article_content": 600,
|
|
"article_content_fetch": 900,
|
|
"league_news": 900,
|
|
"fifa_worldcup_news": 600,
|
|
"dqd_worldcup": 600,
|
|
"ai_comment_dispatch": 1800,
|
|
"match_live_stream": 180,
|
|
"match_odds": 300,
|
|
}
|
|
|
|
|
|
def get_action_timeout(action: str) -> int:
|
|
return ACTION_TIMEOUTS.get(str(action or "").strip(), DEFAULT_ACTION_TIMEOUT)
|
|
|
|
|
|
class CrawlerActionLock:
|
|
"""Cross-process non-blocking lock for one crawler action."""
|
|
|
|
def __init__(self, action: str, lock_dir: Path | None = None):
|
|
safe_action = "".join(ch if ch.isalnum() or ch in ("-", "_") else "_" for ch in str(action or "all"))
|
|
self.lock_dir = Path(lock_dir or os.environ.get("DQD_LOCK_DIR") or (Path(__file__).resolve().parent / "data" / "locks"))
|
|
self.path = self.lock_dir / f"{safe_action}.lock"
|
|
self._fp = None
|
|
self._locked = False
|
|
|
|
def acquire(self) -> bool:
|
|
self.lock_dir.mkdir(parents=True, exist_ok=True)
|
|
self._fp = open(self.path, "a+b")
|
|
try:
|
|
self._fp.seek(0)
|
|
if os.name == "nt":
|
|
import msvcrt
|
|
msvcrt.locking(self._fp.fileno(), msvcrt.LK_NBLCK, 1)
|
|
else:
|
|
import fcntl
|
|
fcntl.flock(self._fp.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB)
|
|
self._locked = True
|
|
self._fp.seek(0)
|
|
self._fp.truncate()
|
|
self._fp.write(f"pid={os.getpid()} started_at={int(time.time())}\n".encode("utf-8"))
|
|
self._fp.flush()
|
|
return True
|
|
except (BlockingIOError, OSError):
|
|
self.release()
|
|
return False
|
|
|
|
def release(self):
|
|
if not self._fp:
|
|
return
|
|
try:
|
|
if self._locked and os.name == "nt":
|
|
import msvcrt
|
|
self._fp.seek(0)
|
|
msvcrt.locking(self._fp.fileno(), msvcrt.LK_UNLCK, 1)
|
|
elif self._locked:
|
|
import fcntl
|
|
fcntl.flock(self._fp.fileno(), fcntl.LOCK_UN)
|
|
finally:
|
|
self._fp.close()
|
|
self._fp = None
|
|
self._locked = False
|
|
|
|
|
|
def _build_summary_text(action: str, raw_result: dict) -> str:
|
|
if not isinstance(raw_result, dict):
|
|
return f"{action} 执行完成"
|
|
for key in ("summary_text", "summary", "msg", "error"):
|
|
value = str(raw_result.get(key) or "").strip()
|
|
if value:
|
|
return value
|
|
if "count" in raw_result:
|
|
return f"{action} 执行完成: {raw_result.get('count', 0)}"
|
|
return f"{action} 执行完成"
|
|
|
|
|
|
def _normalize_task_result(action: str, raw_result, execution_success: bool) -> dict:
|
|
if isinstance(raw_result, list):
|
|
raw_result = {"results": raw_result}
|
|
raw_result = raw_result if isinstance(raw_result, dict) else {}
|
|
|
|
results = raw_result.get("results")
|
|
if isinstance(results, list):
|
|
candidate_count = sum(int(item.get("candidate_count") or item.get("count") or 0) for item in results if isinstance(item, dict))
|
|
saved_count = sum(int(item.get("saved_count") or item.get("count") or 0) for item in results if isinstance(item, dict))
|
|
raw_success = all(bool(item.get("success", False)) for item in results) if results else bool(execution_success)
|
|
else:
|
|
candidate_count = int(raw_result.get("candidate_count") or 0)
|
|
saved_count = int(raw_result.get("saved_count") or 0)
|
|
if not candidate_count and "pending" in raw_result:
|
|
candidate_count = int(raw_result.get("pending") or 0)
|
|
if not candidate_count and "total" in raw_result:
|
|
candidate_count = int(raw_result.get("total") or 0)
|
|
if not candidate_count and "fail" in raw_result:
|
|
candidate_count = int(raw_result.get("count") or 0) + int(raw_result.get("fail") or 0)
|
|
if not candidate_count and "count" in raw_result:
|
|
candidate_count = int(raw_result.get("count") or 0)
|
|
if not saved_count and "count" in raw_result:
|
|
saved_count = int(raw_result.get("count") or 0)
|
|
raw_success = bool(raw_result.get("success", execution_success))
|
|
|
|
run_stats = ErrorCollector.get().get_run_stats()
|
|
return {
|
|
"action": action,
|
|
"task_name": raw_result.get("task") or action,
|
|
"command": "crawler",
|
|
"params": action,
|
|
"policy": TASK_ALERT_POLICY.get(action, "saved_required"),
|
|
"execution_success": bool(execution_success),
|
|
"success": raw_success,
|
|
"candidate_count": candidate_count,
|
|
"saved_count": saved_count,
|
|
"http_error_count": int(run_stats.get("http_error_count") or 0),
|
|
"severe_http_statuses": list(run_stats.get("severe_http_statuses") or []),
|
|
"error_message": str(raw_result.get("error") or ""),
|
|
"summary_text": _build_summary_text(action, raw_result),
|
|
}
|
|
|
|
|
|
async def cmd_standings():
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
results = await runner.crawl_all_leagues("standings")
|
|
_print_results(results, "积分榜")
|
|
return results
|
|
|
|
|
|
async def cmd_schedule():
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
results = await runner.crawl_all_leagues("schedule")
|
|
_print_results(results, "赛程")
|
|
return results
|
|
|
|
|
|
async def cmd_menu():
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
result = await runner.crawl_match_menu()
|
|
if result["success"]:
|
|
print(f"✅ 比赛菜单采集成功: {result['count']} 条")
|
|
else:
|
|
print(f"❌ 比赛菜单采集失败: {result.get('error')}")
|
|
return result
|
|
|
|
|
|
async def cmd_news():
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
result = await runner.crawl_news()
|
|
if result["success"]:
|
|
print(f"✅ 新闻采集成功: {result['count']} 条, {result.get('elapsed', 0):.1f}s")
|
|
else:
|
|
print(f"❌ 新闻采集失败: {result.get('error')}")
|
|
return result
|
|
|
|
|
|
async def cmd_league_news():
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
result = await runner.crawl_league_news()
|
|
if result["success"]:
|
|
print(result.get("summary", f"✅ 联赛资讯采集成功: {result['count']} 条"))
|
|
else:
|
|
print(f"❌ 联赛资讯采集失败: {result.get('error')}")
|
|
return result
|
|
|
|
|
|
async def cmd_article_content():
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
result = await runner.crawl_article_content()
|
|
if result["success"]:
|
|
print(result.get("summary", f"✅ 文章内容补全: {result['count']} 篇"))
|
|
else:
|
|
print(f"❌ 文章内容补全失败: {result.get('error')}")
|
|
return result
|
|
|
|
|
|
async def cmd_video():
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
result = await runner.crawl_videos()
|
|
if result["success"]:
|
|
print(f"✅ 视频采集成功: {result['count']} 条, {result.get('elapsed', 0):.1f}s")
|
|
else:
|
|
print(f"❌ 视频采集失败: {result.get('error')}")
|
|
return result
|
|
|
|
|
|
async def cmd_content():
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
result = await runner.crawl_article_details()
|
|
if result["success"]:
|
|
msg = result.get("msg", "")
|
|
if msg:
|
|
print(f"✅ {msg}")
|
|
else:
|
|
print(f"✅ 文章详情补全: {result['count']}/{result.get('total', 0)} 篇, {result.get('elapsed', 0):.1f}s")
|
|
else:
|
|
print(f"❌ 文章详情补全失败: {result.get('error')}")
|
|
return result
|
|
|
|
|
|
async def cmd_live():
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
result = await runner.crawl_live()
|
|
if result["success"]:
|
|
print(f"✅ 实时比赛更新: {result['count']}/{result.get('pending', 0)} 场, {result.get('elapsed', 0):.1f}s")
|
|
if result.get("msg"):
|
|
print(f" {result['msg']}")
|
|
else:
|
|
print(f"❌ 实时比赛采集失败: {result.get('error')}")
|
|
return result
|
|
|
|
|
|
async def cmd_lottery():
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
result = await runner.crawl_lottery()
|
|
if result["success"]:
|
|
msg = result.get("msg", "")
|
|
if msg:
|
|
print(f"✅ {msg}")
|
|
else:
|
|
print(f"✅ 六合彩开奖采集成功: {result['count']} 条, {result.get('elapsed', 0):.1f}s")
|
|
else:
|
|
print(f"❌ 六合彩开奖采集失败: {result.get('error')}")
|
|
return result
|
|
|
|
|
|
async def _cmd_page_match(crawl_func_name: str, label: str):
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
result = await getattr(runner, crawl_func_name)()
|
|
if result["success"]:
|
|
print(f"✅ {label}赛程采集成功: 新增 {result.get('inserted', 0)}, 更新 {result.get('updated', 0)}, 未变 {result.get('unchanged', 0)}, 共 {result['count']} 场, {result.get('elapsed', 0):.1f}s")
|
|
else:
|
|
print(f"❌ {label}赛程采集失败: {result.get('error')}")
|
|
return result
|
|
|
|
|
|
async def cmd_csl_match():
|
|
return await _cmd_page_match("crawl_csl_matches", "中超")
|
|
|
|
|
|
async def cmd_nba_match():
|
|
return await _cmd_page_match("crawl_nba_matches", "NBA")
|
|
|
|
|
|
async def cmd_cba_match():
|
|
return await _cmd_page_match("crawl_cba_matches", "CBA")
|
|
|
|
|
|
async def cmd_epl_match():
|
|
return await _cmd_page_match("crawl_epl_matches", "英超")
|
|
|
|
|
|
async def cmd_bundesliga_match():
|
|
return await _cmd_page_match("crawl_bundesliga_matches", "德甲")
|
|
|
|
|
|
async def cmd_laliga_match():
|
|
return await _cmd_page_match("crawl_laliga_matches", "西甲")
|
|
|
|
|
|
async def cmd_seriea_match():
|
|
return await _cmd_page_match("crawl_seriea_matches", "意甲")
|
|
|
|
|
|
async def cmd_ligue1_match():
|
|
return await _cmd_page_match("crawl_ligue1_matches", "法甲")
|
|
|
|
|
|
async def cmd_ucl_match():
|
|
return await _cmd_page_match("crawl_ucl_matches", "欧冠")
|
|
|
|
|
|
async def cmd_uel_match():
|
|
return await _cmd_page_match("crawl_uel_matches", "欧联")
|
|
|
|
|
|
async def cmd_tennis_match():
|
|
return await _cmd_page_match("crawl_tennis_matches", "网球")
|
|
|
|
|
|
async def cmd_esports_match():
|
|
return await _cmd_page_match("crawl_esports_matches", "电竞")
|
|
|
|
|
|
async def cmd_sports_match():
|
|
return await _cmd_page_match("crawl_sports_matches", "体坛")
|
|
|
|
|
|
async def cmd_truth_social():
|
|
import asyncio
|
|
from truth_social import run
|
|
from src.core.config import load_config
|
|
cfg = load_config().database
|
|
db_config = {
|
|
"host": cfg.host, "port": cfg.port, "user": cfg.username,
|
|
"password": cfg.password, "database": cfg.database, "charset": cfg.charset,
|
|
}
|
|
return await asyncio.to_thread(run, db_config)
|
|
|
|
|
|
async def cmd_crypto_news():
|
|
import asyncio
|
|
from crypto_news import run
|
|
from src.core.config import load_config
|
|
cfg = load_config().database
|
|
db_config = {
|
|
"host": cfg.host, "port": cfg.port, "user": cfg.username,
|
|
"password": cfg.password, "database": cfg.database, "charset": cfg.charset,
|
|
}
|
|
return await asyncio.to_thread(run, db_config)
|
|
|
|
|
|
async def cmd_lottery_news():
|
|
import asyncio
|
|
from lottery_news import run
|
|
from src.core.config import load_config
|
|
cfg = load_config().database
|
|
db_config = {
|
|
"host": cfg.host, "port": cfg.port, "user": cfg.username,
|
|
"password": cfg.password, "database": cfg.database, "charset": cfg.charset,
|
|
}
|
|
return await asyncio.to_thread(run, db_config)
|
|
|
|
|
|
async def cmd_taiwan_lottery_news():
|
|
import asyncio
|
|
from taiwan_lottery_news import run
|
|
from src.core.config import load_config
|
|
cfg = load_config().database
|
|
db_config = {
|
|
"host": cfg.host, "port": cfg.port, "user": cfg.username,
|
|
"password": cfg.password, "database": cfg.database, "charset": cfg.charset,
|
|
}
|
|
return await asyncio.to_thread(run, db_config)
|
|
|
|
|
|
async def cmd_hkjc_lottery_news():
|
|
import asyncio
|
|
from hkjc_lottery_news import run
|
|
from src.core.config import load_config
|
|
cfg = load_config().database
|
|
db_config = {
|
|
"host": cfg.host, "port": cfg.port, "user": cfg.username,
|
|
"password": cfg.password, "database": cfg.database, "charset": cfg.charset,
|
|
}
|
|
return await asyncio.to_thread(run, db_config)
|
|
|
|
|
|
async def cmd_fifa_worldcup_news():
|
|
import asyncio
|
|
from fifa_worldcup_news import run
|
|
from src.core.config import load_config
|
|
cfg = load_config().database
|
|
db_config = {
|
|
"host": cfg.host, "port": cfg.port, "user": cfg.username,
|
|
"password": cfg.password, "database": cfg.database, "charset": cfg.charset,
|
|
"prefix": cfg.prefix,
|
|
}
|
|
return await asyncio.to_thread(run, db_config)
|
|
|
|
|
|
async def cmd_dqd_worldcup():
|
|
import asyncio
|
|
from dqd_worldcup import run
|
|
from src.core.config import load_config
|
|
cfg = load_config().database
|
|
db_config = {
|
|
"host": cfg.host, "port": cfg.port, "user": cfg.username,
|
|
"password": cfg.password, "database": cfg.database, "charset": cfg.charset,
|
|
"prefix": cfg.prefix,
|
|
}
|
|
return await asyncio.to_thread(run, db_config)
|
|
|
|
|
|
async def cmd_lottery_community():
|
|
import asyncio
|
|
from lottery_community import run
|
|
from src.core.config import load_config
|
|
cfg = load_config().database
|
|
db_config = {
|
|
"host": cfg.host, "port": cfg.port, "user": cfg.username,
|
|
"password": cfg.password, "database": cfg.database, "charset": cfg.charset,
|
|
"prefix": cfg.prefix,
|
|
}
|
|
return await asyncio.to_thread(run, db_config)
|
|
|
|
|
|
async def cmd_article_content_fetch():
|
|
import asyncio
|
|
from article_fetcher import run
|
|
from src.core.config import load_config
|
|
cfg = load_config().database
|
|
db_config = {
|
|
"host": cfg.host, "port": cfg.port, "user": cfg.username,
|
|
"password": cfg.password, "database": cfg.database, "charset": cfg.charset,
|
|
}
|
|
return await asyncio.to_thread(run, db_config)
|
|
|
|
|
|
async def cmd_nba_news():
|
|
import asyncio
|
|
from nba_news import run
|
|
from src.core.config import load_config
|
|
cfg = load_config().database
|
|
db_config = {
|
|
"host": cfg.host, "port": cfg.port, "user": cfg.username,
|
|
"password": cfg.password, "database": cfg.database, "charset": cfg.charset,
|
|
}
|
|
return await asyncio.to_thread(run, db_config)
|
|
|
|
|
|
async def cmd_cba_news():
|
|
import asyncio
|
|
from cba_news import run
|
|
from src.core.config import load_config
|
|
cfg = load_config().database
|
|
db_config = {
|
|
"host": cfg.host, "port": cfg.port, "user": cfg.username,
|
|
"password": cfg.password, "database": cfg.database, "charset": cfg.charset,
|
|
}
|
|
return await asyncio.to_thread(run, db_config)
|
|
|
|
|
|
async def cmd_ai_comment_dispatch():
|
|
import asyncio
|
|
from ai_comment_dispatch import run
|
|
from src.core.config import load_config
|
|
|
|
cfg = load_config().database
|
|
db_config = {
|
|
"host": cfg.host, "port": cfg.port, "user": cfg.username,
|
|
"password": cfg.password, "database": cfg.database, "charset": cfg.charset,
|
|
"prefix": cfg.prefix,
|
|
}
|
|
return await asyncio.to_thread(run, db_config)
|
|
|
|
|
|
async def cmd_match_live_stream():
|
|
import asyncio
|
|
from match_live_stream import run
|
|
from src.core.config import load_config
|
|
|
|
cfg = load_config().database
|
|
db_config = {
|
|
"host": cfg.host,
|
|
"port": cfg.port,
|
|
"user": cfg.username,
|
|
"password": cfg.password,
|
|
"database": cfg.database,
|
|
"charset": cfg.charset,
|
|
"prefix": cfg.prefix,
|
|
}
|
|
return await asyncio.to_thread(run, db_config)
|
|
|
|
|
|
async def cmd_match_odds():
|
|
import asyncio
|
|
from match_odds import run
|
|
from src.core.config import load_config
|
|
|
|
cfg = load_config().database
|
|
db_config = {
|
|
"host": cfg.host,
|
|
"port": cfg.port,
|
|
"user": cfg.username,
|
|
"password": cfg.password,
|
|
"database": cfg.database,
|
|
"charset": cfg.charset,
|
|
"prefix": cfg.prefix,
|
|
}
|
|
result = await asyncio.to_thread(run, db_config)
|
|
prefix = "✅" if result.get("success") else "❌"
|
|
print(f"{prefix} {result.get('summary_text', '第三方赛事赔率抓取完成')}")
|
|
if result.get("error"):
|
|
print(f" {result['error']}")
|
|
return result
|
|
|
|
|
|
async def cmd_match_data():
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
result = await runner.crawl_match_data()
|
|
if result["success"]:
|
|
print(f"✅ 赛事数据同步成功: {result['count']} 场比赛, {result.get('rounds', 0)} 轮次, {result.get('elapsed', 0):.1f}s")
|
|
if result.get("msg"):
|
|
print(f" {result['msg']}")
|
|
else:
|
|
print(f"❌ 赛事数据同步失败: {result.get('error')}")
|
|
return result
|
|
|
|
|
|
async def cmd_match_finish():
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
result = await runner.crawl_match_finish()
|
|
if result["success"]:
|
|
print(f"✅ 超时比赛收尾完成: {result['count']} 场, 详情成功 {result.get('detail_ok', 0)} 场, 跳过 {result.get('skipped', 0)} 场, {result.get('elapsed', 0):.1f}s")
|
|
if result.get("msg"):
|
|
print(f" {result['msg']}")
|
|
else:
|
|
print(f"❌ 超时比赛收尾失败: {result.get('error')}")
|
|
return result
|
|
|
|
|
|
async def cmd_live_detail():
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
result = await runner.crawl_live_detail()
|
|
if result["success"]:
|
|
print(f"✅ 实时比赛详情更新: {result['count']}/{result.get('pending', 0)} 场, "
|
|
f"结束 {result.get('finished', 0)} 场, 失败 {result.get('failed', 0)} 场, "
|
|
f"文字直播 {result.get('live_text', 0)} 条, {result.get('elapsed', 0):.1f}s")
|
|
if result.get("msg"):
|
|
print(f" {result['msg']}")
|
|
else:
|
|
print(f"❌ 实时比赛详情更新失败: {result.get('error')}")
|
|
return result
|
|
|
|
|
|
async def cmd_lottery_draw():
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
result = await runner.crawl_lottery_draw()
|
|
if result["success"]:
|
|
print(f"✅ 彩种开奖数据采集成功: {result['count']} 条, {result.get('elapsed', 0):.1f}s")
|
|
if result.get("msg"):
|
|
print(f" {result['msg']}")
|
|
else:
|
|
print(f"❌ 彩种开奖数据采集失败: {result.get('error')}")
|
|
return result
|
|
|
|
|
|
async def cmd_lottery_draw_force():
|
|
"""强制全量刷新所有彩种最新开奖数据(跳过时间检查,用于断采后补数据)"""
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
result = await runner.crawl_lottery_draw(force=True)
|
|
if result["success"]:
|
|
print(f"✅ 彩种开奖数据强制刷新成功: {result['count']} 条, {result.get('elapsed', 0):.1f}s")
|
|
else:
|
|
print(f"❌ 彩种开奖数据强制刷新失败: {result.get('error')}")
|
|
result["task"] = "lottery_draw_force"
|
|
return result
|
|
|
|
|
|
async def cmd_error_report():
|
|
"""企业微信发送未恢复的定时任务告警"""
|
|
from src.core.alert_manager import CrontabAlertManager
|
|
manager = CrontabAlertManager()
|
|
result = await manager.dispatch_open_alerts(limit=50)
|
|
print(result.get("summary_text", "告警分发完成"))
|
|
return result
|
|
|
|
|
|
async def cmd_all():
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
print("=" * 60)
|
|
print("采集积分榜")
|
|
print("=" * 60)
|
|
r1 = await runner.crawl_all_leagues("standings")
|
|
_print_results(r1, "积分榜")
|
|
|
|
print()
|
|
print("=" * 60)
|
|
print("采集赛程")
|
|
print("=" * 60)
|
|
r2 = await runner.crawl_all_leagues("schedule")
|
|
_print_results(r2, "赛程")
|
|
return {
|
|
"results": [
|
|
{"task": "standings_batch", "success": all(r.get("success", False) for r in r1), "count": sum(int(r.get("count", 0) or 0) for r in r1)},
|
|
{"task": "schedule_batch", "success": all(r.get("success", False) for r in r2), "count": sum(int(r.get("count", 0) or 0) for r in r2)},
|
|
]
|
|
}
|
|
|
|
|
|
async def cmd_single(season_id: int):
|
|
from src.scheduler.task_runner import TaskRunner
|
|
async with TaskRunner() as runner:
|
|
result = await runner.crawl_standings(season_id, f"season_{season_id}")
|
|
if result["success"]:
|
|
print(f"✅ 采集成功: {result['count']} 条, {result['elapsed']:.1f}s")
|
|
else:
|
|
print(f"❌ 采集失败: {result.get('error')}")
|
|
return result
|
|
|
|
|
|
async def cmd_cron():
|
|
from src.scheduler.cron_scheduler import CronScheduler
|
|
scheduler = CronScheduler()
|
|
await scheduler.start()
|
|
print("定时调度器已启动,按 Ctrl+C 停止")
|
|
try:
|
|
while True:
|
|
await asyncio.sleep(1)
|
|
except (KeyboardInterrupt, asyncio.CancelledError):
|
|
await scheduler.stop()
|
|
print("调度器已停止")
|
|
|
|
|
|
async def cmd_init_db():
|
|
from src.storage.database import Database
|
|
async with Database() as db:
|
|
await db.ensure_tables()
|
|
print("✅ 数据库表初始化完成")
|
|
return {"success": True, "candidate_count": 1, "saved_count": 1, "summary_text": "数据库表初始化完成"}
|
|
|
|
|
|
async def cmd_test():
|
|
"""测试各通道连通性"""
|
|
from src.engine.hybrid_engine import HybridEngine
|
|
print("=" * 60)
|
|
print("懂球帝爬虫连通性测试")
|
|
print("=" * 60)
|
|
|
|
async with HybridEngine() as engine:
|
|
# 测试 API 通道
|
|
print("\n1. 测试 API 通道 (积分榜 - 中超)...")
|
|
try:
|
|
start = time.time()
|
|
data = await engine.get_standings(26322)
|
|
elapsed = time.time() - start
|
|
rounds = data.get("content", {}).get("rounds", [])
|
|
if rounds:
|
|
teams = rounds[0].get("content", {}).get("data", [])
|
|
print(f" ✅ 成功! {len(teams)} 支球队, {elapsed:.2f}s")
|
|
else:
|
|
print(f" ⚠️ 返回空数据, {elapsed:.2f}s")
|
|
except Exception as e:
|
|
print(f" ❌ 失败: {e}")
|
|
|
|
# 测试比赛菜单
|
|
print("\n2. 测试比赛菜单...")
|
|
try:
|
|
start = time.time()
|
|
data = await engine.get_match_menu()
|
|
elapsed = time.time() - start
|
|
if data.get("errCode") == 0:
|
|
items = data.get("data", {}).get("list", [])
|
|
print(f" ✅ 成功! {len(items)} 个比赛类型, {elapsed:.2f}s")
|
|
else:
|
|
print(f" ⚠️ 返回错误: {data}, {elapsed:.2f}s")
|
|
except Exception as e:
|
|
print(f" ❌ 失败: {e}")
|
|
|
|
# 测试数据库
|
|
print("\n3. 测试数据库连接...")
|
|
try:
|
|
from src.storage.database import Database
|
|
async with Database() as db:
|
|
result = await db.fetchone("SELECT 1 AS ok")
|
|
if result and result.get("ok") == 1:
|
|
print(" ✅ 数据库连接正常")
|
|
else:
|
|
print(" ❌ 数据库连接异常")
|
|
except Exception as e:
|
|
print(f" ❌ 数据库连接失败: {e}")
|
|
|
|
print("\n" + "=" * 60)
|
|
print("引擎统计:")
|
|
stats = engine.stats
|
|
for ch, s in stats.get("channels", {}).items():
|
|
print(f" {ch}: 成功={s['success']}, 失败={s['fail']}")
|
|
print("=" * 60)
|
|
return {"success": True, "candidate_count": 3, "saved_count": 3, "summary_text": "连通性测试完成"}
|
|
|
|
|
|
def _parse_extra_args(argv):
|
|
"""解析 --log-id=N --crontab-id=N 参数"""
|
|
extra = {}
|
|
for arg in argv:
|
|
if arg.startswith('--log-id='):
|
|
extra['log_id'] = int(arg.split('=', 1)[1])
|
|
elif arg.startswith('--crontab-id='):
|
|
extra['crontab_id'] = int(arg.split('=', 1)[1])
|
|
return extra
|
|
|
|
|
|
class _CrontabLogWriter:
|
|
"""实时写入执行日志到数据库"""
|
|
|
|
def __init__(self, log_id: int, crontab_id: int):
|
|
self.log_id = log_id
|
|
self.crontab_id = crontab_id
|
|
self._conn = None
|
|
self._cfg = None
|
|
|
|
async def connect(self):
|
|
import aiomysql
|
|
self._cfg = load_config().database
|
|
self._conn = await aiomysql.connect(
|
|
host=self._cfg.host, port=self._cfg.port, user=self._cfg.username,
|
|
password=self._cfg.password, db=self._cfg.database, charset=self._cfg.charset,
|
|
)
|
|
|
|
async def flush_output(self, output: str, elapsed: float):
|
|
"""实时刷新 output 到日志表(status 保持 0=执行中)"""
|
|
async with self._conn.cursor() as cur:
|
|
await cur.execute(
|
|
f"UPDATE `{self._cfg.prefix}dev_crontab_log` SET output=%s, elapsed=%s WHERE id=%s",
|
|
(output[:5000], elapsed, self.log_id)
|
|
)
|
|
await self._conn.commit()
|
|
|
|
async def finish(self, success: bool, output: str, elapsed: float):
|
|
"""最终更新状态"""
|
|
async with self._conn.cursor() as cur:
|
|
status = 1 if success else 2
|
|
await cur.execute(
|
|
f"UPDATE `{self._cfg.prefix}dev_crontab_log` SET status=%s, output=%s, elapsed=%s WHERE id=%s",
|
|
(status, output[:5000], elapsed, self.log_id)
|
|
)
|
|
if success:
|
|
await cur.execute(
|
|
f"UPDATE `{self._cfg.prefix}dev_crontab` SET last_time=%s, `time`=%s, error='' WHERE id=%s",
|
|
(int(time.time()), elapsed, self.crontab_id)
|
|
)
|
|
else:
|
|
await cur.execute(
|
|
f"UPDATE `{self._cfg.prefix}dev_crontab` SET error=%s WHERE id=%s",
|
|
(output[:500], self.crontab_id)
|
|
)
|
|
await self._conn.commit()
|
|
|
|
async def close(self):
|
|
if self._conn:
|
|
self._conn.close()
|
|
|
|
|
|
def _print_results(results, task_name):
|
|
success = sum(1 for r in results if r.get("success"))
|
|
total = len(results)
|
|
print(f"\n{task_name}采集完成: {success}/{total} 成功")
|
|
for r in results:
|
|
s = "✅" if r.get("success") else "❌"
|
|
league = r.get("league", "?")
|
|
count = r.get("count", 0)
|
|
elapsed = r.get("elapsed", 0)
|
|
err = r.get("error", "")
|
|
if r.get("success"):
|
|
print(f" {s} {league}: {count} 条, {elapsed:.1f}s")
|
|
else:
|
|
print(f" {s} {league}: {err}")
|
|
|
|
|
|
def main():
|
|
load_config()
|
|
setup_logger()
|
|
|
|
if len(sys.argv) < 2:
|
|
print(__doc__)
|
|
sys.exit(0)
|
|
|
|
cmd = sys.argv[1].lower()
|
|
extra = _parse_extra_args(sys.argv[2:])
|
|
|
|
cmd_map = {
|
|
"standings": cmd_standings,
|
|
"schedule": cmd_schedule,
|
|
"menu": cmd_menu,
|
|
"news": cmd_news,
|
|
"league_news": cmd_league_news,
|
|
"article_content": cmd_article_content,
|
|
"video": cmd_video,
|
|
"content": cmd_content,
|
|
"live": cmd_live,
|
|
"lottery": cmd_lottery,
|
|
"match_data": cmd_match_data,
|
|
"csl_match": cmd_csl_match,
|
|
"nba_match": cmd_nba_match,
|
|
"cba_match": cmd_cba_match,
|
|
"epl_match": cmd_epl_match,
|
|
"bundesliga_match": cmd_bundesliga_match,
|
|
"laliga_match": cmd_laliga_match,
|
|
"seriea_match": cmd_seriea_match,
|
|
"ligue1_match": cmd_ligue1_match,
|
|
"ucl_match": cmd_ucl_match,
|
|
"uel_match": cmd_uel_match,
|
|
"tennis_match": cmd_tennis_match,
|
|
"esports_match": cmd_esports_match,
|
|
"sports_match": cmd_sports_match,
|
|
"truth_social": cmd_truth_social,
|
|
"crypto_news": cmd_crypto_news,
|
|
"lottery_news": cmd_lottery_news,
|
|
"taiwan_lottery_news": cmd_taiwan_lottery_news,
|
|
"hkjc_lottery_news": cmd_hkjc_lottery_news,
|
|
"fifa_worldcup_news": cmd_fifa_worldcup_news,
|
|
"dqd_worldcup": cmd_dqd_worldcup,
|
|
"lottery_community": cmd_lottery_community,
|
|
"nba_news": cmd_nba_news,
|
|
"cba_news": cmd_cba_news,
|
|
"ai_comment_dispatch": cmd_ai_comment_dispatch,
|
|
"match_live_stream": cmd_match_live_stream,
|
|
"match_odds": cmd_match_odds,
|
|
"article_content_fetch": cmd_article_content_fetch,
|
|
"match_finish": cmd_match_finish,
|
|
"live_detail": cmd_live_detail,
|
|
"lottery_draw": cmd_lottery_draw,
|
|
"lottery_draw_force": cmd_lottery_draw_force,
|
|
"error_report": cmd_error_report,
|
|
"all": cmd_all,
|
|
"cron": cmd_cron,
|
|
"init-db": cmd_init_db,
|
|
"test": cmd_test,
|
|
}
|
|
|
|
if cmd == "single":
|
|
plain_args = [a for a in sys.argv[2:] if not a.startswith('--')]
|
|
if not plain_args:
|
|
print("用法: python main.py single <season_id>")
|
|
sys.exit(1)
|
|
coro = cmd_single(int(plain_args[0]))
|
|
elif cmd in cmd_map:
|
|
coro = cmd_map[cmd]()
|
|
else:
|
|
print(f"未知命令: {cmd}")
|
|
print(__doc__)
|
|
sys.exit(1)
|
|
|
|
if 'log_id' in extra and 'crontab_id' in extra:
|
|
asyncio.run(_run_with_log(coro, extra['log_id'], extra['crontab_id'], cmd))
|
|
else:
|
|
result = asyncio.run(coro)
|
|
normalized = _normalize_task_result(cmd, result, True)
|
|
if not normalized.get("success", True):
|
|
sys.exit(1)
|
|
|
|
|
|
async def _run_with_log(coro, log_id: int, crontab_id: int, action: str):
|
|
"""带实时日志刷新的异步执行"""
|
|
from loguru import logger as _logger
|
|
|
|
writer = _CrontabLogWriter(log_id, crontab_id)
|
|
await asyncio.wait_for(writer.connect(), timeout=15)
|
|
action_lock = CrawlerActionLock(action)
|
|
if not action_lock.acquire():
|
|
if hasattr(coro, "close"):
|
|
coro.close()
|
|
summary = f"{action} 已有任务执行中,跳过本次"
|
|
await writer.finish(True, summary, 0)
|
|
await writer.close()
|
|
return
|
|
|
|
ErrorCollector.get().reset_run_stats()
|
|
|
|
log_buf = io.StringIO()
|
|
stdout_buf = io.StringIO()
|
|
start_time = time.time()
|
|
execution_success = True
|
|
task_done = False
|
|
task_result = None
|
|
|
|
sink_id = _logger.add(log_buf, format="{time:HH:mm:ss} | {level:<8} | {message}", level="INFO")
|
|
|
|
def get_output():
|
|
parts = []
|
|
s = stdout_buf.getvalue()
|
|
if s:
|
|
parts.append(s.strip())
|
|
l = log_buf.getvalue()
|
|
if l:
|
|
parts.append(l.strip())
|
|
return "\n".join(parts) if parts else ""
|
|
|
|
async def flush_loop():
|
|
last_hash = 0
|
|
while not task_done:
|
|
await asyncio.sleep(3)
|
|
current = get_output()
|
|
h = len(current)
|
|
if h != last_hash:
|
|
last_hash = h
|
|
elapsed = round(time.time() - start_time, 2)
|
|
try:
|
|
await writer.flush_output(current, elapsed)
|
|
except Exception:
|
|
pass
|
|
|
|
flush_task = asyncio.create_task(flush_loop())
|
|
try:
|
|
with redirect_stdout(stdout_buf):
|
|
task_result = await asyncio.wait_for(coro, timeout=get_action_timeout(action))
|
|
except asyncio.TimeoutError:
|
|
execution_success = False
|
|
task_result = {"success": False, "error": "timeout", "summary_text": f"{action} 执行超时"}
|
|
stdout_buf.write(f"\n错误: {action} 执行超过 {get_action_timeout(action)} 秒,已终止")
|
|
except Exception as e:
|
|
execution_success = False
|
|
task_result = {"success": False, "error": str(e), "summary_text": f"执行异常: {e}"}
|
|
stdout_buf.write(f"\n错误: {e}")
|
|
finally:
|
|
task_done = True
|
|
_logger.remove(sink_id)
|
|
await asyncio.sleep(0)
|
|
flush_task.cancel()
|
|
try:
|
|
await flush_task
|
|
except asyncio.CancelledError:
|
|
pass
|
|
|
|
elapsed = round(time.time() - start_time, 2)
|
|
output = get_output()
|
|
normalized = _normalize_task_result(action, task_result, execution_success)
|
|
try:
|
|
await writer.finish(execution_success, output or normalized["summary_text"] or '执行完成', elapsed)
|
|
try:
|
|
from src.core.alert_manager import CrontabAlertManager
|
|
await CrontabAlertManager().record_crawler_result(action, crontab_id, log_id, normalized, output)
|
|
except Exception as alert_error:
|
|
print(f"告警记录失败: {alert_error}")
|
|
finally:
|
|
action_lock.release()
|
|
await writer.close()
|
|
await ErrorCollector.get().close()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|