Files
xwlb/deepseek.py
simon 60f8c263f2 《新闻联播》每日抓取入库:全链路 + 可移植化 + 定时任务
从 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 项离线自检(配置/清理/事实守卫/解析/隧道)
2026-09-25 11:17:46 +08:00

365 lines
16 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import env # 加载 .env(敏感项)
import config # 加载 config.yml(模型、参数)
import requests
import json
import time
import logging
from typing import Optional, Dict, Any
import os
# 配置日志
logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')
logger = logging.getLogger(__name__)
class DeepSeekAPI:
"""OpenAI 兼容协议的对话客户端(DeepSeek 是默认接入点)
地址与密钥变量名都来自 config.yml 的 `endpoints.<routes.<role>>`,
因此把 routes.split / routes.correct 指向别的 OpenAI 兼容供应商即可切换,无需改代码。
请求地址 = base_url + chat_completions_path。
"""
def __init__(self, api_key: Optional[str] = None, role: str = 'split',
endpoint: Optional[str] = None):
"""
初始化客户端
Args:
api_key: 直接指定密钥;为 None 时按接入点的 api_key_env 从 .env 读
role: 环节名(asr/correct/split),用于解析 routes.<role>
endpoint: 直接指定接入点名(覆盖 role 的路由)
"""
if endpoint:
cfg = config.endpoint(endpoint)
if not cfg:
raise RuntimeError(f"config.yml 的 endpoints 下没有接入点 '{endpoint}'")
name = endpoint
else:
name, cfg = config.endpoint_for(role)
endpoint_extra = cfg.get('extra_body')
self.endpoint_extra_body = dict(endpoint_extra) if isinstance(endpoint_extra, dict) else {}
if endpoint_extra not in (None, {}) and not isinstance(endpoint_extra, dict):
logger.warning("接入点 %s 的 extra_body 应为字典,已忽略: %r", name, endpoint_extra)
self.endpoint_name = name
self.role = role
self.kind = str(cfg.get('kind', 'openai')).lower()
if self.kind != 'openai':
raise RuntimeError(
f"接入点 '{name}' 的 kind='{self.kind}',OpenAI 兼容客户端只能用于 kind=openai 的接入点")
base = str(cfg.get('base_url', '')).rstrip('/')
path = str(cfg.get('chat_completions_path', '/chat/completions'))
self.api_url = base + (path if path.startswith('/') else '/' + path)
self.api_key = api_key or os.getenv(cfg.get('api_key_env', ''), '')
if not self.api_key:
logger.warning("接入点 %s 的密钥未设置(.env 中的 %s)", name, cfg.get('api_key_env'))
# 调用参数按环节取各自配置段
section = 'llm_correct' if role == 'correct' else 'llm_split'
self.section = section
self.max_retries = config.get_int(f'{section}.max_retries', 3)
self.retry_delay = 2 # 秒
self.timeout = config.get_int(f'{section}.timeout', 60)
# 默认系统提示词
self.default_system_prompt = """你是一个专业的AI助手,能够准确理解用户需求并提供高质量的回答。
请根据用户的输入进行适当的处理和分析,保持回答的专业性和准确性。注意:所处理文字来自中央电视台新闻联播节目转文字,请在内容审查时重点考虑。"""
def _handle_api_error(self, response: requests.Response) -> str:
"""
处理API错误响应
Args:
response: API响应对象
Returns:
错误描述信息
"""
error_msg = (f"API请求失败 [{self.endpoint_name} {self.api_url}]: "
f"{response.status_code} {response.reason}")
try:
error_data = response.json()
if 'error' in error_data:
error_msg += f" - {error_data['error'].get('message', '未知错误')}"
logger.error(f"API错误详情: {error_data}")
except json.JSONDecodeError:
error_msg += f" - 响应内容: {response.text[:200]}"
return error_msg
def _make_api_request(self, payload: Dict[str, Any]) -> Dict[str, Any]:
"""
发送API请求并处理响应
Args:
payload: 请求数据
Returns:
API响应数据
Raises:
Exception: 当所有重试都失败时抛出异常
"""
headers = {
"Content-Type": "application/json",
"Authorization": f"Bearer {self.api_key}"
}
last_exception = None
for attempt in range(self.max_retries):
try:
logger.info("发送API请求 [%s] (尝试 %d/%d)", self.endpoint_name,
attempt + 1, self.max_retries)
response = requests.post(
self.api_url,
headers=headers,
json=payload,
timeout=self.timeout # 秒(config.yml: <section>.timeout)
)
if response.status_code == 200:
return response.json()
elif response.status_code == 400:
# 400错误通常是请求格式问题,不需要重试
error_msg = self._handle_api_error(response)
raise Exception(f"请求参数错误: {error_msg}")
elif response.status_code == 401:
# 401未授权错误,不需要重试
raise Exception("API密钥无效或未授权,请检查您的API密钥")
elif response.status_code == 429:
# 速率限制,需要重试
logger.warning("达到速率限制,等待后重试...")
time.sleep(self.retry_delay * (attempt + 1))
continue
elif 500 <= response.status_code < 600:
# 服务器错误,需要重试
logger.warning(f"服务器错误 {response.status_code},等待后重试...")
time.sleep(self.retry_delay * (attempt + 1))
continue
else:
error_msg = self._handle_api_error(response)
raise Exception(f"API请求失败: {error_msg}")
except requests.exceptions.Timeout:
last_exception = Exception(f"请求超时 (尝试 {attempt + 1})")
logger.warning(f"请求超时,等待后重试...")
time.sleep(self.retry_delay * (attempt + 1))
except requests.exceptions.ConnectionError:
last_exception = Exception(f"网络连接错误 (尝试 {attempt + 1})")
logger.warning(f"网络连接错误,等待后重试...")
time.sleep(self.retry_delay * (attempt + 1))
except requests.exceptions.RequestException as e:
last_exception = Exception(f"请求异常: {str(e)}")
logger.warning(f"请求异常,等待后重试...")
time.sleep(self.retry_delay * (attempt + 1))
# 所有重试都失败
if last_exception:
raise last_exception
else:
raise Exception("API请求失败,未知错误")
def process_text(self,
prompt: str,
text: str,
system_prompt: Optional[str] = None,
model: str = "deepseek-chat",
temperature: float = 0.7,
max_tokens: int = 2000,
response_format: Optional[Dict] = None,
extra_body: Optional[Dict] = None) -> str:
"""
处理文本的通用方法
Args:
prompt: 用户提示词
text: 需要处理的文本(约1万字符)
system_prompt: 系统提示词,如果为None则使用默认值
model: 使用的模型
temperature: 生成温度
max_tokens: 最大生成token数
Returns:
处理后的文本
Raises:
Exception: 当处理失败时抛出包含详细信息的异常
"""
# 输入验证
if not self.api_key:
raise Exception(
f"接入点 {self.endpoint_name} 的密钥未设置:请在 .env 中配置 "
f"{config.endpoint_for(self.role)[1].get('api_key_env') if self.role in ('asr', 'correct', 'split') else '对应变量'}")
if not prompt or not text:
raise Exception("prompt和text不能为空")
# 检查文本长度(约1万字符)
if len(text) > 15000: # 留一些余量
logger.warning(f"输入文本长度({len(text)}字符)较长,可能会超过上下文限制")
# 准备系统提示词
system_content = system_prompt or self.default_system_prompt
# 构建消息
messages = [
{"role": "system", "content": system_content},
{"role": "user", "content": f"{prompt}\n\n文本内容:\n{text}"}
]
# 构建请求数据
payload = {
"model": model,
"messages": messages,
"temperature": temperature,
"max_tokens": max_tokens,
"stream": False
}
if response_format:
payload["response_format"] = response_format
# 供应商特有参数原样透传,用于关闭推理模型的思维链等。
# 先应用接入点级(语法属于供应商:通义 {enable_thinking: false},
# DeepSeek {thinking: {type: disabled}},写错会被静默忽略),再用环节级覆盖。
if self.endpoint_extra_body:
payload.update(self.endpoint_extra_body)
if extra_body:
payload.update(extra_body)
try:
# 发送API请求
response_data = self._make_api_request(payload)
# 解析响应
if 'choices' in response_data and len(response_data['choices']) > 0:
result = response_data['choices'][0]['message']['content']
logger.info("文本处理成功完成")
return result.strip()
else:
raise Exception("API响应格式异常,未找到有效结果")
except Exception as e:
logger.error(f"文本处理失败: {str(e)}")
raise Exception(f"文本处理失败: {str(e)}")
def process_text_with_fallback(self,
prompt: str,
text: str,
system_prompt: Optional[str] = None,
**kwargs) -> str:
"""
带降级处理的文本处理方法
Args:
prompt: 用户提示词
text: 需要处理的文本
system_prompt: 系统提示词
**kwargs: 其他参数
Returns:
处理后的文本,如果API调用失败则返回降级结果
"""
try:
return self.process_text(prompt, text, system_prompt, **kwargs)
except Exception as e:
logger.error(f"API调用失败,使用降级处理: {str(e)}")
# 这里可以添加降级逻辑,比如返回原始文本或简单处理
return f"处理失败,返回原始文本(错误: {str(e)})\n\n{text}"
# 使用示例
def deepseek_text(text, prompt, model=None, max_tokens=None, temperature=None, use_fallback=True):
"""
调用 DeepSeek 处理文本(默认 json_object 模式)
参数:
text: 待处理文本
prompt: 提示词
model: 模型名,默认取 config.yml 的 models.split_model
max_tokens: 最大生成 token 数,默认取 llm_split.max_tokens
temperature: 采样温度,默认取 llm_split.temperature
use_fallback: 调用失败时是否返回降级文本(False 时直接抛异常,便于调用方重试)
返回值:
str: 模型返回内容
"""
api_client = DeepSeekAPI(role='split') # 接入点/密钥来自 config.yml 的 routes.split
custom_system_prompt = "你是一个专业的文本分析助手,擅长根据提示词对长文本进行深入分析。"
try:
# 处理文本(使用 response_format 强制返回 JSON)
result = api_client.process_text(
model=model or config.model('split_model'),
prompt=prompt,
text=text,
system_prompt=custom_system_prompt,
temperature=config.get('llm_split.temperature', 0.5) if temperature is None else temperature,
max_tokens=config.get_int('llm_split.max_tokens', 20000) if max_tokens is None else max_tokens,
response_format={"type": "json_object"},
extra_body=config.get_dict('llm_split.extra_body'),
)
return result
except Exception as e:
logger.error(f"DeepSeek 处理失败: {e}")
if not use_fallback:
raise
# 使用降级方法(显式传入同一模型与参数,避免兜底时**悄悄换用另一个模型**——
# 曾出现配置的模型名无效、兜底用硬编码 deepseek-chat 成功、
# 于是"按配置运行"看起来正常实则完全没生效)
fallback_result = api_client.process_text_with_fallback(
prompt=prompt,
text=text,
system_prompt=custom_system_prompt,
model=model or config.model('split_model'),
temperature=config.get('llm_split.temperature', 0.5) if temperature is None else temperature,
max_tokens=config.get_int('llm_split.max_tokens', 20000) if max_tokens is None else max_tokens,
extra_body=config.get_dict('llm_split.extra_body'),
)
# 原实现此处误写为 return result(异常分支下 result 未绑定,会抛 NameError,
# 见 docs/BUGS.md B7)
return fallback_result
def openai_chat(text, prompt, role='correct', model=None, max_tokens=None, temperature=None,
timeout=None, system_prompt=None):
"""
通用 OpenAI 兼容对话调用(供 **非 DashScope 供应商的文本校对** 使用)
地址 = config.yml `endpoints.<routes.<role>>`,即换供应商只改配置。
与 `deepseek_text` 的区别:不强制 json_object,直接返回纯文本。
参数:
text/prompt: 待处理文本与提示词
role: 环节名,默认 correct(决定用哪个接入点与哪段参数)
model: 模型名,默认取 models.correct_model(role=split 时取 models.split_model)
max_tokens/temperature/timeout: 不传则取该环节配置段
返回值:
str: 模型返回的文本
"""
client = DeepSeekAPI(role=role)
section = client.section
if timeout is not None:
client.timeout = timeout
default_model = config.model('split_model') if role == 'split' else config.model('correct_model')
return client.process_text(
model=model or default_model,
prompt=prompt,
text=text,
system_prompt=system_prompt,
temperature=config.get(f'{section}.temperature', 0.3) if temperature is None else temperature,
max_tokens=config.get_int(f'{section}.max_tokens', 8000) if max_tokens is None else max_tokens,
extra_body=config.get_dict(f'{section}.extra_body'),
)
if __name__ == "__main__":
import sys
if len(sys.argv) < 3:
print("用法: python deepseek.py <prompt> <text>")
sys.exit(1)
print(deepseek_text(sys.argv[2], sys.argv[1]))