从 CCTV 主页抓取《新闻联播》,下载 → 转 MP3/WAV → 静音切分 → ASR 识别 → LLM 校对 → 切分为单条新闻 → 入库 MySQL。 主要内容: - 全链路:getVideo5 抓取下载、audioRead 转写、deepseek 校对与切分、newsProcess 入库 - 可移植化:配置分层,.env 只放密钥、config.yml 放模型/接入点/路由/参数 - 可换供应商:endpoints(kind/base_url/api_key_env/extra_body)+ routes 按环节选路 - 数据保真:数值事实守卫,校对改动数字/年份/届次则整片回退 ASR 原文; 识别不完整不发布该日精编,避免半天内容被当成完整一天 - 定时任务:systemd 每天 21:00,失败 21:30 / 22:00 重试; 只缺切分时只重跑切分(省掉全部 ASR),用 state/.asr_complete_* 标记判定阶段 - 隧道自愈:13306 不通时自动执行 autossh.sh(所有入口共用,systemd 托管时只等待) - 中间产物每日清理;97 项离线自检(配置/清理/事实守卫/解析/隧道)
271 lines
11 KiB
Python
271 lines
11 KiB
Python
"""
|
||
xwlb_daily 表结构如下:
|
||
+--------------+---------+------+-----+---------+----------------+
|
||
| Field | Type | Null | Key | Default | Extra |
|
||
+--------------+---------+------+-----+---------+----------------+
|
||
| nid | int(11) | NO | PRI | NULL | auto_increment |
|
||
| news_days | date | NO | | NULL | |
|
||
| daily_sub_id | int(11) | NO | | NULL | |
|
||
| news_raw | text | NO | | NULL | |
|
||
| news_improve | text | NO | | NULL | |
|
||
| news_title | text | NO | | NULL | |
|
||
+--------------+---------+------+-----+---------+----------------+
|
||
xwlb_daily_ext 表结构如下:
|
||
+--------------+--------------+------+-----+---------+----------------+
|
||
| Field | Type | Null | Key | Default | Extra |
|
||
+--------------+--------------+------+-----+---------+----------------+
|
||
| extid | int(11) | NO | PRI | NULL | auto_increment |
|
||
| news_date | date | NO | | NULL | |
|
||
| sub_id | tinyint(4) | NO | | NULL | |
|
||
| news_title | varchar(256) | NO | | NULL | |
|
||
| news_content | text | NO | | NULL | |
|
||
+--------------+--------------+------+-----+---------+----------------+
|
||
|
||
流程:取给定日期的所有 news_improve 内容,按 daily_sub_id 顺序拼接为一个字符串,
|
||
交给 DeepSeek 切分为独立新闻并起标题,写入 xwlb_daily_ext。
|
||
|
||
本次修复(详见 docs/BUGS.md):
|
||
- B3:LLM 返回的 ```json 围栏 / 截断 / 非 JSON 内容按「剥围栏 → 宽解析 → 重试 → 明确报错」处理,
|
||
不再静默丢弃当天全部文本;
|
||
- B5:写入前先删除该日期的旧记录,重跑不产生重复;
|
||
- B13:是否「已处理」改为按 xwlb_daily_ext 是否存在记录判断(原 `COUNT(*)>5` 会把
|
||
只切出 ≤5 条的日期永久判为未处理)。
|
||
"""
|
||
import json
|
||
import logging
|
||
import re
|
||
|
||
from mysqlHandle import MySQLDB
|
||
from deepseek import deepseek_text
|
||
import config
|
||
|
||
logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')
|
||
logger = logging.getLogger(__name__)
|
||
|
||
SPLIT_PROMPT = ("###请根据下面新闻内容的文本逻辑 \n"
|
||
" - 帮我分割成各个独立的新闻内容(注意:不要修改新闻本身,仅分割文本),并给每个新闻总结一个标题; \n"
|
||
" - 如果遇到'国内快讯'、'国际快讯'或'联播快讯',也请根据每个条快讯分割为一个新闻以及新闻标题; \n"
|
||
" - 返回json格式。json格式包含:news_id,news_title,news_content; news_id从1开始递增。")
|
||
|
||
# 首次解析失败后的强化提示词(进一步约束输出格式)
|
||
RETRY_PROMPT = SPLIT_PROMPT + ("\n\n重要:只输出一个 JSON 数组,不要输出 markdown 代码块、不要输出任何解释文字。"
|
||
"数组每一项形如 {\"news_id\": 1, \"news_title\": \"标题\", \"news_content\": \"正文\"}。")
|
||
|
||
# news_to_db 的执行结果
|
||
STATUS_WRITTEN = 'written' # 本次完成切分并写入
|
||
STATUS_SKIPPED = 'skipped' # 已有精编记录,按设计跳过
|
||
STATUS_FAILED = 'failed' # 无可用文本 / 切分失败 / 写入不完整
|
||
|
||
|
||
def normalize_date(date_str):
|
||
"""把 YYYYMMDD / YYYY-MM-DD 统一为 YYYY-MM-DD"""
|
||
s = str(date_str or '').strip()
|
||
if re.fullmatch(r'\d{8}', s):
|
||
return f"{s[:4]}-{s[4:6]}-{s[6:8]}"
|
||
if re.fullmatch(r'\d{4}-\d{2}-\d{2}', s):
|
||
return s
|
||
raise ValueError(f"无法识别的日期格式: {date_str!r}(应为 YYYYMMDD 或 YYYY-MM-DD)")
|
||
|
||
|
||
def _get_ext_count(db, target_date):
|
||
try:
|
||
row = db.query_one("xwlb_daily_ext", "COUNT(*) as count", "news_date = %s", (target_date,))
|
||
return int(row['count']) if row else 0
|
||
except Exception as e:
|
||
logger.error(f"查询 xwlb_daily_ext 失败: {e}")
|
||
return 0
|
||
|
||
|
||
def get_news_improve_by_date(target_date):
|
||
"""
|
||
获取指定日期的所有 news_improve 内容,按 daily_sub_id 顺序拼接
|
||
|
||
参数:
|
||
target_date: 目标日期,格式 YYYY-MM-DD
|
||
返回值:
|
||
str | None: 拼接后的文本;该日期没有任何有效文本时返回 None
|
||
"""
|
||
try:
|
||
db = MySQLDB()
|
||
try:
|
||
result = db.query_data(
|
||
table="xwlb_daily",
|
||
columns="news_improve",
|
||
where="news_days = %s order by daily_sub_id ASC",
|
||
params=(target_date,))
|
||
finally:
|
||
db.close()
|
||
|
||
if not result:
|
||
return None
|
||
# 过滤空行:历史数据中存在识别失败留下的空白记录(B1),不能把它们拼进正文
|
||
parts = [row['news_improve'].strip() for row in result
|
||
if row.get('news_improve') and row['news_improve'].strip()]
|
||
if not parts:
|
||
return None
|
||
return '\n'.join(parts)
|
||
except Exception as e:
|
||
logger.error(f"查询失败: {e}")
|
||
return None
|
||
|
||
|
||
# --------------------------------------------------------------------------- 解析
|
||
|
||
def _strip_code_fence(text):
|
||
"""去掉 ```json ... ``` 包裹(B3 中最常见的失败形态)"""
|
||
s = (text or '').strip()
|
||
if s.startswith('```'):
|
||
s = re.sub(r'^```[a-zA-Z0-9_-]*\s*', '', s)
|
||
s = re.sub(r'\s*```\s*$', '', s)
|
||
return s.strip()
|
||
|
||
|
||
def _loads_lenient(text):
|
||
"""先整体解析;失败则从第一个 { 或 [ 处做 raw_decode(容忍前后多余文字)"""
|
||
try:
|
||
return json.loads(text)
|
||
except json.JSONDecodeError:
|
||
match = re.search(r'[\[{]', text)
|
||
if not match:
|
||
raise
|
||
obj, _ = json.JSONDecoder().raw_decode(text[match.start():])
|
||
return obj
|
||
|
||
|
||
def _to_news_list(data):
|
||
"""把 LLM 返回的多种形态统一成列表"""
|
||
if isinstance(data, list):
|
||
return data
|
||
if isinstance(data, dict):
|
||
for value in data.values():
|
||
if isinstance(value, list):
|
||
return value
|
||
values = list(data.values())
|
||
# 兼容 {"1": {...}, "2": {...}};但不能把单条新闻的字段值当成列表(原实现踩过这个坑)
|
||
if values and all(isinstance(v, dict) for v in values):
|
||
return values
|
||
raise ValueError("JSON 对象中未找到新闻列表")
|
||
raise ValueError(f"不支持的 DeepSeek 响应格式: {type(data).__name__}")
|
||
|
||
|
||
def extract_news_rows(response_text, target_date):
|
||
"""
|
||
解析 DeepSeek 返回并规整为待写入的行
|
||
|
||
返回值:
|
||
list[dict]: 可直接交给 insert_many 的行
|
||
异常:
|
||
ValueError: 无法解析或没有任何有效新闻
|
||
"""
|
||
if not response_text or not str(response_text).strip():
|
||
raise ValueError("LLM 响应为空")
|
||
|
||
data = _loads_lenient(_strip_code_fence(str(response_text)))
|
||
items = _to_news_list(data)
|
||
|
||
rows = []
|
||
for idx, item in enumerate(items, 1):
|
||
if not isinstance(item, dict):
|
||
logger.warning(f"跳过非对象条目: {str(item)[:60]}")
|
||
continue
|
||
content = str(item.get('news_content') or item.get('content') or '').strip()
|
||
if not content:
|
||
continue
|
||
title = str(item.get('news_title') or item.get('title') or '').strip()
|
||
try:
|
||
sub_id = int(item.get('news_id', idx))
|
||
except (TypeError, ValueError):
|
||
sub_id = idx
|
||
rows.append({
|
||
"news_date": target_date,
|
||
"sub_id": sub_id,
|
||
"news_title": title[:256], # varchar(256)
|
||
"news_content": content,
|
||
})
|
||
|
||
if not rows:
|
||
raise ValueError("解析结果中没有任何有效新闻")
|
||
return rows
|
||
|
||
|
||
# --------------------------------------------------------------------------- 入库
|
||
|
||
def news_to_db(target_date, force=False):
|
||
"""
|
||
取当天原文 → DeepSeek 切分 → 覆盖写入 xwlb_daily_ext
|
||
|
||
参数:
|
||
target_date: 日期(YYYYMMDD 或 YYYY-MM-DD)
|
||
force: 即使已有精编记录也重新切分(默认 False,避免重复消耗 token)
|
||
返回值:
|
||
str: STATUS_WRITTEN / STATUS_SKIPPED / STATUS_FAILED
|
||
(注意:不是 bool;'skipped' 与 'written' 都表示数据已就绪)
|
||
"""
|
||
target_date = normalize_date(target_date)
|
||
|
||
if not force:
|
||
db = MySQLDB()
|
||
try:
|
||
existing = _get_ext_count(db, target_date)
|
||
finally:
|
||
db.close()
|
||
if existing > 0:
|
||
logger.info(f"日期 {target_date} 已有 {existing} 条精编记录,跳过(需重跑请加 force)")
|
||
return STATUS_SKIPPED
|
||
|
||
result = get_news_improve_by_date(target_date)
|
||
if result is None:
|
||
logger.warning(f"日期 {target_date} 没有新闻内容,跳过")
|
||
return STATUS_FAILED
|
||
logger.info(f"日期 {target_date} 的新闻内容长度:{len(result)} 字符")
|
||
|
||
attempts = [(SPLIT_PROMPT, {}),
|
||
(RETRY_PROMPT, {"max_tokens": config.get_int('llm_split.retry_max_tokens', 32000)})]
|
||
rows = None
|
||
last_error = None
|
||
for attempt, (prompt, extra) in enumerate(attempts, 1):
|
||
response = None
|
||
try:
|
||
# use_fallback=False:切分环节自带重试与解析校验,
|
||
# 关掉"返回降级文本"的兜底,模型名/权限等配置错误会直接暴露而不是静默出数据
|
||
response = deepseek_text(result, prompt, use_fallback=False, **extra)
|
||
rows = extract_news_rows(response, target_date)
|
||
logger.info(f"第 {attempt} 次切分成功,得到 {len(rows)} 条新闻")
|
||
break
|
||
except Exception as e:
|
||
last_error = e
|
||
logger.error(f"第 {attempt} 次切分/解析失败: {e}")
|
||
if response is not None:
|
||
logger.error(f"DeepSeek 原始返回(前 300 字): {str(response)[:300]}")
|
||
rows = None
|
||
|
||
if not rows:
|
||
# 不写库、不删除已有数据,交给下次重跑
|
||
logger.error(f"❌ 日期 {target_date} 切分失败,未写入数据库: {last_error}")
|
||
return STATUS_FAILED
|
||
|
||
db = MySQLDB()
|
||
try:
|
||
deleted = db.execute("DELETE FROM xwlb_daily_ext WHERE news_date = %s", (target_date,))
|
||
written = db.insert_many("xwlb_daily_ext", rows)
|
||
finally:
|
||
db.close()
|
||
|
||
if written != len(rows):
|
||
logger.error(f"❌ 日期 {target_date} 写入不完整:期望 {len(rows)} 条,实际 {written} 条")
|
||
return STATUS_FAILED
|
||
logger.info(f"✓ 日期 {target_date} 写入 {written} 条新闻(替换旧记录 {deleted} 条)")
|
||
return STATUS_WRITTEN
|
||
|
||
|
||
if __name__ == "__main__":
|
||
import sys
|
||
from datetime import datetime
|
||
|
||
date_arg = sys.argv[1] if len(sys.argv) > 1 else datetime.now().strftime('%Y-%m-%d')
|
||
force_flag = '--force' in sys.argv
|
||
config.log_summary()
|
||
status = news_to_db(date_arg, force=force_flag)
|
||
label = {'written': '成功写入', 'skipped': '已有精编记录,已跳过', 'failed': '失败'}.get(status, status)
|
||
print(f"处理结果: {label}({status})")
|
||
raise SystemExit(0 if status != STATUS_FAILED else 1) |