From bc9aef5a2f00dbc63cb853c797c77c27ccb2432d Mon Sep 17 00:00:00 2001 From: hajimi Date: Sun, 21 Jun 2026 12:15:40 +0800 Subject: [PATCH] feat: add docker match live stream crawler --- docker/crawler/config/crawler_tasks.yaml | 5 + docker/crawler/main.py | 22 ++ docker/crawler/match_live_stream.py | 260 ++++++++++++++++++ .../crawler/tests/test_match_live_stream.py | 44 +++ docs/sql/insert_crontab_match_live_stream.sql | 20 ++ 5 files changed, 351 insertions(+) create mode 100644 docker/crawler/match_live_stream.py create mode 100644 docker/crawler/tests/test_match_live_stream.py create mode 100644 docs/sql/insert_crontab_match_live_stream.sql diff --git a/docker/crawler/config/crawler_tasks.yaml b/docker/crawler/config/crawler_tasks.yaml index 635cc21..eba8150 100644 --- a/docker/crawler/config/crawler_tasks.yaml +++ b/docker/crawler/config/crawler_tasks.yaml @@ -154,3 +154,8 @@ tasks: action: ai_comment_dispatch cron: "*/10 * * * *" active: true + - key: match_live_stream + name: 赛事直播流抓取 + action: match_live_stream + cron: "*/1 * * * *" + active: true diff --git a/docker/crawler/main.py b/docker/crawler/main.py index 775c374..dd9cb21 100644 --- a/docker/crawler/main.py +++ b/docker/crawler/main.py @@ -30,6 +30,7 @@ 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 article_content_fetch # 抓取资讯文章正文内容 python main.py all # 采集积分榜 + 赛程 python main.py single # 采集单个联赛积分榜 @@ -59,6 +60,7 @@ TASK_ALERT_POLICY = { "lottery_draw": "saved_if_candidates", "lottery_draw_force": "saved_if_candidates", "ai_comment_dispatch": "saved_if_candidates", + "match_live_stream": "saved_if_candidates", } DEFAULT_ACTION_TIMEOUT = 600 @@ -86,6 +88,7 @@ ACTION_TIMEOUTS = { "fifa_worldcup_news": 600, "dqd_worldcup": 600, "ai_comment_dispatch": 1800, + "match_live_stream": 180, } @@ -522,6 +525,24 @@ async def cmd_ai_comment_dispatch(): 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_data(): from src.scheduler.task_runner import TaskRunner async with TaskRunner() as runner: @@ -837,6 +858,7 @@ def main(): "nba_news": cmd_nba_news, "cba_news": cmd_cba_news, "ai_comment_dispatch": cmd_ai_comment_dispatch, + "match_live_stream": cmd_match_live_stream, "article_content_fetch": cmd_article_content_fetch, "match_finish": cmd_match_finish, "live_detail": cmd_live_detail, diff --git a/docker/crawler/match_live_stream.py b/docker/crawler/match_live_stream.py new file mode 100644 index 0000000..af4fb5d --- /dev/null +++ b/docker/crawler/match_live_stream.py @@ -0,0 +1,260 @@ +#!/usr/bin/env python3 +from __future__ import annotations + +import json +import os +import re +import time +from dataclasses import dataclass +from pathlib import Path +from typing import Any, Dict, List + +import pymysql +import requests + + +DETAIL_URL_TEMPLATE = "https://json.ncctrials.com/room/{room_id}/detail.json?v={ts}" +FETCH_BATCH_SIZE = 30 +FETCH_TIMEOUT = 15 +FETCH_ERROR_LIMIT = 500 +ROOM_ID_RE = re.compile(r"/room/(\d+)") +DETAIL_JSONP_RE = re.compile(r"detail\((\{.*\})\)\s*$", re.S) +DEFAULT_HEADERS = { + "User-Agent": ( + "Mozilla/5.0 (Windows NT 10.0; Win64; x64) " + "AppleWebKit/537.36 (KHTML, like Gecko) " + "Chrome/137.0.0.0 Safari/537.36" + ), + "Accept": "application/json,text/javascript,*/*;q=0.1", + "Referer": "https://yyzb1.tv/", +} + + +def parse_match_live_detail_html(body: str) -> dict: + match = DETAIL_JSONP_RE.search((body or "").strip()) + if not match: + raise ValueError("detail jsonp not found") + + payload = json.loads(match.group(1)) + stream = payload.get("data", {}).get("stream", {}) + parsed = { + "play_url_flv": str(stream.get("flv") or "").strip(), + "play_url_hd_flv": str(stream.get("hdFlv") or "").strip(), + "play_url_m3u8": str(stream.get("m3u8") or "").strip(), + "play_url_hd_m3u8": str(stream.get("hdM3u8") or "").strip(), + } + if not any(parsed.values()): + raise ValueError("stream urls empty") + return parsed + + +def choose_next_fetch_at(now_ts: int, has_existing_stream: bool) -> int: + return int(now_ts) + (600 if has_existing_stream else 60) + + +def _safe_int(value: Any, default: int = 0) -> int: + try: + return int(value) + except (TypeError, ValueError): + return default + + +def _trim_error(message: Any) -> str: + return str(message or "").strip()[:FETCH_ERROR_LIMIT] + + +def _derive_room_id(source_url: str) -> str: + match = ROOM_ID_RE.search(str(source_url or "").strip()) + if not match: + raise ValueError("room id not found in source_url") + return match.group(1) + + +def _build_detail_url(source_url: str, now_ts: int) -> str: + room_id = _derive_room_id(source_url) + return DETAIL_URL_TEMPLATE.format(room_id=room_id, ts=int(now_ts)) + + +def _has_existing_stream(row: Dict[str, Any]) -> bool: + return any( + str(row.get(field) or "").strip() + for field in ("play_url_m3u8", "play_url_hd_m3u8", "play_url_flv", "play_url_hd_flv") + ) + + +def _fetch_stream_payload(source_url: str, now_ts: int) -> Dict[str, str]: + detail_url = _build_detail_url(source_url, now_ts) + response = requests.get(detail_url, timeout=FETCH_TIMEOUT, headers=DEFAULT_HEADERS) + response.raise_for_status() + return parse_match_live_detail_html(response.text) + + +class TaskFileLock: + def __init__(self, name: str, lock_dir: Path | None = None): + safe_name = "".join(ch if ch.isalnum() or ch in ("-", "_") else "_" for ch in str(name or "task")) + 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_name}.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) -> None: + 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 + + +@dataclass +class MatchLiveStreamRepository: + db_config: Dict[str, Any] + + def __post_init__(self) -> None: + self.prefix = str(self.db_config.get("prefix") or "la_") + self.conn = pymysql.connect( + host=self.db_config["host"], + port=int(self.db_config["port"]), + user=self.db_config["user"], + password=self.db_config["password"], + database=self.db_config["database"], + charset=self.db_config.get("charset") or "utf8mb4", + cursorclass=pymysql.cursors.DictCursor, + autocommit=False, + ) + + def close(self) -> None: + self.conn.close() + + def load_due_rows(self, now_ts: int, limit: int = FETCH_BATCH_SIZE) -> List[Dict[str, Any]]: + sql = ( + f"SELECT id, match_id, title, source_url, play_url_flv, play_url_hd_flv, play_url_m3u8, play_url_hd_m3u8 " + f"FROM `{self.prefix}match_live` " + f"WHERE `delete_time` IS NULL AND `next_fetch_at` <= %s " + f"ORDER BY `next_fetch_at` ASC, `id` ASC LIMIT %s" + ) + with self.conn.cursor() as cur: + cur.execute(sql, (int(now_ts), int(limit))) + return list(cur.fetchall() or []) + + def mark_success(self, row_id: int, stream_urls: Dict[str, str], now_ts: int) -> None: + sql = ( + f"UPDATE `{self.prefix}match_live` SET " + "`play_url_flv`=%s, `play_url_hd_flv`=%s, `play_url_m3u8`=%s, `play_url_hd_m3u8`=%s, " + "`fetch_status`=1, `fetch_error`='', `last_fetch_at`=%s, `next_fetch_at`=%s, `update_time`=%s " + "WHERE `id`=%s" + ) + next_fetch_at = choose_next_fetch_at(now_ts, has_existing_stream=True) + with self.conn.cursor() as cur: + cur.execute( + sql, + ( + stream_urls["play_url_flv"], + stream_urls["play_url_hd_flv"], + stream_urls["play_url_m3u8"], + stream_urls["play_url_hd_m3u8"], + int(now_ts), + next_fetch_at, + int(now_ts), + int(row_id), + ), + ) + self.conn.commit() + + def mark_failed(self, row_id: int, error_message: str, now_ts: int, has_existing_stream: bool) -> None: + sql = ( + f"UPDATE `{self.prefix}match_live` SET " + "`fetch_status`=2, `fetch_error`=%s, `last_fetch_at`=%s, `next_fetch_at`=%s, `update_time`=%s " + "WHERE `id`=%s" + ) + with self.conn.cursor() as cur: + cur.execute( + sql, + ( + _trim_error(error_message), + int(now_ts), + choose_next_fetch_at(now_ts, has_existing_stream=has_existing_stream), + int(now_ts), + int(row_id), + ), + ) + self.conn.commit() + + +def run(db_config: Dict[str, Any]) -> Dict[str, Any]: + lock = TaskFileLock("match_live_stream") + if not lock.acquire(): + return { + "success": True, + "candidate_count": 0, + "saved_count": 0, + "summary_text": "match_live_stream 已在执行中,跳过本次", + } + + repo = MatchLiveStreamRepository(db_config) + now_ts = int(time.time()) + success_count = 0 + failed_count = 0 + rows: List[Dict[str, Any]] = [] + + try: + rows = repo.load_due_rows(now_ts, limit=FETCH_BATCH_SIZE) + for row in rows: + fetch_ts = int(time.time()) + has_existing_stream = _has_existing_stream(row) + try: + stream_urls = _fetch_stream_payload(str(row.get("source_url") or ""), fetch_ts) + repo.mark_success(_safe_int(row.get("id")), stream_urls, fetch_ts) + success_count += 1 + except Exception as exc: # noqa: BLE001 + repo.mark_failed(_safe_int(row.get("id")), str(exc), fetch_ts, has_existing_stream) + failed_count += 1 + + candidate_count = len(rows) + return { + "success": True, + "candidate_count": candidate_count, + "saved_count": success_count, + "count": success_count, + "failed_count": failed_count, + "summary_text": ( + f"赛事直播流抓取完成: 待处理 {candidate_count} 条, " + f"成功 {success_count} 条, 失败 {failed_count} 条" + ), + } + finally: + repo.close() + lock.release() diff --git a/docker/crawler/tests/test_match_live_stream.py b/docker/crawler/tests/test_match_live_stream.py new file mode 100644 index 0000000..a7bca80 --- /dev/null +++ b/docker/crawler/tests/test_match_live_stream.py @@ -0,0 +1,44 @@ +import json +import sys +import unittest +from pathlib import Path + + +ROOT = Path(__file__).resolve().parents[1] +if str(ROOT) not in sys.path: + sys.path.insert(0, str(ROOT)) + +from match_live_stream import choose_next_fetch_at, parse_match_live_detail_html + + +class MatchLiveStreamTest(unittest.TestCase): + def test_parse_json_ncctrials_detail_extracts_four_urls(self): + payload = { + "code": 200, + "msg": "ok", + "data": { + "stream": { + "flv": "http://a/flv", + "hdFlv": "http://a/hdflv", + "m3u8": "http://a/m3u8", + "hdM3u8": "http://a/hdm3u8", + } + }, + } + + parsed = parse_match_live_detail_html(f"detail({json.dumps(payload, ensure_ascii=False)})") + + self.assertEqual(parsed["play_url_m3u8"], "http://a/m3u8") + self.assertEqual(parsed["play_url_hd_m3u8"], "http://a/hdm3u8") + self.assertEqual(parsed["play_url_flv"], "http://a/flv") + self.assertEqual(parsed["play_url_hd_flv"], "http://a/hdflv") + + def test_choose_next_fetch_at_uses_60s_for_empty_streams(self): + self.assertEqual(choose_next_fetch_at(now_ts=1000, has_existing_stream=False), 1060) + + def test_choose_next_fetch_at_uses_600s_for_existing_streams(self): + self.assertEqual(choose_next_fetch_at(now_ts=1000, has_existing_stream=True), 1600) + + +if __name__ == "__main__": + unittest.main() diff --git a/docs/sql/insert_crontab_match_live_stream.sql b/docs/sql/insert_crontab_match_live_stream.sql new file mode 100644 index 0000000..8311b69 --- /dev/null +++ b/docs/sql/insert_crontab_match_live_stream.sql @@ -0,0 +1,20 @@ +INSERT INTO `la_dev_crontab` +(`name`, `type`, `system`, `remark`, `command`, `params`, `status`, `expression`, `create_time`, `update_time`) +SELECT + '赛事直播流抓取', + 1, + 0, + 'Docker crawler 每分钟抓取赛事直播线路播放地址', + 'crawler', + 'match_live_stream', + 1, + '* * * * *', + UNIX_TIMESTAMP(), + UNIX_TIMESTAMP() +WHERE NOT EXISTS ( + SELECT 1 + FROM `la_dev_crontab` + WHERE `command` = 'crawler' + AND `params` = 'match_live_stream' + AND `delete_time` IS NULL +);