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}")