从 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 项离线自检(配置/清理/事实守卫/解析/隧道)
684 lines
28 KiB
Python
684 lines
28 KiB
Python
import datetime
|
||
import env # 加载 .env(敏感项)
|
||
import config # 加载 config.yml(模型、参数、路径、开关)
|
||
import os
|
||
import re
|
||
import time
|
||
import dashscope
|
||
import pydub
|
||
from pydub import AudioSegment
|
||
from pydub.silence import split_on_silence
|
||
from dashscope.audio.asr import Recognition
|
||
from dashscope import Generation
|
||
from http import HTTPStatus
|
||
from mysqlHandle import MySQLDB
|
||
from deepseek import openai_chat
|
||
import logging
|
||
|
||
logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# 确保 dashscope SDK 使用 config.yml 里配置的接入点地址
|
||
# (SDK 在 import 时读环境变量,config 加载阶段已写入;这里再显式覆盖一次,不依赖 import 顺序)
|
||
config.apply_dashscope_endpoints(dashscope)
|
||
# 设置环境变量
|
||
# os.environ["DASHSCOPE_API_KEY"] = "sk-your-dashscope-key"
|
||
|
||
def convert_mp3_to_wav(mp3_path, output_wav_path):
|
||
"""
|
||
将MP3文件转换为16kHz单声道WAV格式,这是Qwen3-ASR-Flash模型的推荐格式
|
||
参数:
|
||
mp3_path (str): MP3文件路径
|
||
output_wav_path (str): 输出WAV文件路径
|
||
返回值:
|
||
str: 转换后的WAV文件路径
|
||
"""
|
||
logger.info(f"开始转换MP3到WAV: {mp3_path}")
|
||
# 加载MP3文件
|
||
audio = AudioSegment.from_file(mp3_path, format="mp3")
|
||
# 转换为 ASR 要求的采样率(config.yml: asr.sample_rate)、单声道
|
||
audio = audio.set_frame_rate(config.get_int('asr.sample_rate', 16000)).set_channels(1)
|
||
# 导出为WAV格式
|
||
audio.export(output_wav_path, format="wav")
|
||
logger.info(f"✓ MP3转换完成: {output_wav_path}")
|
||
#print(f"✓ MP3转换完成: {output_wav_path}")
|
||
return output_wav_path
|
||
|
||
def split_audio_by_fixed_duration(audio_path, chunk_duration, output_folder):
|
||
"""
|
||
将音频文件按固定时长分割成多个片段
|
||
参数:
|
||
audio_path (str): 音频文件路径
|
||
chunk_duration (int): 分片时长(毫秒)
|
||
output_folder (str): 输出文件夹路径
|
||
返回值:
|
||
list: 分片文件路径列表
|
||
"""
|
||
# 加载音频文件
|
||
audio = AudioSegment.from_file(audio_path)
|
||
# 计算总时长(毫秒)
|
||
total_duration = len(audio)
|
||
# 分片数
|
||
num_chunks = total_duration // chunk_duration + 1
|
||
# 存储分片文件路径
|
||
chunks = []
|
||
|
||
# 创建输出文件夹
|
||
os.makedirs(output_folder, exist_ok=True)
|
||
|
||
logger.info(f"开始音频分割,总时长: {total_duration/1000:.1f}秒,将分割为{num_chunks}个片段")
|
||
|
||
for i in range(num_chunks):
|
||
# 计算当前分片的起始和结束时间
|
||
start_time = i * chunk_duration
|
||
end_time = (i + 1) * chunk_duration
|
||
# 提取分片音频
|
||
chunk = audio[start_time:end_time]
|
||
# 生成文件名
|
||
chunk_name = f"chunk_{i}.wav"
|
||
chunk_path = os.path.join(output_folder, chunk_name)
|
||
# 导出分片音频
|
||
chunk.export(chunk_path, format="wav")
|
||
chunks.append(chunk_path)
|
||
|
||
# 打印处理进度
|
||
progress = (i + 1) / num_chunks * 100
|
||
logger.info(f"✓ 已完成分片 {i+1}/{num_chunks} ({progress:.1f}%)")
|
||
|
||
logger.info(f"✓ 音频分割完成,共生成{len(chunks)}个分片文件")
|
||
return chunks
|
||
|
||
def split_audio_by_smart_silence(audio_path, min_silence_len, silence_thresh, output_folder):
|
||
"""
|
||
将音频文件按智能静音检测方式分割成多个片段,每段不超过3分钟
|
||
参数:
|
||
audio_path (str): 音频文件路径
|
||
min_silence_len (int): 最小静音长度(毫秒)
|
||
silence_thresh (int): 静音阈值(dBFS)
|
||
output_folder (str): 输出文件夹路径
|
||
返回值:
|
||
list: 分片文件路径列表
|
||
"""
|
||
# 加载音频文件
|
||
audio = AudioSegment.from_file(audio_path, format="wav")
|
||
# 按静音分割
|
||
segments = split_on_silence(
|
||
audio,
|
||
# 静音超过该长度则分割(config.yml: audio_split.min_silence_ms)
|
||
min_silence_len=min_silence_len,
|
||
# 静音阈值(config.yml: audio_split.silence_thresh_db)
|
||
silence_thresh=silence_thresh,
|
||
# 保留静音部分(config.yml: audio_split.keep_silence_ms)
|
||
keep_silence=config.get_int('audio_split.keep_silence_ms', 400)
|
||
)
|
||
|
||
logger.info(f"✓ 静音分割完成,共{len(segments)}个初始片段")
|
||
|
||
# 合并过短的片段
|
||
merged_segments = []
|
||
current_segment = None
|
||
for segment in segments:
|
||
if current_segment is None:
|
||
current_segment = segment
|
||
else:
|
||
# 合并当前片段和新片段
|
||
temp_segment = current_segment + segment
|
||
# 如果合并后的片段超过上限则单独保存(config.yml: audio_split.max_chunk_ms)
|
||
if len(temp_segment) > config.get_int('audio_split.max_chunk_ms', 180000):
|
||
merged_segments.append(current_segment)
|
||
current_segment = segment
|
||
else:
|
||
current_segment = temp_segment
|
||
# 添加最后一个片段
|
||
if current_segment is not None:
|
||
merged_segments.append(current_segment)
|
||
|
||
logger.info(f"✓ 片段合并完成,共{len(merged_segments)}个最终片段")
|
||
|
||
# 存储分片文件路径
|
||
chunks = []
|
||
|
||
# 创建输出文件夹
|
||
os.makedirs(output_folder, exist_ok=True)
|
||
|
||
logger.info(f"开始导出音频片段到: {output_folder}")
|
||
|
||
for i, segment in enumerate(merged_segments):
|
||
# 生成文件名
|
||
chunk_name = f"chunk_{i}.wav"
|
||
chunk_path = os.path.join(output_folder, chunk_name)
|
||
# 导出分片音频
|
||
segment.export(chunk_path, format="wav")
|
||
chunks.append(chunk_path)
|
||
|
||
# 打印处理进度
|
||
progress = (i + 1) / len(merged_segments) * 100
|
||
logger.info(f"✓ 已完成分片 {i+1}/{len(merged_segments)} ({progress:.1f}%)")
|
||
|
||
logger.info(f"✓ 智能静音分割完成,共生成{len(chunks)}个分片文件")
|
||
return chunks
|
||
|
||
|
||
def transcribe_audio(audio_path):
|
||
"""
|
||
仅执行语音识别(ASR),**不再在同一函数里做 LLM 校对**
|
||
|
||
原实现把「识别」与「qwen 校对」串在同一个 try 中,校对超时会被当成识别失败,
|
||
导致已识别成功的文本被整体丢弃(生产库中已累积 101 行空记录,见 docs/BUGS.md B1)。
|
||
|
||
参数:
|
||
audio_path (str): 音频文件路径(必须是16kHz单声道WAV)
|
||
返回值:
|
||
str: 识别文本;失败返回 ''(由调用方决定是否跳过入库)
|
||
"""
|
||
try:
|
||
# 确保音频文件存在
|
||
if not os.path.exists(audio_path):
|
||
logger.error(f"音频文件不存在: {audio_path}")
|
||
return ''
|
||
dashscope.api_key = config.api_key('asr') # 密钥变量名由 config.yml 的 endpoints.*.api_key_env 决定
|
||
if not dashscope.api_key:
|
||
logger.error("接入点 %s 的密钥未配置(.env 中的 %s),无法识别音频",
|
||
config.route('asr'), config.api_key_env('asr'))
|
||
return ''
|
||
asr_endpoint, asr_cfg = config.endpoint_for('asr')
|
||
if str(asr_cfg.get('kind', '')).lower() != 'dashscope':
|
||
logger.error("ASR 接入点 '%s' 的 kind=%s:实时语音识别仅支持 kind=dashscope 的接入点,"
|
||
"更换 ASR 供应商需要新增适配器", asr_endpoint, asr_cfg.get('kind'))
|
||
return ''
|
||
# 创建识别对象
|
||
recognition = Recognition(
|
||
model=config.model('asr_model'), # config.yml: models.asr_model
|
||
format='wav',
|
||
sample_rate=config.get_int('asr.sample_rate', 16000),
|
||
language_hints=config.get_list('asr.language_hints', ['zh', 'en']), # 中文和英文
|
||
callback=None
|
||
)
|
||
|
||
# 调用识别
|
||
logger.info(f"开始识别音频: {audio_path}")
|
||
result = recognition.call(audio_path)
|
||
if result.status_code == HTTPStatus.OK:
|
||
# 提取识别结果
|
||
sentence = result.get_sentence()
|
||
text = merge_transcripts(sentence)
|
||
logger.info(f"✓ {audio_path} 识别成功,文本长度: {len(text)}")
|
||
return text
|
||
logger.error(f"❌ ASR 任务失败: {getattr(result, 'message', result)}")
|
||
return ''
|
||
except Exception as e:
|
||
logger.error(f"ASR 识别异常: {audio_path}: {e}")
|
||
return ''
|
||
|
||
|
||
def _extract_llm_text(response):
|
||
"""
|
||
兼容 dashscope 的两种返回形态,返回模型生成的纯文本(取不到则返回 '')
|
||
|
||
实测(dashscope 1.27.7 / qwen-plus / 约 700 字输入):status=200 时
|
||
`output.text` 为 None,真正的内容在 `output.choices[0].message.content`;
|
||
而短输入时 `output.text` 有值。旧代码只读 output.text,于是每次都拿到 None
|
||
(生产日志中「文本修正返回 None」共出现于每个分片,见 docs/BUGS.md B2)。
|
||
另外 401 等失败场景 `output` 为 None,需要一并防御。
|
||
"""
|
||
output = getattr(response, 'output', None)
|
||
if output is None:
|
||
return ''
|
||
text = getattr(output, 'text', None)
|
||
if text and text.strip():
|
||
return text.strip()
|
||
choices = getattr(output, 'choices', None) or []
|
||
if choices:
|
||
choice = choices[0]
|
||
message = choice.get('message') if isinstance(choice, dict) else getattr(choice, 'message', None)
|
||
if message:
|
||
content = message.get('content') if isinstance(message, dict) else getattr(message, 'content', None)
|
||
if content and content.strip():
|
||
return content.strip()
|
||
return ''
|
||
|
||
def merge_transcripts(transcripts):
|
||
"""
|
||
将多段识别文本合并成完整句子(保留原始段落逻辑,用空格连接)
|
||
参数:
|
||
transcripts (list): 识别结果列表,每个元素为字典{'text': '识别文本'}
|
||
返回:
|
||
str: 合并后的完整文本
|
||
"""
|
||
# 输入参数检查
|
||
if not transcripts:
|
||
return ""
|
||
|
||
# 确保transcripts是可迭代对象
|
||
if not hasattr(transcripts, '__iter__'):
|
||
return ""
|
||
|
||
try:
|
||
# 提取所有有效的text字段
|
||
texts = []
|
||
for t in transcripts:
|
||
try:
|
||
# 检查是否为字典类型且包含text字段
|
||
if isinstance(t, dict) and 'text' in t and t['text']:
|
||
text = t['text']
|
||
# 确保text是字符串类型
|
||
if isinstance(text, str) and text.strip():
|
||
texts.append(text.strip())
|
||
except (KeyError, TypeError, AttributeError):
|
||
# 忽略单个元素的处理错误,继续处理其他元素
|
||
continue
|
||
|
||
# 用空格连接所有段落(根据实际需求可调整连接符)
|
||
return " ".join(texts) if texts else ""
|
||
|
||
except Exception as e:
|
||
logger.error(f"合并转录文本时发生错误: {e}")
|
||
return ""
|
||
def text_correction(text, max_retries=None, timeout=None):
|
||
"""
|
||
使用通义千问模型修正文本中的错误和标点符号
|
||
|
||
与原实现的区别:
|
||
- 显式传入 timeout(原实现依赖 SDK 默认 300s,生产日志中 57 次读超时全发生在这里);
|
||
- 按结果取文本时兼容 `output.text` 与 `output.choices[0].message.content`(B2 根因);
|
||
- 失败时**抛异常**,由 analyze_and_correct_text 决定回退原文,不再静默返回 None;
|
||
- max_tokens 由 30000 降为可配置的 8000(校对输出不会超过输入量级)。
|
||
|
||
参数:
|
||
text (str): 需要修正的文本
|
||
max_retries (int): 重试次数,默认取环境变量 LLM_CORRECT_RETRIES=2
|
||
timeout (int): 单次调用超时秒数,默认取环境变量 LLM_CORRECT_TIMEOUT=90
|
||
返回值:
|
||
str: 修正后的文本
|
||
异常:
|
||
RuntimeError: 重试耗尽仍失败
|
||
"""
|
||
max_retries = max_retries if max_retries is not None else config.get_int('llm_correct.max_retries', 2)
|
||
timeout = timeout if timeout is not None else config.get_int('llm_correct.timeout', 90)
|
||
max_tokens = config.get_int('llm_correct.max_tokens', 8000)
|
||
model = config.model('correct_model') # config.yml: models.correct_model
|
||
endpoint_name, endpoint_cfg = config.endpoint_for('correct')
|
||
kind = str(endpoint_cfg.get('kind', 'dashscope')).lower()
|
||
|
||
logger.info("开始文本修正...")
|
||
|
||
# 构建修正提示词
|
||
correction_prompt = """请仔细检查以下文本,修正其中的错误:
|
||
1. 错别字和语法错误
|
||
2. 标点符号使用错误
|
||
3. 语句不通顺的地方
|
||
4. 逻辑不清晰的部分
|
||
5. 严禁改动任何事实信息:数字、年份、日期、届次、数量、机构名、人名、地名、专有名词(如"十五五")必须与原文完全一致,即使你认为原文有误也不要修改。
|
||
|
||
请直接返回修正后的完整文本,不要添加任何解释说明。"""
|
||
|
||
# 构建消息列表
|
||
messages = [
|
||
{"role": "system", "content": "你是一个专业的文本校对助手,擅长修正文本中的各种错误。"},
|
||
{"role": "user", "content": correction_prompt},
|
||
{"role": "user", "content": text}
|
||
]
|
||
|
||
last_error = "未知错误"
|
||
for attempt in range(1, max_retries + 1):
|
||
try:
|
||
logger.info(f"调用文本校对模型(接入点={endpoint_name} kind={kind} 第 {attempt}/{max_retries} 次,"
|
||
f"model={model},timeout={timeout}s)...")
|
||
if kind == 'openai':
|
||
# 非 DashScope 供应商:走 OpenAI 兼容的 /chat/completions
|
||
corrected = openai_chat(text, correction_prompt, role='correct',
|
||
model=model, max_tokens=max_tokens, timeout=timeout)
|
||
if corrected:
|
||
corrected = corrected.strip()
|
||
logger.info(f"✓ 文本修正完成,长度 {len(text)} -> {len(corrected)}")
|
||
return corrected
|
||
last_error = "返回内容为空"
|
||
logger.warning(f"⚠️ 文本修正返回空内容: {last_error}")
|
||
else:
|
||
# DashScope:用 SDK(dashscope.api_key 已按接入点设置)
|
||
dashscope.api_key = config.api_key('correct')
|
||
response = Generation.call(
|
||
model=model,
|
||
messages=messages,
|
||
max_tokens=max_tokens,
|
||
temperature=config.get('llm_correct.temperature', 0.1), # 较低温度以提高确定性
|
||
top_p=config.get('llm_correct.top_p', 0.5),
|
||
timeout=timeout,
|
||
)
|
||
if response.status_code != HTTPStatus.OK:
|
||
last_error = f"status={response.status_code} message={getattr(response, 'message', '')}"
|
||
logger.warning(f"❌ 文本修正API调用失败: {last_error}")
|
||
else:
|
||
corrected = _extract_llm_text(response)
|
||
if corrected:
|
||
logger.info(f"✓ 文本修正完成,长度 {len(text)} -> {len(corrected)}")
|
||
return corrected
|
||
last_error = "响应中无可用文本(output/choices 均为空)"
|
||
logger.warning(f"⚠️ 文本修正返回空内容: {last_error}")
|
||
except Exception as e:
|
||
last_error = f"{type(e).__name__}: {e}"
|
||
logger.warning(f"⚠️ 文本修正异常(第 {attempt}/{max_retries} 次): {last_error}")
|
||
|
||
if attempt < max_retries:
|
||
time.sleep(min(2 ** attempt, 8))
|
||
|
||
raise RuntimeError(f"文本修正失败(已重试 {max_retries} 次): {last_error}")
|
||
|
||
_CN_NUM_CHARS = '零〇一二两三四五六七八九十百千万亿'
|
||
_NUM_TOKEN_RE = re.compile(r'\d+(?:\.\d+)?|[零〇一二两三四五六七八九十百千万亿]+')
|
||
_CN_DIGITS = {'零': 0, '〇': 0, '一': 1, '二': 2, '两': 2, '三': 3, '四': 4,
|
||
'五': 5, '六': 6, '七': 7, '八': 8, '九': 9}
|
||
_CN_UNITS = {'十': 10, '百': 100, '千': 1000, '万': 10000, '亿': 100000000}
|
||
|
||
|
||
def _split_cn_numeral_run(run):
|
||
"""把中文数字串按「连续数字位」切成若干数词,处理"十五五/十四五"这类缩写
|
||
|
||
'十五五' -> ['十五', '五'];'十四五' -> ['十四', '五'];'一百二十三' -> ['一百二十三']
|
||
"""
|
||
segments, buf, prev_digit = [], '', False
|
||
for ch in run:
|
||
is_digit = ch in _CN_DIGITS
|
||
if is_digit and prev_digit and buf:
|
||
segments.append(buf)
|
||
buf = ''
|
||
buf += ch
|
||
prev_digit = is_digit
|
||
if buf:
|
||
segments.append(buf)
|
||
return segments
|
||
|
||
|
||
def _cn_to_number(text):
|
||
"""单个中文数词转数值(十/百/千/万/亿),无法确定时返回 None"""
|
||
total = section = number = 0
|
||
for ch in text:
|
||
if ch in _CN_DIGITS:
|
||
number = _CN_DIGITS[ch]
|
||
elif ch in _CN_UNITS:
|
||
unit = _CN_UNITS[ch]
|
||
if unit >= 10000:
|
||
section = (section + number) * unit
|
||
total += section
|
||
section = 0
|
||
else:
|
||
section += (number or 1) * unit
|
||
number = 0
|
||
else:
|
||
return None
|
||
return total + section + number
|
||
|
||
|
||
def _number_signature(text):
|
||
"""
|
||
提取文本的数值事实签名:返回 (计数字典, 数字拼接串)
|
||
|
||
阿拉伯数字与中文数字统一为数值,忽略书写形式差异:
|
||
例:`7月23号` 与 `七月二十三` 都得到 {7:1, 23:1} / "723",不算改动;
|
||
而 `2027年`→`2024年`、`十一届`→`九届`、`十五五`→`十四五` 会被检出。
|
||
"""
|
||
counter = {}
|
||
concat = []
|
||
# "百分之X" 与 "X%" 等价,先去掉"百分之"避免把"百"当成数值 100
|
||
text = (text or '').replace('百分之', '')
|
||
for token in _NUM_TOKEN_RE.findall(text):
|
||
if token[0].isdigit():
|
||
keys = [token.lstrip('0') or '0']
|
||
else:
|
||
keys = []
|
||
for segment in _split_cn_numeral_run(token):
|
||
value = _cn_to_number(segment)
|
||
keys.append(str(value) if value is not None else segment)
|
||
for key in keys:
|
||
counter[key] = counter.get(key, 0) + 1
|
||
concat.append(key)
|
||
return counter, ''.join(concat)
|
||
|
||
|
||
def _number_drift(raw_text, corrected_text):
|
||
"""返回 (原文独有, 修正后独有) 的数值清单;两者皆空表示数值事实未被改动"""
|
||
before, before_concat = _number_signature(raw_text)
|
||
after, after_concat = _number_signature(corrected_text)
|
||
# 计数一致,或仅切分/书写形式不同导致拼接串一致(如 2026 ↔ 二零二六),都视为未改动
|
||
if before == after or before_concat == after_concat:
|
||
return [], []
|
||
only_raw = sorted(t for t in before if before[t] > after.get(t, 0))
|
||
only_new = sorted(t for t in after if after[t] > before.get(t, 0))
|
||
return only_raw, only_new
|
||
|
||
|
||
def analyze_and_correct_text(text):
|
||
"""
|
||
分析文本并自动修正错误;**任何失败都回退原文**,保证不丢已识别文本;
|
||
**数值事实被改动时同样回退原文**(LLM 曾把 2027 年改成 2024 年、
|
||
十五五改成十四五、第十一届改成第九届,见 docs/BUGS.md B2 补充说明)。
|
||
|
||
可用 config.yml 的 llm_correct.enabled=0 完全关闭校对(默认开启)。
|
||
|
||
参数:
|
||
text (str): 待分析和修正的文本
|
||
返回值:
|
||
str: 修正后的文本;关闭/失败/改动事实时返回原文(绝不会是 None 或空串)
|
||
"""
|
||
if not text or not text.strip():
|
||
return text or ''
|
||
|
||
if not config.get_bool('llm_correct.enabled', True):
|
||
logger.info("LLM 校对已关闭(config.yml: llm_correct.enabled=0),保留 ASR 原文")
|
||
return text
|
||
|
||
logger.info("开始文本分析和修正流程...")
|
||
try:
|
||
corrected_text = text_correction(text)
|
||
if corrected_text and corrected_text.strip():
|
||
only_raw, only_new = _number_drift(text, corrected_text)
|
||
if only_raw or only_new:
|
||
logger.warning(f"⚠️ 校对改动了数值事实,已回退 ASR 原文"
|
||
f"(原文独有={only_raw} 修正后独有={only_new})")
|
||
return text
|
||
logger.info(f"原始文本长度: {len(text)},修正后: {len(corrected_text)}(数值事实一致)")
|
||
return corrected_text
|
||
logger.warning("文本修正返回空内容,使用原始文本")
|
||
except Exception as e:
|
||
logger.warning(f"文本修正失败,保留 ASR 原文(不影响已识别内容): {e}")
|
||
|
||
logger.info(f"保留原始文本长度: {len(text)}")
|
||
return text
|
||
def analyze_text(text, prompt):
|
||
"""
|
||
使用通义千问模型分析文本
|
||
参数:
|
||
text (str): 待分析文本
|
||
prompt (str): 分析提示词
|
||
返回值:
|
||
str: 分析结果
|
||
"""
|
||
logger.info("开始文本分析...")
|
||
# 设置系统提示
|
||
system_prompt = "你是一个专业的文本分析助手,擅长根据提示词对长文本进行深入分析。"
|
||
# 构建消息列表
|
||
messages = [
|
||
{"role": "system", "content": system_prompt},
|
||
{"role": "user", "content": prompt},
|
||
{"role": "user", "content": text}
|
||
]
|
||
|
||
logger.info("调用通义千问模型进行文本分析...")
|
||
# 调用DashScope文本生成接口
|
||
response = Generation.call(
|
||
model=config.model('correct_model'),
|
||
messages=messages,
|
||
max_tokens=8190, # 控制生成文本的最大长度
|
||
temperature=0.3, # 控制生成文本的确定性
|
||
top_p=0.7 # 控制生成文本的多样性
|
||
)
|
||
|
||
# 检查API调用是否成功
|
||
if response.status_code != 200:
|
||
logger.error(f"❌ API调用失败: {response.message}")
|
||
raise Exception(f"API调用失败: {response.message}")
|
||
|
||
logger.info("✓ 文本分析完成")
|
||
# 返回分析结果(兼容 output.text 与 choices[0].message.content 两种形态)
|
||
return _extract_llm_text(response)
|
||
|
||
def _upsert_daily_chunk(db, date_str, sub_id, raw, improved, title=''):
|
||
"""
|
||
按 (news_days, daily_sub_id) 覆盖写入 xwlb_daily
|
||
|
||
原实现无条件 INSERT,重跑会追加一整组 daily_sub_id=0..N 的记录
|
||
(库中 2026-06-15 全部 11 个分片已重复 2 份,见 docs/BUGS.md B5)。
|
||
这里先删同键记录再插入,无需依赖唯一索引即可实现幂等。
|
||
"""
|
||
db.execute(
|
||
"DELETE FROM xwlb_daily WHERE news_days = %s AND daily_sub_id = %s",
|
||
(date_str, sub_id)
|
||
)
|
||
return db.insert_data("xwlb_daily", {
|
||
"news_days": date_str,
|
||
"daily_sub_id": sub_id,
|
||
"news_raw": raw,
|
||
"news_improve": improved,
|
||
"news_title": title,
|
||
})
|
||
|
||
|
||
# ---- 识别完整性标记 ----
|
||
# 全部识别成功后写一个标记文件;切分失败需要重试时,重试逻辑靠它区分
|
||
# "识别已完成、只缺切分"(可只重跑切分)与"识别只跑了一半"(必须重跑全链路)。
|
||
ASR_MARKER_PREFIX = '.asr_complete_'
|
||
|
||
|
||
def asr_marker_path(date_str):
|
||
"""标记文件路径:项目下 state/ 目录
|
||
|
||
不放在 audio_processing/ 里,因为清理时会把中间产物删空;这个标记是
|
||
"该日识别已全部完成"的长期凭证,必须活过清理,否则无法与
|
||
"识别只跑了一半但切分照样写了 ext" 区分开。
|
||
"""
|
||
d8 = str(date_str).replace('-', '')
|
||
return os.path.join(config.BASE_DIR, 'state', f'{ASR_MARKER_PREFIX}{d8}')
|
||
|
||
|
||
def write_asr_marker(date_str, chunks):
|
||
try:
|
||
path = asr_marker_path(date_str)
|
||
os.makedirs(os.path.dirname(path) or '.', exist_ok=True)
|
||
with open(path, 'w', encoding='utf-8') as f:
|
||
f.write(f'chunks={chunks}\nts={datetime.datetime.now().isoformat(timespec="seconds")}\n')
|
||
except OSError as e:
|
||
logger.warning(f"写入识别完成标记失败(不影响本次结果): {e}")
|
||
|
||
|
||
def remove_asr_marker(date_str):
|
||
try:
|
||
os.remove(asr_marker_path(date_str))
|
||
except OSError:
|
||
pass
|
||
|
||
|
||
def process_long_audio(mp3_path, output_folder, date_str):
|
||
"""
|
||
处理长音频文件:分割 → 逐片识别 → 校对 → 落库
|
||
|
||
与原实现的区别:
|
||
- 识别失败/空结果**不再写入空行**,只记 ERROR 并跳过(B1);
|
||
- 每片按 (日期, 分片序号) 覆盖写入,重跑不产生重复(B5);
|
||
- 分片文件在 finally 中清理,失败分支不再残留;
|
||
- WAV 按日期命名并在结束时清理(原实现固定 input.wav,跨天互相覆盖且长期堆积)。
|
||
|
||
参数:
|
||
mp3_path (str): MP3文件路径
|
||
output_folder (str): 输出文件夹路径
|
||
date_str (str): 日期,格式 YYYYMMDD
|
||
返回值:
|
||
dict: {'chunks': 分片总数, 'ok': 成功入库数, 'failed': 失败分片数}
|
||
"""
|
||
logger.info("开始处理长音频...")
|
||
output_folder = os.fspath(output_folder)
|
||
os.makedirs(output_folder, exist_ok=True) # 原实现未建目录,目录缺失时导出会失败
|
||
|
||
# 转换MP3为WAV格式
|
||
logger.info("步骤1/3: 转换MP3为WAV格式")
|
||
wav_path = convert_mp3_to_wav(
|
||
mp3_path, os.path.join(output_folder, f"input_{date_str}.wav")
|
||
)
|
||
|
||
# 分割音频(智能静音分割,参数见 config.yml: audio_split.*)
|
||
logger.info("步骤2/3: 智能静音分割音频")
|
||
chunks = split_audio_by_smart_silence(
|
||
wav_path,
|
||
config.get_int('audio_split.min_silence_ms', 700),
|
||
config.get('audio_split.silence_thresh_db', -40),
|
||
output_folder,
|
||
)
|
||
|
||
logger.info(f"步骤3/3: 开始识别音频分片,共{len(chunks)}个分片")
|
||
ok = 0
|
||
failed = 0
|
||
db = MySQLDB() # 使用默认参数连接数据库(整个流程复用一条连接)
|
||
try:
|
||
for i, chunk_path in enumerate(chunks):
|
||
try:
|
||
logger.info(f"识别进度: {i+1}/{len(chunks)} ({((i+1)/len(chunks)*100):.1f}%)")
|
||
raw = transcribe_audio(chunk_path)
|
||
if not raw or not raw.strip():
|
||
failed += 1
|
||
logger.error(f"❌ 分片 {i+1}/{len(chunks)} 识别为空,跳过入库: {chunk_path}")
|
||
continue
|
||
improved = analyze_and_correct_text(raw)
|
||
_upsert_daily_chunk(db, date_str, i, raw, improved, '')
|
||
ok += 1
|
||
except Exception as e:
|
||
failed += 1
|
||
logger.error(f"识别失败: {chunk_path}, 错误: {e}")
|
||
finally:
|
||
# 无论成功失败都清理分片(原实现仅在成功分支删除)
|
||
try:
|
||
os.remove(chunk_path)
|
||
except OSError:
|
||
pass
|
||
|
||
# 重跑时切分点可能略有变化:若本次分片数少于库中已有记录,清理尾部残留行,
|
||
# 否则旧 sub_id 的正文会被 get_news_improve_by_date 一起拼进当天文本
|
||
try:
|
||
stale = db.execute(
|
||
"DELETE FROM xwlb_daily WHERE news_days = %s AND daily_sub_id >= %s",
|
||
(date_str, len(chunks))
|
||
)
|
||
if stale:
|
||
logger.info(f"清理 {date_str} 多余的旧分片记录 {stale} 行(本次共 {len(chunks)} 片)")
|
||
except Exception as e:
|
||
logger.warning(f"清理旧分片记录失败(不影响本次结果): {e}")
|
||
finally:
|
||
db.close()
|
||
try:
|
||
os.remove(wav_path)
|
||
except OSError:
|
||
pass
|
||
|
||
if failed:
|
||
logger.error(f"⚠️ {date_str} 长音频处理完成:成功 {ok}/{len(chunks)} 片,失败 {failed} 片(该日期内容不完整)")
|
||
remove_asr_marker(date_str)
|
||
else:
|
||
logger.info(f"✓ 长音频处理完成:成功 {ok}/{len(chunks)} 片")
|
||
# 全部识别成功才写"识别已完成"标记:重试逻辑据此判断
|
||
# 「只缺切分」还是「识别本身没跑完」,避免把半天内容当完整数据切分入库
|
||
write_asr_marker(date_str, len(chunks))
|
||
return {'chunks': len(chunks), 'ok': ok, 'failed': failed}
|
||
|
||
# 使用示例:python audioRead.py [YYYYMMDD]
|
||
if __name__ == "__main__":
|
||
import sys
|
||
from datetime import datetime
|
||
|
||
date_str = sys.argv[1] if len(sys.argv) > 1 else datetime.now().strftime('%Y%m%d')
|
||
config.log_summary()
|
||
mp3_path = os.path.join(config.video_dir(), f"{date_str}.mp3")
|
||
output_folder = config.audio_dir()
|
||
|
||
try:
|
||
result = process_long_audio(mp3_path, output_folder, date_str)
|
||
print(f"处理结果: {result}")
|
||
except Exception as e:
|
||
print(f"处理失败: {e}") |